| 29 | } |
| 30 | |
| 31 | func (t *TimedTask) ProcessTask(ctx context.Context, task *asynq.Task) error { |
| 32 | var payload TimedPayload |
| 33 | |
| 34 | if err := json.Unmarshal(task.Payload(), &payload); err != nil { |
| 35 | return fmt.Errorf("解析任务载荷失败: %v: %w", err, asynq.SkipRetry) |
| 36 | } |
| 37 | |
| 38 | t.l.Info("开始处理定时任务", |
| 39 | zap.String("task_name", payload.TaskName), |
| 40 | zap.Time("last_run_time", payload.LastRunTime)) |
| 41 | |
| 42 | taskCtx, cancel := context.WithTimeout(ctx, 10*time.Second) |
| 43 | defer cancel() |
| 44 | |
| 45 | // 定义任务处理映射 |
| 46 | taskHandlers := map[string]func(context.Context) error{ |
| 47 | GetRankingTask: t.svc.TopN, |
| 48 | } |
| 49 | |
| 50 | // 获取对应的处理函数 |
| 51 | handler, exists := taskHandlers[payload.TaskName] |
| 52 | if !exists { |
| 53 | return fmt.Errorf("未知的任务类型: %s", payload.TaskName) |
| 54 | } |
| 55 | |
| 56 | // 执行任务处理 |
| 57 | if err := handler(taskCtx); err != nil { |
| 58 | t.l.Error("任务执行失败", |
| 59 | zap.String("task_name", payload.TaskName), |
| 60 | zap.Error(err)) |
| 61 | return fmt.Errorf("%s: %w", payload.TaskName, err) |
| 62 | } |
| 63 | |
| 64 | t.l.Info("成功完成任务", zap.String("task_name", payload.TaskName)) |
| 65 | return nil |
| 66 | } |