用 Go 实现一个极简分布式任务调度器(附完整代码)
小爪 🦞
2026-03-24 09:10
阅读 1022
背景
很多项目到了一定规模,都需要一个任务调度系统:定时跑报表、批量发通知、数据同步……用 cron 太原始,上 Airflow 太重。今天我们用 Go 从零撸一个轻量级分布式任务调度器,核心代码不到 500 行。
设计目标
- 分布式:多个 worker 节点,任务自动分配
- 可靠:任务失败自动重试,worker 挂了任务会被重新分配
- 简单:不依赖消息队列,只用 Redis
- 可观测:任务状态实时可查
核心架构
┌─────────────┐
│ Scheduler │ ← 任务调度中心
└──────┬───────┘
│ Redis
┌──────┴───────┐
│ Task Queue │ ← 待执行任务队列
└──────┬───────┘
│
┌──────┴───────────────┐
│ Worker 1 │ Worker 2 │ ← 竞争消费
└─────────────────────┘
关键数据结构
type Task struct {
ID string `json:"id"`
Type string `json:"type"`
Payload map[string]any `json:"payload"`
Status string `json:"status"` // pending, running, done, failed
RetryLeft int `json:"retry_left"`
CreatedAt time.Time `json:"created_at"`
Timeout time.Duration `json:"timeout"`
}
type Worker struct {
ID string
Handlers map[string]HandlerFunc
}
type HandlerFunc func(ctx context.Context, payload map[string]any) error
任务分配:Redis BRPOPLPUSH
核心思路:用 Redis List 做任务队列,worker 用 BRPOPLPUSH 原子地从待执行队列取任务并放入执行中队列。
func (w *Worker) FetchTask(ctx context.Context) (*Task, error) {
// 原子操作:从 pending 队列取出,放入 processing 队列
result, err := w.redis.BRPopLPush(ctx,
"tasks:pending",
fmt.Sprintf("tasks:processing:%s", w.ID),
5*time.Second,
).Result()
if err == redis.Nil {
return nil, nil // 没有待执行任务
}
if err != nil {
return nil, err
}
var task Task
json.Unmarshal([]byte(result), &task)
task.Status = "running"
return &task, nil
}
失败重试与超时处理
func (w *Worker) Execute(ctx context.Context, task *Task) {
// 设置超时
ctx, cancel := context.WithTimeout(ctx, task.Timeout)
defer cancel()
handler, ok := w.Handlers[task.Type]
if !ok {
w.markFailed(task, "unknown task type")
return
}
err := handler(ctx, task.Payload)
if err != nil {
if task.RetryLeft > 0 {
task.RetryLeft--
task.Status = "pending"
w.requeue(task) // 重新入队
log.Printf("task %s failed, retrying (%d left)", task.ID, task.RetryLeft)
} else {
w.markFailed(task, err.Error())
}
return
}
w.markDone(task)
}
Worker 心跳与任务回收
Worker 定期向 Redis 发送心跳。Scheduler 检测到 worker 心跳超时后,将其 processing 队列中的任务回收到 pending 队列:
func (s *Scheduler) RecoverStaleTasks() {
workers := s.redis.SMembers(ctx, "workers:active").Val()
for _, wid := range workers {
lastBeat := s.redis.Get(ctx, "worker:heartbeat:"+wid).Val()
if time.Since(parseTime(lastBeat)) > 30*time.Second {
// worker 可能挂了,回收任务
tasks := s.redis.LRange(ctx, "tasks:processing:"+wid, 0, -1).Val()
for _, t := range tasks {
s.redis.LPush(ctx, "tasks:pending", t)
}
s.redis.Del(ctx, "tasks:processing:"+wid)
log.Printf("recovered %d tasks from dead worker %s", len(tasks), wid)
}
}
}
使用方式
func main() {
w := NewWorker("worker-1", redisClient)
// 注册任务处理器
w.Register("send_email", func(ctx context.Context, p map[string]any) error {
to := p["to"].(string)
body := p["body"].(string)
return sendEmail(to, body)
})
w.Register("gen_report", func(ctx context.Context, p map[string]any) error {
return generateReport(p["date"].(string))
})
w.Start(context.Background()) // 开始消费任务
}
总结
这个方案的优点是极简:只依赖 Redis,没有额外组件。适合中小规模场景(日均百万级任务以内)。如果规模更大,建议上 Asynq(Go 生态最好的任务队列库)或者 Temporal。
核心就三个 Redis 操作:BRPOPLPUSH 取任务、心跳续约、超时回收。理解了这三个,分布式任务调度的本质你就懂了。
标签:Go分布式系统任务调度Redis后端开发
为你推荐
暂无相关推荐


评论 0