监控
This commit is contained in:
@@ -114,14 +114,10 @@ func (s *MonitoringService) GetOverview(ctx context.Context) (*OverviewStats, er
|
||||
|
||||
// GetQueueStatuses 获取所有队列状态
|
||||
func (s *MonitoringService) GetQueueStatuses(ctx context.Context) ([]*QueueStatus, error) {
|
||||
taskTypes := []string{
|
||||
constants.TaskTypePollingRealname,
|
||||
constants.TaskTypePollingCarddata,
|
||||
constants.TaskTypePollingPackage,
|
||||
}
|
||||
taskTypes := pollingMonitorTaskTypes()
|
||||
|
||||
result := make([]*QueueStatus, 0, len(taskTypes))
|
||||
now := time.Now().Unix()
|
||||
now := time.Now()
|
||||
|
||||
for _, taskType := range taskTypes {
|
||||
status := &QueueStatus{
|
||||
@@ -142,43 +138,43 @@ func (s *MonitoringService) GetQueueStatuses(ctx context.Context) ([]*QueueStatu
|
||||
|
||||
if s.queueMgr != nil {
|
||||
status.QueueSize, _ = s.queueMgr.GetTotalQueueDepth(ctx, taskType)
|
||||
status.DueCount, _ = s.queueMgr.GetTotalDueCount(ctx, taskType, now)
|
||||
status.AvgWaitTime, _ = s.queueMgr.GetAverageWaitTime(ctx, taskType, now, 10)
|
||||
} else {
|
||||
status.QueueSize, _ = s.redis.ZCard(ctx, queueKey).Result()
|
||||
status.DueCount, _ = s.redis.ZCount(ctx, queueKey, "-inf", formatInt64(now.Unix())).Result()
|
||||
status.AvgWaitTime = s.getLegacyAverageWaitTime(ctx, queueKey, now.Unix())
|
||||
}
|
||||
|
||||
// 获取到期数量(score <= now)
|
||||
status.DueCount, _ = s.redis.ZCount(ctx, queueKey, "-inf", formatInt64(now)).Result()
|
||||
|
||||
// 获取手动触发队列待处理数
|
||||
manualKey := constants.RedisPollingManualQueueKey(taskType)
|
||||
status.ManualPending, _ = s.redis.LLen(ctx, manualKey).Result()
|
||||
|
||||
// 计算平均等待时间(取最早的10个任务的平均等待时间)
|
||||
earliest, err := s.redis.ZRangeWithScores(ctx, queueKey, 0, 9).Result()
|
||||
if err == nil && len(earliest) > 0 {
|
||||
var totalWait float64
|
||||
for _, z := range earliest {
|
||||
waitTime := float64(now) - z.Score
|
||||
if waitTime > 0 {
|
||||
totalWait += waitTime
|
||||
}
|
||||
}
|
||||
status.AvgWaitTime = totalWait / float64(len(earliest))
|
||||
}
|
||||
|
||||
result = append(result, status)
|
||||
}
|
||||
|
||||
return result, nil
|
||||
}
|
||||
|
||||
// getLegacyAverageWaitTime 从旧版非分片队列计算平均等待时间。
|
||||
func (s *MonitoringService) getLegacyAverageWaitTime(ctx context.Context, queueKey string, now int64) float64 {
|
||||
earliest, err := s.redis.ZRangeWithScores(ctx, queueKey, 0, 9).Result()
|
||||
if err != nil || len(earliest) == 0 {
|
||||
return 0
|
||||
}
|
||||
var totalWait float64
|
||||
for _, z := range earliest {
|
||||
waitTime := float64(now) - z.Score
|
||||
if waitTime > 0 {
|
||||
totalWait += waitTime
|
||||
}
|
||||
}
|
||||
return totalWait / float64(len(earliest))
|
||||
}
|
||||
|
||||
// GetTaskStatuses 获取所有任务统计
|
||||
func (s *MonitoringService) GetTaskStatuses(ctx context.Context) ([]*TaskStats, error) {
|
||||
taskTypes := []string{
|
||||
constants.TaskTypePollingRealname,
|
||||
constants.TaskTypePollingCarddata,
|
||||
constants.TaskTypePollingPackage,
|
||||
}
|
||||
taskTypes := pollingMonitorTaskTypes()
|
||||
|
||||
result := make([]*TaskStats, 0, len(taskTypes))
|
||||
|
||||
@@ -269,11 +265,26 @@ func (s *MonitoringService) getTaskTypeName(taskType string) string {
|
||||
return "流量检查"
|
||||
case constants.TaskTypePollingPackage:
|
||||
return "套餐检查"
|
||||
case constants.TaskTypePollingProtect:
|
||||
return "保护期检查"
|
||||
case constants.TaskTypePollingCardStatus:
|
||||
return "卡状态检查"
|
||||
default:
|
||||
return taskType
|
||||
}
|
||||
}
|
||||
|
||||
// pollingMonitorTaskTypes 返回后台监控需要展示的全部轮询任务类型。
|
||||
func pollingMonitorTaskTypes() []string {
|
||||
return []string{
|
||||
constants.TaskTypePollingRealname,
|
||||
constants.TaskTypePollingCarddata,
|
||||
constants.TaskTypePollingPackage,
|
||||
constants.TaskTypePollingProtect,
|
||||
constants.TaskTypePollingCardStatus,
|
||||
}
|
||||
}
|
||||
|
||||
// parseInt64 解析字符串为 int64
|
||||
func parseInt64(s string) int64 {
|
||||
var result int64
|
||||
|
||||
Reference in New Issue
Block a user