Go分布式任务调度项目复盘:从Cron到任务依赖DAG的架构演进
2026/7/23 7:42:59 网站建设 项目流程

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

在单机上工作正常。但问题随之而来:

  1. 任务之间隐式依赖——aggregate必须在clean完成后执行。clean执行超时不会通知aggregate等待。
  2. 某次clean跑了2.5小时(正常1小时),aggregate在4:00准时启动——读到了未清洗完的数据。
  3. 无法水平扩展——所有任务都在一台机器上,数据量增长后单机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 }

任务执行的分布式Worker:

type 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, attempt+1, 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任务是否需要执行?方案:提供两种策略——FailFast(B失败,整个DAG停止)和BestEffort(B失败,忽略B继续执行不依赖B的任务)。

边界三:Redis分布式锁的死锁预防。如果Worker在执行任务时crash,锁会保持在Redis中。方案:锁设置5分钟TTL,Worker通过心跳持续续约。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个之前不存在但实际存在的依赖关系。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询