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

分布式任务调度系统需要在多节点环境下保证任务不被重复执行、故障节点任务能被自动接管。Go语言凭借goroutine轻量并发和channel通信模型,适合构建高吞吐调度引擎。etcd作为强一致性的分布式KV存储,提供Lease租约和Watch机制,是实现分布式锁和任务分发的理想基础设施。

下面从架构设计到核心代码,完整实现一个基于Go和etcd的分布式定时任务调度器。

系统架构设计

调度系统由三个核心组件构成:

  • Scheduler:调度引擎,每个节点运行一个实例,通过cron表达式触发任务,竞争分布式锁获取执行权
  • Worker Pool:goroutine工作池,接收调度任务执行,控制并发度,管理超时和重试
  • etcd Registry:任务注册表和节点注册表,存储任务定义和节点存活状态,通过Watch通知变更

任务执行流程:所有Scheduler节点同时收到cron触发,各自向etcd尝试获取该任务的分布式锁,只有一个节点获锁成功并执行任务,其余节点放弃。持锁节点通过Lease保持心跳,执行完成后释放锁。若持锁节点宕机,Lease过期锁自动释放,下次触发时其他节点竞争获锁。

etcd分布式锁实现

etcd分布式锁基于Revision排序+Lease租约实现。多个客户端在同一个prefix下创建带租约的KV,revision最小者获锁成功,其余客户端Watch前一个KV的删除事件排队等锁。

package scheduler

import (
    "context"
    "fmt"
    "log"
    "time"

    "go.etcd.io/etcd/client/v3"
    "go.etcd.io/etcd/client/v3/concurrency"
)

type DistributedLock struct {
    client *clientv3.Client
    ttl    int
}

func NewDistributedLock(endpoints []string, ttl int) (*DistributedLock, error) {
    cli, err := clientv3.New(clientv3.Config{
        Endpoints:   endpoints,
        DialTimeout: 5 * time.Second,
    })
    if err != nil {
        return nil, fmt.Errorf("connect etcd failed: %w", err)
    }
    return &DistributedLock{client: cli, ttl: ttl}, nil
}

func (dl *DistributedLock) TryLock(ctx context.Context, key string) (func(), error) {
    session, err := concurrency.NewSession(dl.client,
        concurrency.WithTTL(dl.ttl),
        concurrency.WithContext(ctx),
    )
    if err != nil {
        return nil, fmt.Errorf("create session failed: %w", err)
    }

    mutex := concurrency.NewMutex(session, key)
    err = mutex.TryLock(ctx)
    if err != nil {
        session.Close()
        if err == concurrency.ErrLocked {
            return nil, nil
        }
        return nil, fmt.Errorf("try lock failed: %w", err)
    }

    log.Printf("[Lock] acquired: %s", key)

    unlock := func() {
        ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
        defer cancel()
        mutex.Unlock(ctx)
        session.Close()
        log.Printf("[Lock] released: %s", key)
    }

    return unlock, nil
}

concurrency.NewMutex底层实现:在prefix下创建一个带租约的KV(revision递增),检查自己的revision是否是当前最小的,如果是则获锁成功。否则Watch前一个revision对应KV的删除事件,被通知后再次检查。

Worker Pool与任务执行引擎

Worker Pool使用buffered channel实现任务队列,固定数量的goroutine从channel消费任务执行。

type Task struct {
    ID       string
    Name     string
    CronExpr string
    Handler  func(ctx context.Context) error
    Timeout  time.Duration
    Retry    int
}

type WorkerPool struct {
    taskCh   chan *Task
    resultCh chan *TaskResult
    workers  int
    quit     chan struct{}
}

func NewWorkerPool(workerCount, queueSize int) *WorkerPool {
    return &WorkerPool{
        taskCh:   make(chan *Task, queueSize),
        resultCh: make(chan *TaskResult, queueSize),
        workers:  workerCount,
        quit:     make(chan struct{}),
    }
}

func (wp *WorkerPool) worker(id int) {
    for {
        select {
        case task := <-wp.taskCh:
            result := wp.executeTask(task)
            wp.resultCh <- result
        case <-wp.quit:
            return
        }
    }
}

func (wp *WorkerPool) executeTask(task *Task) *TaskResult {
    start := time.Now()
    ctx, cancel := context.WithTimeout(context.Background(), task.Timeout)
    defer cancel()

    var lastErr error
    for attempt := 0; attempt <= task.Retry; attempt++ {
        if attempt > 0 {
            time.Sleep(time.Duration(attempt) * time.Second)
        }
        err := task.Handler(ctx)
        if err == nil {
            return &TaskResult{TaskID: task.ID, Success: true, Duration: time.Since(start)}
        }
        lastErr = err
        if ctx.Err() == context.DeadlineExceeded {
            break
        }
    }
    return &TaskResult{TaskID: task.ID, Success: false, Error: lastErr, Duration: time.Since(start)}
}

Cron调度器与分布式锁集成

调度器使用robfig/cron库解析cron表达式,每次触发时先竞争分布式锁,获锁成功后提交到WorkerPool执行。

type Scheduler struct {
    cron  *cron.Cron
    lock  *DistributedLock
    pool  *WorkerPool
    tasks map[string]*Task
}

func (s *Scheduler) RegisterTask(task *Task) error {
    lockKey := fmt.Sprintf("/scheduler/lock/%s", task.ID)
    s.tasks[task.ID] = task

    _, err := s.cron.AddFunc(task.CronExpr, func() {
        s.tryExecute(task, lockKey)
    })
    if err != nil {
        return fmt.Errorf("register cron failed: %w", err)
    }
    return nil
}

func (s *Scheduler) tryExecute(task *Task, lockKey string) {
    ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
    defer cancel()

    unlock, err := s.lock.TryLock(ctx, lockKey)
    if err != nil || unlock == nil {
        return
    }

    if !s.pool.Submit(task) {
        unlock()
        return
    }

    go func() {
        time.Sleep(task.Timeout + 5*time.Second)
        unlock()
    }()
}

节点故障检测与任务接管

当持锁节点宕机,etcd的Lease在TTL到期后自动删除相关KV,分布式锁随之释放。下一次cron触发时,存活节点竞争获锁,任务由新节点接管执行。这个过程依靠etcd的租约机制自动完成。

节点存活检测通过etcd的KeepAlive实现。每个节点启动时在/scheduler/nodes/下创建带租约的KV,定期发送KeepAlive保活。

func (s *Scheduler) registerNode(nodeID string) error {
    key := fmt.Sprintf("/scheduler/nodes/%s", nodeID)

    resp, err := s.lock.client.Grant(context.Background(), 5)
    if err != nil {
        return err
    }
    leaseID := resp.ID

    _, err = s.lock.client.Put(context.Background(), key, nodeID,
        clientv3.WithLease(leaseID))
    if err != nil {
        return err
    }

    ch, err := s.lock.client.KeepAlive(context.Background(), leaseID)
    if err != nil {
        return err
    }

    go func() {
        for {
            select {
            case ka := <-ch:
                if ka == nil {
                    log.Printf("[Node] keepalive expired, node %s offline", nodeID)
                    return
                }
            case <-s.quit:
                s.lock.client.Revoke(context.Background(), leaseID)
                return
            }
        }
    }()
    return nil
}

TTL设为5秒意味着节点宕机后最多5秒锁释放,最迟下一次cron触发时任务被其他节点接管。TTL不宜过短(网络抖动可能导致误判),也不宜过长(故障接管延迟过大)。5-10秒是常见取值。

任务分片与负载分担

对于高吞吐场景(如百万级定时任务),单靠分布式锁竞争会导致etcd压力骤增。任务分片将任务按hash分配到不同节点,每个节点只负责自己分片内的任务,无需锁竞争。

func shouldExecute(taskID string, nodePartition, totalPartitions int) bool {
    h := fnv.New32a()
    h.Write([]byte(taskID))
    return int(h.Sum32())%totalPartitions == nodePartition
}

节点列表维护在etcd中,当节点加入或离开时,通过Watch感知变化重新计算分片。Consistent Hashing可避免节点数变化时大批量任务重新分片,仅影响相邻分片的任务。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/go-yu-yan-fen-bu-shi-ren-wu-diao-du-xi-tong-she-ji-yu-etcd/

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

相关推荐