分布式任务调度的核心需求
后端系统中的定时任务、异步任务、批量任务需要一个可靠的调度框架来管理。单机cron方案在面对多实例部署、任务去重、失败重试、动态调整等需求时力不从心。分布式任务调度系统需要解决的核心问题:任务的精确触发、多实例竞争执行时的去重与一致性、执行状态的可靠追踪、以及异常场景的自动恢复。
Go语言因其goroutine的轻量级并发模型和channel天然的异步通信机制,非常适合构建高性能的任务调度系统。以下从架构设计到代码实现完整展开。
系统架构与核心组件
调度系统分为四个核心模块:
– Scheduler:调度引擎,维护任务注册表,根据cron表达式或延迟时间触发任务执行
– Dispatcher:任务分发器,将待执行任务推入消息队列,实现调度与执行的解耦
– Worker:任务执行器,从队列拉取任务并执行,支持并发控制和超时管理
– Registry:分布式协调器,基于etcd实现任务注册、Leader选举、分布式锁
调度引擎实现
调度引擎负责管理所有注册任务的触发时机:
package scheduler
import (
"context"
"sync"
"time"
"github.com/robfig/cron/v3"
)
type TaskSpec struct {
ID string
Name string
CronExpr string
Timeout time.Duration
MaxRetry int
Payload json.RawMessage
}
type Scheduler struct {
cron *cron.Cron
tasks map[string]*TaskSpec
mu sync.RWMutex
dispatch func(taskID string, payload json.RawMessage)
ctx context.Context
cancel context.CancelFunc
}
func NewScheduler(dispatcher func(string, json.RawMessage)) *Scheduler {
ctx, cancel := context.WithCancel(context.Background())
return &Scheduler{
cron: cron.New(cron.WithSeconds()),
tasks: make(map[string]*TaskSpec),
dispatch: dispatcher,
ctx: ctx,
cancel: cancel,
}
}
func (s *Scheduler) Register(spec *TaskSpec) error {
s.mu.Lock()
defer s.mu.Unlock()
if _, exists := s.tasks[spec.ID]; exists {
return fmt.Errorf("task %s already registered", spec.ID)
}
entryID, err := s.cron.AddFunc(spec.CronExpr, func() {
s.dispatch(spec.ID, spec.Payload)
})
if err != nil {
return fmt.Errorf("invalid cron expression: %w", err)
}
s.tasks[spec.ID] = spec
log.Printf("registered task %s with cron [%s]", spec.ID, spec.CronExpr)
return nil
}
func (s *Scheduler) Start() {
s.cron.Start()
}
func (s *Scheduler) Stop() {
s.cancel()
stopCtx := s.cron.Stop()
<-stopCtx.Done()
}
分布式锁与任务去重
多实例部署时,同一个定时任务会被所有实例的调度引擎同时触发。必须通过分布式锁确保只有一个实例执行任务:
package lock
import (
"context"
"time"
"go.etcd.io/etcd/client/v3/concurrency"
)
type DistLock struct {
client *clientv3.Client
}
func (dl *DistLock) TryLock(ctx context.Context, taskID string, ttl time.Duration) (func(), error) {
session, err := concurrency.NewSession(dl.client,
concurrency.WithTTL(int(ttl.Seconds())),
concurrency.WithContext(ctx),
)
if err != nil {
return nil, err
}
mutex := concurrency.NewMutex(session, "/task-lock/"+taskID)
lockCtx, cancel := context.WithTimeout(ctx, 3*time.Second)
defer cancel()
if err := mutex.TryLock(lockCtx); err != nil {
session.Close()
return nil, nil
}
return func() {
mutex.Unlock(ctx)
session.Close()
}, nil
}
func (w *Worker) executeWithLock(task *Task) {
unlock, err := w.lock.TryLock(w.ctx, task.ID, task.Timeout+30*time.Second)
if err != nil || unlock == nil {
log.Printf("task %s skipped", task.ID)
return
}
defer unlock()
w.executeTask(task)
}
关键设计决策:使用TryLock而非Lock,避免Worker阻塞等待。获取锁失败意味着其他实例已开始执行,直接跳过即可。锁的TTL设置为任务超时时间加30秒,防止实例崩溃后锁永不释放。
任务执行器与重试策略
Worker从消息队列拉取任务后执行,支持并发控制、超时管理和指数退避重试:
type Worker struct {
sem chan struct{}
handlers map[string]Handler
lock *lock.DistLock
ctx context.Context
}
type Handler func(ctx context.Context, payload json.RawMessage) error
func (w *Worker) processMessage(msg amqp.Delivery) {
var task Task
json.Unmarshal(msg.Body, &task)
handler, ok := w.handlers[task.Type]
if !ok {
msg.Nack(false, false)
return
}
ctx, cancel := context.WithTimeout(w.ctx, task.Timeout)
defer cancel()
err := w.executeWithRetry(ctx, handler, &task)
if err != nil {
if task.RetryCount < task.MaxRetry {
w.requeueWithDelay(&task)
} else {
w.markFailed(&task, err)
}
msg.Nack(false, false)
return
}
msg.Ack(false)
}
func (w *Worker) executeWithRetry(ctx context.Context, handler Handler, task *Task) error {
var lastErr error
for attempt := 0; attempt <= task.RetryCount; attempt++ {
if attempt > 0 {
delay := time.Duration(math.Pow(2, float64(attempt))) * time.Second
if delay > 5*time.Minute {
delay = 5 * time.Minute
}
select {
case <-time.After(delay):
case <-ctx.Done():
return ctx.Err()
}
}
lastErr = handler(ctx, task.Payload)
if lastErr == nil {
return nil
}
task.RetryCount++
}
return lastErr
}
任务状态追踪与可观测性
每个任务的执行状态需要持久化存储并支持查询。在Redis中维护任务状态的实时快照,在数据库中保存历史执行记录:
type TaskStatus struct {
ID string `json:"id"`
TaskID string `json:"task_id"`
Status string `json:"status"`
WorkerID string `json:"worker_id"`
StartTime time.Time `json:"start_time"`
EndTime time.Time `json:"end_time"`
RetryCount int `json:"retry_count"`
Error string `json:"error,omitempty"`
}
func (r *RedisTracker) UpdateStatus(status *TaskStatus) error {
key := fmt.Sprintf("task:status:%s", status.ID)
data, _ := json.Marshal(status)
return r.client.Set(r.ctx, key, data, 2*time.Hour).Err()
}
func (db *DBTracker) ArchiveStatus(status *TaskStatus) error {
_, err := db.Exec(`
INSERT INTO task_executions
(id, task_id, status, worker_id, start_time, end_time, retry_count, error)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
`, status.ID, status.TaskID, status.Status, status.WorkerID,
status.StartTime, status.EndTime, status.RetryCount, status.Error)
return err
}
配合Prometheus指标暴露,可以构建完整的任务执行大盘:调度触发次数、执行成功/失败率、平均执行耗时、队列积压深度、Worker并发利用率等关键指标实时可观测。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/go-yu-yan-fen-bu-shi-ren-wu-diao-du-xi-tong-she-ji-yu-shi/