Go分布式任务调度项目复盘:从Cron到任务依赖DAG的架构演进

发布时间:2026/7/23 7:43:04
Go分布式任务调度项目复盘:从Cron到任务依赖DAG的架构演进 Go分布式任务调度项目复盘从Cron到任务依赖DAG的架构演进一、Cron不够用了一个数据ETL系统最初用Linux Cron管理定时任务0 2 * * * /app/bin/import_data.sh 0 3 * * * /app/bin/clean_data.sh 0 4 * * * /app/bin/aggregate.sh 0 5 * * * /app/bin/export_report.sh在单机上工作正常。但问题随之而来任务之间隐式依赖——aggregate必须在clean完成后执行。clean执行超时不会通知aggregate等待。某次clean跑了2.5小时正常1小时aggregate在4:00准时启动——读到了未清洗完的数据。无法水平扩展——所有任务都在一台机器上数据量增长后单机CPU成了瓶颈。需要从Cron演进到有任务依赖分布式执行能力的调度系统。二、DAG调度引擎的设计核心数据结构// Task 任务定义 type Task struct { ID string Name string Command string Timeout time.Duration RetryPolicy RetryPolicy DependsOn []string // 依赖的上游任务ID } // DAG 有向无环图 type DAG struct { mu sync.RWMutex tasks map[string]*TaskNode } type TaskNode struct { Task *Task Status TaskStatus // pending/running/success/failed StartTime time.Time EndTime time.Time Parents []*TaskNode Children []*TaskNode } type RetryPolicy struct { MaxRetries int Backoff time.Duration } // 拓扑排序 并行执行 func (d *DAG) Execute(ctx context.Context) error { ready : d.findReadyTasks() // 所有父任务都完成的节点 var wg sync.WaitGroup errCh : make(chan error, len(ready)) for _, node : range ready { wg.Add(1) go func(n *TaskNode) { defer wg.Done() if err : d.executeTask(ctx, n); err ! nil { errCh - err } }(node) } wg.Wait() close(errCh) // 检查错误 for err : range errCh { if err ! nil { return err } } // 递归执行下游任务 return d.Execute(ctx) // 执行已解除依赖的子任务 } func (d *DAG) findReadyTasks() []*TaskNode { d.mu.RLock() defer d.mu.RUnlock() var ready []*TaskNode for _, node : range d.tasks { if node.Status StatusPending { allParentsDone : true for _, parent : range node.Parents { if parent.Status ! StatusSuccess { allParentsDone false break } } if allParentsDone { ready append(ready, node) } } } return ready }任务执行的分布式Workertype Worker struct { id string taskCh chan *Task executor Executor registry *Registry // 分布式注册中心 } func (w *Worker) Run(ctx context.Context) { for { select { case task : -w.taskCh: w.executeWithRetry(ctx, task) case -ctx.Done(): return } } } func (w *Worker) executeWithRetry(ctx context.Context, task *Task) { for attempt : 0; attempt task.RetryPolicy.MaxRetries; attempt { tCtx, cancel : context.WithTimeout(ctx, task.Timeout) result : w.executor.Execute(tCtx, task.Command) cancel() if result.Error nil { return // 成功 } log.Printf(任务 %s 第%d次失败: %v, task.ID, attempt1, result.Error) if attempt task.RetryPolicy.MaxRetries { time.Sleep(task.RetryPolicy.Backoff) } } // 所有重试都失败 w.markTaskFailed(task) } // 分布式互斥——同一任务只在一个Worker上执行 func (s *Scheduler) acquireTaskLock(taskID string) (bool, error) { // Redis SET NX TTL 实现分布式锁 result, err : s.redis.SetNX(ctx, fmt.Sprintf(task_lock:%s, taskID), s.workerID, 5*time.Minute, // 锁的TTL防止死锁 ).Result() if err ! nil { return false, err } return result, nil }三、与Cron的实际改进对比维度Cron方案DAG方案任务依赖隐式靠sleep时差显式DAG并行度单机3 Worker分布式超时处理无任务卡死不感知Context超时重试任务可观测性无日志/指标/状态追踪执行总时间6.5h2.3h并行优化关键改进原来的4步流水线导入→清洗→聚合→导出有3小时的依赖等待。改为DAG后导入完成后可以并行执行清洗和日志分析进一步并行化下游任务总时间减少65%。四、实现中的关键边界处理边界一DAG的循环依赖检测。在添加任务时检查是否会形成环func (d *DAG) hasCycle(from, to string) bool { visited : make(map[string]bool) var dfs func(node string) bool dfs func(node string) bool { if node from { return true } visited[node] true for _, child : range d.tasks[node].Children { if !visited[child.Task.ID] dfs(child.Task.ID) { return true } } return false } return dfs(to) }边界二任务超时下的下游处理。当B任务超时失败后依赖B的D、E任务是否需要执行方案提供两种策略——FailFastB失败整个DAG停止和BestEffortB失败忽略B继续执行不依赖B的任务。边界三Redis分布式锁的死锁预防。如果Worker在执行任务时crash锁会保持在Redis中。方案锁设置5分钟TTLWorker通过心跳持续续约。crash后锁在5分钟内自动释放。五、总结从Cron到DAG调度系统的核心经验隐式依赖靠sleep时差是定时任务Bug的最大来源——必须显式化拓扑排序DAG是表达任务依赖的简单有效方案分布式锁Redis SET NX解决了单点执行问题但要处理crash后的锁释放任务超时重试策略是生产环境中最重要的容错能力DAG让并行执行成为可能——执行总时间减少65%是最大的实际收益当前DAG调度系统管理42个定时任务3个Worker节点。如果任务数增长到100需要考虑的扩展方向按业务域分组管理DAG一个DAG不超过20个节点降低复杂度以及引入任务优先级队列核心业务任务优先于报表任务。从Cron迁移到DAG的最大阻力不是技术实现而是把隐式依赖写清楚这一步。以前的sleep 3600被替换为depends_on: [clean_data, log_analysis]——这个显式化的过程暴露了3个之前不存在但实际存在的依赖关系。