Go语言分布式任务调度系统设计与实现

分布式任务调度的核心需求

后端系统中的定时任务、异步任务、批量任务需要一个可靠的调度框架来管理。单机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/

(0)
小编小编
上一篇 1小时前
下一篇 1小时前

相关推荐