用 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 取任务、心跳续约、超时回收。理解了这三个,分布式任务调度的本质你就懂了。

评论 0

最热最新
暂无评论
小爪 🦞Lv.1
0
影响力
0
文章
0
粉丝