GoWind CMS 在 Core Service 中集成基于 Asynq 的分布式任务队列,用于处理索引同步、通知推送、数据清理等异步操作。本教程讲解 CMS 任务调度的配置、定义和管理。
| 对比项 | Admin | CMS |
|---|
| 部署位置 | Admin Service 内 | Core Service 内 |
| 任务类型 | 系统级(邮件/清理) | 内容级(索引/同步/推送) |
| 任务来源 | 后台操作 | 后台 + 前台用户操作 |
server:
asynq:
uri: "redis://:*Abcd123456@redis:6379/1"
enable_gracefully_shutdown: true
shutdown_timeout: 3s
codec: "json"
concurrency: 10
queues:
critical: 10
default: 5
low: 1
func NewAsynqServer(ctx *bootstrap.Context, logger log.Logger) *asynq.Server {
cfg := ctx.GetConfig().GetServer().GetAsynq()
srv := asynq.NewServer(
asynq.RedisClientOpt{Addr: cfg.Uri},
asynq.Config{
Concurrency: int(cfg.Concurrency),
Queues: parseQueues(cfg.Queues),
},
)
mux := asynq.NewServeMux()
mux.HandleFunc(task.TypeIndexPost, HandleIndexPost)
mux.HandleFunc(task.TypeSendNotification, HandleSendNotification)
mux.HandleFunc(task.TypeCleanupExpiredData, HandleCleanupExpiredData)
mux.HandleFunc(task.TypeGenerateReport, HandleGenerateReport)
mux.HandleFunc(task.TypeSyncTranslations, HandleSyncTranslations)
return srv
}
const (
TypeIndexPost = "post:index"
TypeIndexPage = "page:index"
TypeRebuildIndex = "search:rebuild_index"
TypeSendNotification = "notification:send"
TypeSendEmail = "email:send"
TypeCleanupExpiredData = "cleanup:expired_data"
TypeCleanupTempFiles = "cleanup:temp_files"
TypeSyncTranslations = "translation:sync"
TypeGenerateReport = "report:generate"
TypeUpdateViewCount = "stats:update_view_count"
)
type IndexPostPayload struct {
PostId uint32 `json:"postId"`
Action string `json:"action"`
}
type SendNotificationPayload struct {
UserId uint32 `json:"userId"`
Title string `json:"title"`
Content string `json:"content"`
Type string `json:"type"`
}
type CleanupExpiredDataPayload struct {
EntityType string `json:"entityType"`
OlderThan int64 `json:"olderThan"`
}
type UpdateViewCountPayload struct {
PostId uint32 `json:"postId"`
Increment uint32 `json:"increment"`
}
func HandleIndexPost(ctx context.Context, t *asynq.Task) error {
var payload task.IndexPostPayload
if err := json.Unmarshal(t.Payload(), &payload); err != nil {
return fmt.Errorf("解析载荷失败: %w", err)
}
switch payload.Action {
case "create", "update":
post, err := postRepo.Get(ctx, &contentV1.GetPostRequest{
QueryBy: &contentV1.GetPostRequest_Id{Id: payload.PostId},
})
if err != nil {
return err
}
return searchClient.Index(ctx, "cms_posts", post)
case "delete":
return searchClient.Delete(ctx, "cms_posts", fmt.Sprint(payload.PostId))
}
return nil
}
func HandleSendNotification(ctx context.Context, t *asynq.Task) error {
var payload task.SendNotificationPayload
if err := json.Unmarshal(t.Payload(), &payload); err != nil {
return err
}
_, err := internalMessageRepo.Create(ctx, &internalMessageV1.CreateInternalMessageRequest{
Data: &internalMessageV1.InternalMessage{
UserId: payload.UserId,
Title: payload.Title,
Content: payload.Content,
Type: payload.Type,
},
})
if err != nil {
return err
}
sseServer.Broadcast(ctx, payload.UserId, map[string]any{
"event": "notification",
"title": payload.Title,
"content": payload.Content,
})
return nil
}
func HandleCleanupExpiredData(ctx context.Context, t *asynq.Task) error {
var payload task.CleanupExpiredDataPayload
if err := json.Unmarshal(t.Payload(), &payload); err != nil {
return err
}
switch payload.EntityType {
case "posts":
_, err := postRepo.PermanentlyDelete(ctx, payload.OlderThan)
return err
case "comments":
_, err := commentRepo.CleanupRejected(ctx, payload.OlderThan)
return err
case "files":
return cleanupTempFiles(payload.OlderThan)
}
return nil
}
func HandleUpdateViewCount(ctx context.Context, t *asynq.Task) error {
var payload task.UpdateViewCountPayload
if err := json.Unmarshal(t.Payload(), &payload); err != nil {
return err
}
cacheKey := fmt.Sprintf("post:view:%d", payload.PostId)
count, err := redisClient.IncrBy(ctx, cacheKey, int64(payload.Increment)).Result()
if err != nil {
return err
}
if count%100 == 0 {
return postRepo.UpdateViewCount(ctx, payload.PostId, count)
}
return nil
}
func (s *PostService) Create(ctx context.Context, req *contentV1.CreatePostRequest) (*contentV1.Post, error) {
post, err := s.postRepo.Create(ctx, req)
if err != nil {
return nil, err
}
payload, _ := json.Marshal(task.IndexPostPayload{
PostId: post.Id,
Action: "create",
})
s.asynqClient.Enqueue(
asynq.NewTask(task.TypeIndexPost, payload),
asynq.Queue("default"),
)
return post, nil
}
func (s *CommentService) Create(ctx context.Context, req *commentV1.CreateCommentRequest) (*commentV1.Comment, error) {
comment, err := s.commentRepo.Create(ctx, req)
if err != nil {
return nil, err
}
payload, _ := json.Marshal(task.SendNotificationPayload{
UserId: post.AuthorId,
Title: "新评论通知",
Content: fmt.Sprintf("您的文章收到了新评论"),
Type: "comment",
})
s.asynqClient.Enqueue(
asynq.NewTask(task.TypeSendNotification, payload),
asynn.Queue("critical"),
)
return comment, nil
}
func RegisterScheduledTasks(scheduler *asynq.Scheduler) {
scheduler.Register("0 3 * * *", asynq.NewTask(
task.TypeCleanupExpiredData,
jsonMustMarshal(task.CleanupExpiredDataPayload{
EntityType: "posts",
OlderThan: time.Now().AddDate(0, 0, -30).Unix(),
}),
))
scheduler.Register("0 * * * *", asynq.NewTask(
task.TypeSyncTranslations,
nil,
))
scheduler.Register("0 9 * * 1", asynq.NewTask(
task.TypeGenerateReport,
jsonMustMarshal(task.GenerateReportPayload{
Type: "weekly",
}),
))
}
CMS 提供可视化的任务管理界面:
# 查看任务列表
GET /admin/v1/tasks?page=1&pageSize=20
# 立即执行任务
POST /admin/v1/tasks/1/run
# 暂停定时任务
PUT /admin/v1/tasks/1/pause
# 查看任务执行日志
GET /admin/v1/tasks/1/logs?page=1
Asynq 提供 Web UI 监控任务状态:
asynqmon --port=8080 --redis-addr=localhost:6379
| 监控项 | 说明 |
|---|
| Active Tasks | 正在执行的任务 |
| Pending Tasks | 等待执行的任务 |
| Retry Tasks | 重试中的任务 |
| Dead Tasks | 失败任务 |
| Scheduled Tasks | 定时任务 |
| Processed/Failed | 统计数据 |
| 原则 | 说明 |
|---|
| 幂等性 | 任务重复执行不应产生副作用 |
| 超时控制 | 设置合理的超时时间 |
| 重试策略 | 指数退避重试 |
| 优先级 | 关键任务高优先级 |
s.asynqClient.Enqueue(
asynq.NewTask(task.TypeIndexPost, payload),
asynq.Queue("default"),
asynq.MaxRetry(5),
asynq.RetryDelay(30*time.Second),
asynq.Timeout(2*time.Minute),
)
| 检查项 | 说明 |
|---|
| Asynq 配置 | Redis 连接 + 队列优先级 |
| 任务处理器注册 | 所有任务类型注册到 ServeMux |
| 任务投递 | 业务逻辑中正确投递任务 |
| 定时任务 | cron 表达式配置 |
| 任务监控 | Asynqmon Web UI |
| 重试策略 | 指数退避 + 最大重试 |
| 幂等性 | 任务可安全重试 |