分布式任务调度系统需要在多节点环境下保证任务不被重复执行、故障节点任务能被自动接管。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/