# Go语言分布式任务调度系统实战:基于asynq构建生产级定时任务管理平台 ## 一、为什么需要分布式任务调度 在微服务架构中,定时任务(Crontab)面临着三个核心痛点: 1. **单点故障**:任务绑定在某台服务器上,机器宕机则任务丢失 2. **缺乏可观测性**:任务是否执行成功、耗时多久、失败原因全靠日志猜测 3. **重试机制薄弱**:传统crontab失败后不会自动重试,也不会告警 本文将使用Go语言和开源库[asynq](https://github.com/hibiken/asynq),从零构建一个支持分布式部署、自动失败重试、拥有Web监控面板、暴露Prometheus指标的生产级任务调度平台。 ## 二、技术选型与架构 ``` ┌─────────────┐ ┌──────────────┐ ┌─────────────────┐ │ Producer │────▶│ Redis队列 │────▶│ Consumer Pool │ │ (API/CLI) │ │ (Streams) │ │ (动态扩缩容) │ └─────────────┘ └──────────────┘ └─────────────────┘ │ │ ▼ ▼ ┌──────────────┐ ┌─────────────────┐ │ Scheduler │ │ Prometheus │ │ (Cron触发) │ │ Metrics + HTTP │ └──────────────┘ └─────────────────┘ ``` **核心依赖:** - `asynq`:任务队列客户端 - `asynqmon`:Web监控面板 - `redis`:作为broker(消息代理) - `prometheus client`:暴露任务指标 ## 三、项目结构 ``` distributed-scheduler/ ├── cmd/ │ ├── producer/ # 任务生产者(API服务) │ └── worker/ # 任务消费者(Worker服务) ├── internal/ │ ├── tasks/ # 任务定义与处理器 │ │ ├── email.go │ │ ├── report.go │ │ └── image.go │ ├── middleware/ # 中间件(日志、指标、重试) │ └── config/ # 配置管理 ├── go.mod ├── go.sum ├── docker-compose.yml └── prometheus.yml ``` ## 四、核心代码实现 ### 4.1 定义任务类型与处理器 ```go // internal/tasks/email.go package tasks import ( "context" "encoding/json" "fmt" "log" "github.com/hibiken/asynq" ) const ( TypeEmailDelivery = "email:deliver" TypeImageResize = "image:resize" TypeReportGenerate = "report:generate" ) // EmailDeliveryPayload 邮件发送任务载荷 type EmailDeliveryPayload struct { UserID int `json:"user_id"` TemplateID string `json:"template_id"` To string `json:"to"` } // NewEmailTask 创建邮件发送任务 func NewEmailTask(p EmailDeliveryPayload) (*asynq.Task, error) { payload, err := json.Marshal(p) if err != nil { return nil, err } return asynq.NewTask(TypeEmailDelivery, payload, asynq.MaxRetry(5)), nil } // HandleEmailDeliveryTask 邮件发送处理器 func HandleEmailDeliveryTask(ctx context.Context, t *asynq.Task) error { var p EmailDeliveryPayload if err := json.Unmarshal(t.Payload(), &p); err != nil { return fmt.Errorf("json.Unmarshal failed: %v: %w", err, asynq.SkipRetry) } log.Printf("[📧] Sending Email to=%s user_id=%d template=%s", p.To, p.UserID, p.TemplateID) // 实际调用的邮件发送逻辑 return sendEmail(p.To, p.TemplateID) } func sendEmail(to, templateID string) error { // 集成SendGrid / SMTP 等 return nil } ``` ```go // internal/tasks/report.go package tasks import ( "context" "encoding/json" "fmt" "log" "github.com/hibiken/asynq" ) type ReportPayload struct { ReportType string `json:"report_type"` UserID int `json:"user_id"` DateRange string `json:"date_range"` } func NewReportTask(p ReportPayload) (*asynq.Task, error) { payload, err := json.Marshal(p) if err != nil { return nil, err } // 报表生成耗时较长,设置超时60s,最多重试3次 return asynq.NewTask(TypeReportGenerate, payload, asynq.MaxRetry(3), asynq.Timeout(60), asynq.Queue("critical"), ), nil } func HandleReportTask(ctx context.Context, t *asynq.Task) error { var p ReportPayload if err := json.Unmarshal(t.Payload(), &p); err != nil { return fmt.Errorf("json.Unmarshal failed: %v: %w", err, asynq.SkipRetry) } log.Printf("[📊] Generating report type=%s user_id=%d range=%s", p.ReportType, p.UserID, p.DateRange) // 实际报表生成逻辑 return generateReport(p) } func generateReport(p ReportPayload) error { return nil } ``` ### 4.2 生产者(任务入队) ```go // cmd/producer/main.go package main import ( "log" "net/http" "github.com/gin-gonic/gin" "github.com/hibiken/asynq" "distributed-scheduler/internal/tasks" ) var client *asynq.Client func init() { client = asynq.NewClient(asynq.RedisClientOpt{ Addr: "localhost:6379", Password: "", DB: 0, }) } func main() { r := gin.Default() // API: 创建邮件发送任务 r.POST("/api/v1/tasks/email", func(c *gin.Context) { var payload tasks.EmailDeliveryPayload if err := c.ShouldBindJSON(&payload); err != nil { c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) return } task, err := tasks.NewEmailTask(payload) if err != nil { c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) return } info, err := client.Enqueue(task) if err != nil { c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) return } log.Printf("[+] Email task enqueued: id=%s queue=%s", info.ID, info.Queue) c.JSON(http.StatusOK, gin.H{ "task_id": info.ID, "queue": info.Queue, "status": "enqueued", }) }) // API: 创建报表生成任务(critical队列) r.POST("/api/v1/tasks/report", func(c *gin.Context) { var payload tasks.ReportPayload if err := c.ShouldBindJSON(&payload); err != nil { c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) return } task, err := tasks.NewReportTask(payload) if err != nil { c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) return } info, err := client.Enqueue(task) if err != nil { c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) return } log.Printf("[+] Report task enqueued: id=%s queue=%s", info.ID, info.Queue) c.JSON(http.StatusOK, gin.H{ "task_id": info.ID, "queue": info.Queue, "status": "enqueued", }) }) r.Run(":8080") } ``` ### 4.3 Worker(任务消费与路由) ```go // cmd/worker/main.go package main import ( "context" "log" "github.com/hibiken/asynq" "github.com/prometheus/client_golang/prometheus" "github.com/prometheus/client_golang/prometheus/promhttp" "distributed-scheduler/internal/middleware" "distributed-scheduler/internal/tasks" ) // Prometheus指标 var ( taskProcessed = prometheus.NewCounterVec( prometheus.CounterOpts{ Name: "asynq_tasks_processed_total", Help: "Total number of tasks processed", }, []string{"task_type", "status"}, ) taskDuration = prometheus.NewHistogramVec( prometheus.HistogramOpts{ Name: "asynq_task_duration_seconds", Help: "Task processing duration in seconds", Buckets: []float64{0.01, 0.05, 0.1, 0.5, 1, 2, 5, 10}, }, []string{"task_type"}, ) ) func init() { prometheus.MustRegister(taskProcessed, taskDuration) } func main() { srv := asynq.NewServer( asynq.RedisClientOpt{Addr: "localhost:6379"}, asynq.Config{ // 并发数:控制同时处理的任务数 Concurrency: 20, // 优先级队列配置 Queues: map[string]int{ "critical": 6, // 60% 资源给critical "default": 3, // 30% 资源给default "low": 1, // 10% 资源给low }, // 严格优先级模式 StrictPriority: true, // 错误回调 ErrorHandler: func(ctx context.Context, task *asynq.Task, err error) { log.Printf("[ERROR] task=%s err=%v", task.Type(), err) taskProcessed.WithLabelValues(task.Type(), "failure").Inc() }, // 重试延迟策略:指数退避 RetryDelayFunc: func(n int, err error, t *asynq.Task) time.Duration { return time.Duration(n*n) * time.Second }, }, ) // 注册任务处理器 mux := asynq.NewServeMux() mux.Use(middleware.Metrics(taskProcessed, taskDuration)) mux.HandleFunc(tasks.TypeEmailDelivery, tasks.HandleEmailDeliveryTask) mux.HandleFunc(tasks.TypeReportGenerate, tasks.HandleReportTask) mux.HandleFunc(tasks.TypeImageResize, tasks.HandleImageResizeTask) // Prometheus指标服务器(与业务共用端口,不同路径) go func() { http.Handle("/metrics", promhttp.Handler()) log.Println("[Prometheus] listening on :9090") log.Fatal(http.ListenAndServe(":9090", nil)) }() log.Println("[Worker] starting, concurrency=20") if err := srv.Run(mux); err != nil { log.Fatalf("could not run server: %v", err) } } ``` ### 4.4 中间件:指标与日志 ```go // internal/middleware/metrics.go package middleware import ( "context" "time" "github.com/hibiken/asynq" "github.com/prometheus/client_golang/prometheus" ) // Metrics 记录任务处理时长和计数中间件 func Metrics(counter *prometheus.CounterVec, histogram *prometheus.HistogramVec) func(handler) handler { return func(next handler) handler { return func(ctx context.Context, task *asynq.Task) error { start := time.Now() err := next(ctx, task) duration := time.Since(start).Seconds() status := "success" if err != nil { status = "failure" } counter.WithLabelValues(task.Type(), status).Inc() histogram.WithLabelValues(task.Type()).Observe(duration) return err } } } ``` ### 4.5 定时任务调度器(Cron) ```go // cmd/worker/scheduler.go package main import ( "log" "github.com/hibiken/asynq" "distributed-scheduler/internal/tasks" ) // StartScheduler 启动定时任务调度器 func StartScheduler() { sched := asynq.NewScheduler(asynq.RedisClientOpt{Addr: "localhost:6379"}, nil) // 每小时执行一次:生成销售报表 reportTask, _ := tasks.NewReportTask(tasks.ReportPayload{ ReportType: "sales_hourly", DateRange: "last_1_hour", }) entryID, err := sched.Register("@every 1h", reportTask, asynq.Queue("critical")) if err != nil { log.Fatalf("failed to register hourly report: %v", err) } log.Printf("[Scheduler] hourly report registered: %s", entryID) // 每天凌晨2点执行:数据画像 portraitTask, _ := tasks.NewReportTask(tasks.ReportPayload{ ReportType: "user_portrait", DateRange: "yesterday", }) entryID, err = sched.Register("0 2 * * *", portraitTask, asynq.Queue("low")) if err != nil { log.Fatalf("failed to register daily portrait: %v", err) } log.Printf("[Scheduler] daily portrait registered: %s", entryID) // 每5分钟执行:未发送邮件重试 retryTask, _ := tasks.NewEmailTask(tasks.EmailDeliveryPayload{ TemplateID: "retry_unsent", }) entryID, err = sched.Register("*/5 * * * *", retryTask) if err != nil { log.Fatalf("failed to register retry cron: %v", err) } log.Printf("[Scheduler] retry cron registered: %s", entryID) if err := sched.Run(); err != nil { log.Fatalf("scheduler failed: %v", err) } } ``` ### 4.6 Docker Compose一键部署 ```yaml # docker-compose.yml version: '3.8' services: redis: image: redis:7-alpine ports: - "6379:6379" command: redis-server --appendonly yes --maxmemory 256mb --maxmemory-policy allkeys-lru producer: build: context: . dockerfile: Dockerfile args: CMD: producer ports: - "8080:8080" depends_on: - redis restart: always worker-1: build: context: . dockerfile: Dockerfile args: CMD: worker environment: - CONCURRENCY=20 depends_on: - redis restart: always deploy: resources: limits: memory: 512M cpus: '1.0' worker-2: build: context: . dockerfile: Dockerfile args: CMD: worker environment: - CONCURRENCY=20 depends_on: - redis restart: always # Web监控面板 asynqmon: image: hibiken/asynqmon:latest ports: - "8081:8080" environment: - REDIS_ADDR=redis:6379 depends_on: - redis prometheus: image: prom/prometheus:latest volumes: - ./prometheus.yml:/etc/prometheus/prometheus.yml ports: - "9090:9090" grafana: image: grafana/grafana:latest ports: - "3000:3000" environment: - GF_SECURITY_ADMIN_PASSWORD=admin ``` ```yaml # prometheus.yml scrape_configs: - job_name: 'asynq-worker' static_configs: - targets: ['worker-1:9090', 'worker-2:9090'] ``` ## 五、关键调优技巧 ### 5.1 Concurrency调优 并发数设置公式:`Concurrency = CPU核数 × (平均I/O等待时间 / 平均CPU时间 + 1)` CPU密集型任务设为`runtime.NumCPU()`,I/O密集型(如HTTP请求)可设为`CPU核数 × 5~10`。 ### 5.2 队列优先级与资源分配 ```go Queues: map[string]int{ "critical": 6, // 核心业务:支付回调、短信验证码 "default": 3, // 普通业务:日志汇总 "low": 1, // 后台任务:数据清理、冷数据归档 }, StrictPriority: true, // 先取critical队列,再取default ``` ### 5.3 失败重试策略 asynq默认使用指数退避,第n次重试延迟为 `2^n` 秒(最长30分钟): - 第1次:2s - 第2次:4s - 第3次:16s - 第4次:64s - 第5次:256s 对于幂等性敏感的任务(如扣款),务必设置 `asynq.MaxRetry(3)` 限制上限。 ### 5.4 任务去重 asynq提供`Unique`选项防止重复入队: ```go task, _ := tasks.NewEmailTask(payload) client.Enqueue(task, asynq.Unique(5*time.Minute), // 5分钟内相同任务不入队 asynq.TaskID("custom-id"), // 自定义ID去重 ) ``` ## 六、生产环境最佳实践 1. **Redis持久化**:必须开启`AOF`,推荐每秒刷盘,避免队列数据丢失 2. **熔断降级**:当任务积压超过阈值时,触发告警并自动限流 3. **队列隔离**:不同业务使用不同队列,避免低优先级任务阻塞高优先级 4. **死信队列**:超过最大重试的任务存入死信队列(RetryLaterer接口),人工介入处理 5. **优雅关闭**:捕获SIGTERM信号,等待当前任务完成再退出 ```go srv.Shutdown() // asynq.Server支持graceful shutdown ``` ## 七、总结 相比传统crontab方案,基于asynq的调度系统提供: - ✅ 分布式部署,无单点故障 - ✅ 自动失败重试 + 指数退避 - ✅ 优先级队列 + 弹性扩缩容 - ✅ Web面板实时监控(asynqmon) - ✅ Prometheus指标 + Grafana可视化 - ✅ 任务去重 + 严格ID保证幂等 整个系统仅依赖Redis,部署成本极低,非常适合Go语言栈团队作为轻量级任务调度方案落地使用。

发表评论 取消回复