This commit is contained in:
ryan
2026-06-08 15:13:59 +08:00
parent f136c9abdb
commit b03d2d7ea6
9 changed files with 201 additions and 73 deletions
+175 -47
View File
@@ -1,6 +1,6 @@
---
name: "new-async-task"
description: "项目专用:指导如何在本项目中新增异步任务(TaskHandler),包括常量定义、处理器实现、注册、Cron 调度、AppendLog 日志规范等完整步骤。"
description: "项目专用:指导如何在本项目中新增异步任务(TaskHandler),包括常量定义、参数传递(TaskParam)、处理器实现、Admin 校验、注册、Cron 调度、AppendLog 日志规范等完整步骤。"
---
# 新增异步任务开发指南
@@ -27,11 +27,12 @@ TaskHandler.Execute (← 你写这一步)
**关键文件**:
- `internal/task/handler.go` — `TaskHandler` 接口、`TaskResult` 类型
- `internal/task/constants.go` — 任务类型常量、`TaskMeta`、`DispatchableTasks`
- `internal/task/constants.go` — 任务类型常量、`TaskMeta`、`TaskParam`、`DispatchableTasks`
- `internal/task/executor.go` — 注册、下发、执行、日志追加
- `internal/task/worker/worker.go` — Worker 启动、处理器注册、路由
- `internal/task/scheduler/scheduler.go` — Cron 定时调度
- `internal/model/task_execution.go` — `TaskExecution` GORM 模型
- `internal/apps/admin/task/routers.go` — 管理 API(下发、查询、重试)
## 新增任务的完整步骤
@@ -53,51 +54,116 @@ const TaskTypeCleanupUploads = "cleanup_unused_uploads"
**3. 在 `DispatchableTasks` 切片中添加 TaskMeta**:
不带参数的任务:
```go
var DispatchableTasks = []TaskMeta{
{
Type: TaskTypeCleanupUploads, // 管理员 API 用的标识
AsynqTask: CleanupUnusedUploadsTask, // Asynq 路由用的标识
Name: "清理未使用上传", // 管理后台显示名
Description: "清理超过1小时未使用的上传文件", // 管理后台显示描述
SupportsTime: false, // 是否支持时间范围参数(暂未使用)
MaxRetry: 3, // Asynq 最大自动重试次数
Queue: QueueDefault, // 投递到的队列名
Retryable: true, // 管理端是否允许手动重试
{
Type: TaskTypeCleanupUploads,
AsynqTask: CleanupUnusedUploadsTask,
Name: "清理未使用上传",
Description: "清理超过1小时未使用的上传文件",
SupportsTime: false,
MaxRetry: 3,
Queue: QueueDefault,
Retryable: true,
}
```
带参数的任务(`Params` 定义前端表单字段):
```go
{
Type: TaskTypeSendEmail,
AsynqTask: SendEmailTask,
Name: "发送邮件",
Description: "异步发送系统邮件",
SupportsTime: false,
MaxRetry: 3,
Queue: QueueDefault,
Retryable: true,
Params: []TaskParam{
{
Name: "to",
Label: "接收邮箱 (To)",
Type: "string",
Required: true,
Placeholder: "receiver@example.com",
Description: "接收邮件的目标邮箱地址",
},
{
Name: "subject",
Label: "邮件主题 (Subject)",
Type: "string",
Required: true,
Placeholder: "请输入邮件主题",
Description: "发送邮件的主题标题",
},
{
Name: "body",
Label: "邮件内容 (Body)",
Type: "text",
Required: true,
Placeholder: "请输入邮件内容(支持 HTML 格式)",
Description: "发送邮件的内容主体",
},
},
}
```
**`TaskParam` 字段说明**:
| 字段 | 作用 |
|------|------|
| `Name` | 参数键名,与 Handler 中 Payload 结构体的 JSON tag 一致 |
| `Label` | 前端表单显示的标签 |
| `Type` | 前端控件类型:`string`(单行输入)、`text`(多行文本)、`number`(数字) |
| `Required` | 前端是否标为必填(仅前端提示,不做服务端校验) |
| `Placeholder` | 前端输入框占位文字 |
| `Description` | 前端显示的参数说明 |
**注意**:`TaskParam` 纯粹是前端表单元数据。服务端不基于它做参数校验——校验逻辑在 Admin dispatch handler 和 Handler Execute 中各写一份。
### 第 2 步:实现 TaskHandler
在 `internal/apps/<module>/` 下创建 `tasks.go`,定义结构体并实现 `TaskHandler` 接口:
在 `internal/apps/<module>/` 下创建 `tasks.go`。
**无参数的任务**(参考 `internal/apps/upload/tasks.go`):
```go
package upload
import (
"context"
"fmt"
"time"
"github.com/linux-do/credit/internal/db"
"github.com/linux-do/credit/internal/model"
"github.com/linux-do/credit/internal/task"
)
// MyTaskHandler 你的任务处理器
type CleanupUnusedUploadsHandler struct{}
// Execute 实现 TaskHandler 接口
// - ctx: 已注入 OTel Span 和 taskID,传给 task.AppendLog 使用
// - payload: 下发时传入的原始字节(可为 nil)
func (h *CleanupUnusedUploadsHandler) Execute(ctx context.Context, payload []byte) (*task.TaskResult, error) {
task.AppendLog(ctx, "开始扫描未使用上传文件,阈值: %s", time.Now().Add(-1*time.Hour).Format(time.RFC3339))
task.AppendLog(ctx, "开始扫描...")
// ... 直接执行业务逻辑,忽略 payload ...
return &task.TaskResult{Message: "完成"}, nil
}
```
**带参数的任务**(参考 `internal/apps/user/tasks.go`):
需要定义自己的 Payload 结构体,在 `Execute` 中反序列化:
```go
// 定义载荷结构体(字段与 TaskParam.Name 对应)
type SendEmailPayload struct {
To string `json:"to"`
Subject string `json:"subject"`
Body string `json:"body"`
}
type SendEmailHandler struct{}
func (h *SendEmailHandler) Execute(ctx context.Context, payload []byte) (*task.TaskResult, error) {
// 第一步:反序列化参数
var req SendEmailPayload
if err := json.Unmarshal(payload, &req); err != nil {
task.AppendLog(ctx, "解析参数失败: %v", err)
return nil, fmt.Errorf("解析参数失败: %w", err)
}
task.AppendLog(ctx, "开始发送邮件到: %s, 主题: %s", req.To, req.Subject)
// ... 业务逻辑 ...
// 如需处理批量数据,使用游标分页模式(参见 upload/tasks.go)
msg := fmt.Sprintf("共处理 %d 个文件,成功删除 %d 个", totalProcessed, totalDeleted)
msg := fmt.Sprintf("邮件成功发送至: %s", req.To)
task.AppendLog(ctx, "%s", msg)
return &task.TaskResult{Message: msg}, nil
}
@@ -107,7 +173,73 @@ func (h *CleanupUnusedUploadsHandler) Execute(ctx context.Context, payload []byt
- 成功:返回 `&task.TaskResult{Message: "摘要", Detail: "可选详细JSON"}`。`ProcessTask` 会将 `Message + "\n" + Detail` 存入 `TaskExecution.Result`。
- 失败:返回 `nil, fmt.Errorf("...")`。`ProcessTask` 会将错误信息存入 `TaskExecution.ErrorMessage`,状态标记为 `failed`,并将 error 返回给 Asynq 触发自动重试。
### 第 3 步:在 worker.go 注册
### 第 3 步:在 Admin dispatch handler 中添加参数校验
打开 `internal/apps/admin/task/routers.go`,在 `DispatchTask` 函数中为新任务类型添加校验逻辑。
下发请求结构体中,前端传入的参数通过 `Payload` 字段(JSON 字符串)传递:
```go
type DispatchTaskRequest struct {
TaskType string `json:"task_type" binding:"required"`
StartTime *time.Time `json:"start_time"`
EndTime *time.Time `json:"end_time"`
UserID *uint64 `json:"user_id"`
Payload string `json:"payload"` // ← 前端传的参数 JSON 字符串
}
```
在 `DispatchTask` 函数中按 task type 分支校验:
```go
var payloadBytes []byte
if req.TaskType == task.TaskTypeSendEmail {
// 校验 payload 非空
if strings.TrimSpace(req.Payload) == "" {
c.JSON(http.StatusBadRequest, util.Err("任务参数 Payload 不能为空"))
return
}
// 解析 JSON 并校验必填字段
var mailPayload struct {
To string `json:"to"`
Subject string `json:"subject"`
Body string `json:"body"`
}
if err := json.Unmarshal([]byte(req.Payload), &mailPayload); err != nil {
c.JSON(http.StatusBadRequest, util.Err("无效的 JSON 格式: "+err.Error()))
return
}
// Trim + 必填校验
mailPayload.To = strings.TrimSpace(mailPayload.To)
mailPayload.Subject = strings.TrimSpace(mailPayload.Subject)
mailPayload.Body = strings.TrimSpace(mailPayload.Body)
if mailPayload.To == "" || mailPayload.Subject == "" || mailPayload.Body == "" {
c.JSON(http.StatusBadRequest, util.Err("to、subject、body 不能为空"))
return
}
// 序列化回 bytes 传给 DispatchTask
payloadBytes, _ = json.Marshal(mailPayload)
} else {
// 无参数任务:直接透传 payload
if req.Payload != "" {
payloadBytes = []byte(req.Payload)
}
}
taskID, err := task.DispatchTask(c.Request.Context(), req.TaskType, payloadBytes, "manual")
```
**参数传递的完整链路**:
```
前端表单(根据 TaskMeta.Params 动态渲染)
→ POST /api/v1/admin/tasks/dispatch { task_type, payload: JSON字符串 }
→ Admin DispatchTask handler 校验 + json.Marshal
→ task.DispatchTask(payloadBytes) 存入 DB + 入队 Asynq
→ ProcessTask → handler.Execute(payload) → handler 内 json.Unmarshal
```
### 第 4 步:在 worker.go 注册
打开 `internal/task/worker/worker.go`,做两件事:
@@ -116,8 +248,7 @@ func (h *CleanupUnusedUploadsHandler) Execute(ctx context.Context, payload []byt
```go
func init() {
task.RegisterHandler(task.CleanupUnusedUploadsTask, &upload.CleanupUnusedUploadsHandler{})
// 添加你的:
task.RegisterHandler(task.YourNewTask, &yourmodule.YourHandler{})
task.RegisterHandler(task.SendEmailTask, &user.SendEmailHandler{})
}
```
@@ -127,21 +258,20 @@ func init() {
mux := asynq.NewServeMux()
mux.Use(taskLoggingMiddleware)
mux.HandleFunc(task.CleanupUnusedUploadsTask, task.ProcessTask)
// 添加你的:
mux.HandleFunc(task.YourNewTask, task.ProcessTask)
mux.HandleFunc(task.SendEmailTask, task.ProcessTask)
```
所有路由都指向同一个 `task.ProcessTask`,它内部根据 Asynq task type 查找对应的 handler 分发执行。
### 第 4 步(可选):添加 Cron 定时调度
### 第 5 步(可选):添加 Cron 定时调度
如果任务需要定时执行,编辑 `internal/task/scheduler/scheduler.go`,在 `StartScheduler()` 中注册:
```go
if _, err = scheduler.Register(
config.Config.Scheduler.CleanupUnusedUploadsTaskCron, // cron 表达式
config.Config.Scheduler.CleanupUnusedUploadsTaskCron,
asynq.NewTask(task.CleanupUnusedUploadsTask, nil),
asynq.Unique(23*time.Hour), // 防止重复执行
asynq.Unique(23*time.Hour),
asynq.MaxRetry(3),
); err != nil {
return
@@ -178,6 +308,7 @@ scheduler:
| 时机 | 示例 |
|------|------|
| 任务开始时 | `task.AppendLog(ctx, "开始扫描,阈值: %s", threshold)` |
| 参数解析成功后 | `task.AppendLog(ctx, "开始发送邮件到: %s, 主题: %s", req.To, req.Subject)` |
| 批量处理的每一批次 | `task.AppendLog(ctx, "本批次处理 %d 条记录", len(batch))` |
| 关键中间状态 | `task.AppendLog(ctx, "已删除对象 %s", filePath)` |
| 遇到可继续的错误时 | `task.AppendLog(ctx, "清理文件失败 [ID:%d]: %v", id, err)` |
@@ -199,6 +330,7 @@ scheduler:
- **OTel 追踪**:自动创建 Span,记录任务类型、Payload 大小、TaskID
- **重试计数**:手动重试时自动递增 `RetryCount`,校验 `RetryCount < MaxRetry`
- **队列路由**:根据 `TaskMeta.Queue` 投递到对应优先级队列
- **前端表单渲染**:`ListTaskTypes` API 返回 `DispatchableTasks`(含 `Params`),前端据此动态渲染参数表单
## 重试机制
@@ -206,14 +338,10 @@ scheduler:
**Asynq 自动重试**:`ProcessTask` 返回 error 时,Asynq 根据入队时设置的 `MaxRetry` 自动重试。但 `TaskExecution` 记录在第一次失败时已标记为 `failed`,后续重试会覆盖同一条记录。
**管理端手动重试**:通过 `POST /api/v1/admin/tasks/executions/{id}/retry` 触发。会创建一个全新的 `TaskExecution` 记录(新 TaskID `retry_{count}_{originalTaskID}`),`RetryCount` 递增。校验条件:状态必须为 `failed`、`Retryable` 为 `true`、`RetryCount < MaxRetry`。
**管理端手动重试**:通过 `POST /api/v1/admin/tasks/executions/{id}/retry` 触发。会创建一个全新的 `TaskExecution` 记录(新 TaskID `retry_{count}_{originalTaskID}`),`RetryCount` 递增,`Payload` 原样复制。校验条件:状态必须为 `failed`、`Retryable` 为 `true`、`RetryCount < MaxRetry`。
## 参考:现有任务处理器
完整的任务处理器示例参见 `internal/apps/upload/tasks.go` 中的 `CleanupUnusedUploadsHandler`,它展示了:
- 游标分页批量处理
- 每批次的 `AppendLog` 记录
- 事务内操作(DB 状态更新 + 外部调用)
- 单条失败时记录日志并 continue 而非整体终止
- 最终汇总结果通过 `TaskResult.Message` 返回
- `internal/apps/upload/tasks.go` — `CleanupUnusedUploadsHandler`:无参数任务,展示游标分页批量处理、每批次 AppendLog、事务内操作、单条失败 continue 不终止。
- `internal/apps/user/tasks.go` — `SendEmailHandler`:带参数任务,展示 Payload 结构体定义、json.Unmarshal 反序列化、参数日志记录。