# 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语言栈团队作为轻量级任务调度方案落地使用。

发表评论 取消回复