第一次见到“ax调度”这四个字的时候我没把它当成什么了不起的概念直到线上任务开始随意丢失、重复执行、高峰时段谁也说不清哪台机器在跑哪个任务才明白调度这件事远不是“写个定时器”那么简单。这篇文章记录的是我在实际项目中开发和运营一套代号为ax的分布式调度系统的完整过程包括核心模型、选型理由、代码实现、生产排障和一年多的运维经验。如果你是后端工程师正在从单机cron往分布式任务调度迁移或者对“任务编排”“延迟队列”“分布式锁”这几个词只有模糊认识这篇文章应该能给你一套可以直接复用的思路而不是零散的组件介绍。1. 一个调度需求的失控先想清楚“调度”到底在调什么1.1 从一条crontab开始的膨胀之路最初我们系统里的定时任务少到可以手写在一个文档里每天凌晨做一次数据汇总每十分钟同步一次缓存超时未支付的订单每五分钟扫描一遍。每个人都能说出每个任务跑在哪台机器上维护成本低到可以忽略。真正的转折点是业务量上来之后的三件事同时发生任务数量突破了三百个服务器从一台变成十几台团队从一个人变成多个人一起往里面加需求。而当时实现定时任务的手段还停留在crontab。每增加一个任务就要在一台服务器上斟酌着改一行配置文件。时间一长混乱的苗头就全冒出来了有人在三台机器上配了同一个任务导致数据重复写入有人把任务的执行机下线了任务就消失得无声无息还有人直接拿测试环境执行生产任务最后不得不人工对账。那段时间我们花在“找任务在哪跑”上的时间比写业务代码的时间还多。1.2 单机cron的三个明显天花板第一是时间不准。cron依赖的是执行节点上的本地时钟NTP同步不到位或者节点负载飙高实际执行时间就会漂移。有些对时间点有硬性要求的任务比如向外部渠道发送整点通知只要漂移上几分钟结果就完全不可用。这个问题的根源不是cron本身而是计时的基准放在了一个不可靠的单点上。第二个是单点故障。任务放在哪台机器上哪台机器就是它的命门。机器一旦宕机任务不会再有人接管恢复之后还得有人手工触发一次“补偿执行”。线上的问题绝大多数是“该跑的没跑”这类静默故障没有报错、没有告警往往要等到下游反馈数据不对才能发现。第三是重试和幂等几乎只能靠外围脚本硬撑。cron本身只负责到点了执行一条命令至于命令是否成功、失败后要不要重跑、重跑会不会产生重复数据全部需要自己手工解决。我见过生产环境里用“再执行一次”来应付失败任务的脚本结果就是批量重复通知客户那是一种完全失控的状态。1.3 我理解的“调度”绕不开的三个要素在动手设计ax调度之前我先强迫自己和团队回答一个问题调度系统到底在“调”什么我的答案是三个词的组合时机、对象、可靠性。时机解决“什么时候触发”。触发方式不只有cron表达式还有延迟触发订单下单后五分钟后关闭、依赖触发上游任务成功之后下游才启动、文件到达触发数据文件落盘后开始处理。如果调度器只支持“给个cron表达式”那它连第一关都没过。对象解决“触发之后执行什么”。任务不是一条孤立的命令它应该包含执行入口、参数体、超时时间、重试上限这些元数据。只有把任务配置化才可能被统一管理、统一监控、统一排障。可靠性解决“执行结果如何闭环”。任务调度出去不等于任务完成必须有状态流转、失败上报、重试机制、幂等保护。没有闭环的调度系统本质上是个带提醒功能的定时器谈不上调度。1.4 何时需要自研调度中间件判断标准可以压缩成三条任务分布在两台以上机器并且需要统一管理任务之间存在顺序依赖或失败重跑需求你需要对一批任务做“全局可控”的操作比如暂停、重跑、批量变更。如果任务量只有几十个节点只有一个真的不需要上调度平台cron加简单的锁文件完全够用。ax不是银弹它解决的是“任务规模和复杂依赖已经超过cron承载力”的那段中间地带。我们当时的情况正好卡在这个阶段说大不大说小不小商用调度平台太重自己写又怕写不出所以然。现在回头看只要抓住时机、对象、可靠性这三个点自研的边界其实比想象中清晰。2. 三个角色、一张ZSet、两层防重复Ax调度的核心模型与选型复盘2.1 三个角色一条主链路ax调度内部的核心模型不复杂只有三层触发器、调度器、执行器。触发器Trigger任务的定义层。每条任务有全局唯一ID包含触发类型、执行参数、超时和重试规则。调度器Scheduler决策层。不断扫描触发时间表发现到期任务后生成一个执行实例写入待执行队列。执行器Executor行动层。从队列里领取实例执行对应Handler上报成功或失败。三层拆开以后最大的收益是调度器和执行器可以独立扩缩容。调度器处理的是“什么时候该跑”的决策执行器处理的是“任务到底怎么跑”的负载。两者之间用一张待执行实例表做缓冲调度器跑得再快也不会直接把压力倒灌给业务执行器再忙也不会影响下一次触发扫描。2.2 为什么第一版用纯DB轮询后来换成了ZSet第一版ax是用数据库轮询的。表结构大概就是task表和task_instance表调度器每过十秒执行一次“select * from task where next_run_time now”。这个方法在任务量小于一千的时候完全够用简单、直观、可排查。但我很快发现一个耐人寻味的问题任务量上来之后为了减少数据库压力轮询周期不敢设太短可周期一长任务的实时性又变差了。换句话说纯DB轮询把“扫描成本”和“调度精度”绑死在一起想提升一个必然牺牲另一个。后来把触发时间表迁到Redis的ZSet效果立刻改善。score存的是下一次触发时间戳调度器只需要用ZRANGEBYSCORE取出所有早于当前时间戳的任务ID。一次操作拿到的就是“眼下该触发的所有任务”完全不需要全表扫描。延迟任务的实现也顺理成章往ZSet里扔一个“当前时间加延迟秒数”的score到期自然被扫到。2.3 为什么没有引入重型消息队列一个很容易犯的错误是一谈到“分布式”就把消息中间件搬出来。我在ax里刻意没有引入Kafka这类重组件原因很实际任务调度不是高吞吐消息流而是低频、高可靠的状态流转。几百上千个任务的规模要求的是“每条任务被认真执行、失败要可见”而不是“每秒承载几十万条消息”。对这个场景来说一张持久化的执行实例状态表配合数据库唯一索引就足以提供可靠性和可排查性引入队列反而会带来消息积压、分区再均衡、消费位点管理这些完全与核心问题无关的复杂度。2.4 分布式锁加唯一索引的双保险调度器为了保证可用性本身也会水平部署。两个节点同时扫到同一个到期任务是绕不开的并发问题。ax在这里做了两层防护。第一层调度器触发之前先对任务ID执行一次分布式锁。我用的是Redis的SET NX EX拿到锁的节点才允许生成执行实例拿不到的节点直接跳过。第二层执行实例表里对“任务ID加触发窗口”建唯一索引。就算分布式锁因为网络抖动失效两个节点同时插入也只有一条能插入成功另一条会报唯一键冲突被当作重复触发丢弃。这两层防护对应的其实是两个层面的问题第一层解决“调度端不要重复触发”第二层解决“存储端不要重复落库”。实测下来跨节点重复触发率降到了零后面的压测数据里也有体现。3. 最小闭环的代码之旅用Go把Ax调度骨架跑起来3.1 先定数据结构再写业务第一步我把任务定义固化成了Trigger结构体。之所以先写结构体是为了让后续所有功能都围绕同一份元数据展开而不是散落在各路代码里。const ( TriggerCron cron TriggerDelay delay ) type Trigger struct { ID string json:id Type string json:type CronExp string json:cron,omitempty DelaySec int json:delay_sec,omitempty Handler string json:handler Payload map[string]string json:payload MaxRetry int json:max_retry TimeoutSec int json:timeout_sec }Trigger里的字段没有一个是多余的。ID用于幂等Type决定触发方式Handler是执行器注册表中的键Payload用来携带业务参数。MaxRetry和TimeoutSec这两个字段特别容易被忽略但它们属于可靠性维度不是可有可无的配置。3.2 调度器一次ZRangeByScore完成到期扫描调度器的主循环用Go表达反而比用文字描述更直观先取到期任务再抢占再派发最后计算下一次执行时间并回写ZSet。func (s *Scheduler) poll(ctx context.Context) error { now : time.Now().Unix() taskIDs, err : s.queue.RangeByScore(ctx, ax:schedule, 0, now, 0, 200) if err ! nil { return err } for _, taskID : range taskIDs { ok, err : s.lock.TryLock(ctx, ax:lock:taskID, 5*time.Second) if !ok || err ! nil { continue } if err : s.dispatch(ctx, taskID); err ! nil { log.Error(dispatch failed, taskID, taskID, err, err) } nxt : calculateNextRunTime(taskID, now) s.queue.Add(ctx, ax:schedule, taskID, nxt) } return nil }这里有几个细节值得注意。锁定时间设成5秒而不是更长是因为任务创建执行实例的操作必须足够快如果5秒还没完成说明调度器本身出了问题主动放弃反而比卡死更好。dispatch失败的时候我没有把任务从ZSet里移除这样任务会在下一个扫描周期被重新触发利用的是“到期即重试”的语义。3.3 执行器状态机由数据库驱动执行器的消费逻辑反而简单从pending实例表中获取一个实例执行Handler根据结果更新状态。我把核心代码保持在二三十行剩下的复杂性全部下沉到状态机里。func (w *Worker) consume(ctx context.Context) { for { inst, err : w.store.AcquirePendingInstance(ctx) if err ! nil { continue } runErr : w.handlers[inst.Handler](ctx, inst.Payload) if runErr ! nil { if inst.Retry inst.MaxRetry { w.store.MarkRetry(inst.ID, inst.Retry1, backoffTime(inst.Retry)) } else { w.store.MarkFailed(inst.ID, runErr.Error()) } continue } w.store.MarkSuccess(inst.ID) } }这里并没有一个复杂的“任务队列”而是通过数据库的pending状态加一条“只取一条可执行实例”的SQL来达到队列效果。这个方案的好处是任务状态全部可查询、可回溯执行器崩溃后pending中的实例会被重新认领不会凭空消失。AcquirePendingInstance会先执行一条带FOR UPDATE SKIP LOCKED的查询把某条pending中的实例锁定为自己专属其他执行器不会碰到它。MarkSuccess、MarkFailed、MarkRetry都是在事务里更新实例状态。backoffTime是重试退避策略我习惯用“2的n次方秒”作为基础最大不超过五分钟。3.4 跑通闭环后我删掉了两样“加戏”的组件第一次跑通时我堆了不少东西一个自制的任务队列、一个复杂的分片协调器、一套事件通知机制。结果跑通了但这些“加戏”的组件很快就成了调试负担。删掉队列是因为实例表本身的pending状态已经足够当队列用删掉分片协调器是因为调度器水平扩展时锁加唯一索引已经解决了重复触发删掉事件通知是因为任务实例的每一次状态变更都落在数据库里没有事件也能随时查得到。注意这里说的“删组件”不是简单的懒人做法而是把不带来业务价值的抽象层拿掉。调度系统的复杂度天然高能少一层概念就少一层排查问题的路径会短很多。这套“先做最小闭环再删多余组件”的方法是我整个ax项目里最值钱的一步。它让我明白调度系统的核心从来不是组件的多少而是那条从“触发”到“执行结果”的状态链路是否完整。4. 崩溃、重试与幂等生产环境里真正吃时间的三个细节4.1 调度器和执行器崩溃了怎么办调度器是无状态的它只负责不断扫描ZSet所以重启非常简单启动后直接从当前时间开始扫描即可。真正麻烦的是执行器崩溃。执行器在拿到一个实例后会将实例标记为running。此时执行器崩了这个实例就永远停留在running状态不做处理任务就静默丢失了。所以ax给执行器加了心跳租约机制执行器启动后注册自己的心跳每隔10秒上报一次心跳超过30秒没有上报调度器就认为该执行器已经失联会把它名下所有running中的实例重新置回pending等待其他执行器认领。一个实例在running状态停留超过超时时间同样会被回收重派。这个机制换来的是“任务可能延迟完成但不会无声消失”。代价是任务可能会有重复执行因为执行器可能在最后一步完成任务、还没来得及上报成功时心跳就断了。这个代价没办法完全消除但可以通过幂等把影响降到最低。4.2 重试策略重试是有成本的补偿措施重试不是简单地在失败后“再来一次”。我在ax里给每次重试计算了退避时间序列大致是第0次重试30秒后 第1次重试1分钟后 第2次重试2分钟后 第3次重试4分钟后 第4次重试8分钟后这个退避序列的目的是避免故障时所有任务像潮水一样同时重试把系统直接打崩。重试次数达到上限后任务进入failed状态并触发告警。一个经常被忽略的细节是重试应该只在“执行结果不确定”时才完整重跑如果任务在执行前就明确失败比如参数校验不通过重试毫无意义直接标记失败即可。所以MaxRetry这个配置要按任务的失败模型来设定不能一刀切。4.3 幂等调度系统唯一主动加的“保护层”重试机制天然会带来重复执行。为了不让重复执行产生数据问题ax在三个层面加了幂等保护而不是只依赖业务方自己处理。第一层是“触发幂等”。前面的分布式锁和唯一索引保证了同一个触发时间窗口只会生成一个执行实例。第二层是“执行幂等”。执行器在真正执行任务前会先查询目标系统当前状态。比如关单任务先查订单状态如果已经是已关闭就直接返回成功推送任务先查询接收方是否存在本次推送的记录存在就跳过。第三层是“记录幂等”。数据库里对业务单据号建立唯一键重复插入直接冲突不产生脏数据。这一层往往是最稳妥的兜底因为前两层都可能在引入新问题时失效唯一索引是数据层最后一道防线。我个人最推荐的幂等落地顺序是先保证任务ID不重复再做业务层幂等判断最后用数据库唯一索引兜底。三层叠加之后重复执行在绝大多数场景下的副作用可以忽略不计。5. 压测与监控没有可观测性的调度系统是盲的5.1 压测数据ax到底扛得住多少任务我做了一次相对完整的压测任务规模设为1000个调度周期覆盖每分钟、每五分钟、每小时执行器并发数从5调到30。统计口径不算严格但足以说明方向性问题。场景任务数调度器节点执行器并发调度成功率平均触发延迟基础场景10001599.9%100ms横向扩展1000310100%80ms高并发执行100033099.8%120ms压测中出现了极低比例的非100%成功率查下来根本原因只有一个某个执行器节点内存被GC卡住心跳超时触发了超时重派。其余场景下重复触发已经归零整体处在可接受范围。5.2 监控指标先定好阈值再上线干活一个调度系统在上线之前至少要盯住这几个指标调度延迟任务到期时间与实际生成执行实例时间的差值。执行时长同一Handler执行耗时的P95以及P99。pending积压数当前待执行实例的数量。重试率重试次数与执行总次数的比例。失联执行器数心跳超时执行器的数量。这些指标全部以计数器或直方图形式暴露。上线第一周我就发现一个规律pending积压数并不是在任务量最多时最高而是在某台执行器因磁盘问题失联时才暴涨。如果没有监控这个问题可能要等业务反馈“任务跑得很慢”才暴露最后的排查成本会大得多。5.3 告警的分级与收敛告警一定要分级且要收敛。ax的告警分成三级警告级、严重级、致命级。警告级是可以容忍但需要留意的比如单节点心跳慢严重级需要人工介入比如重试率连续五分钟超过5%致命级需要立即响应比如pending积压数持续上升且执行器全部失联。如果所有异常都用一个告警级别运维大概率会陷入“狼来了”的麻木状态最终在真正的故障面前视而不见。收敛的做法是把同一任务ID同一天的重复告警合并成一条聚合消息让值班人能快速定位而不是淹没在告警流里。6. 一年运行下来我总结的调度系统设计教训6.1 宁可重复执行也不要静默丢失这是整个ax项目里我最想分享的一条原则。分布式环境下任务执行做到“至少一次”比“精确一次”更现实。精确一次意味着要么引入事务消息要么引入更复杂的协议对大多数业务而言成本远高于收益。宁可让任务偶尔重复也要保证任务一定被执行重复的副作用用幂等来解决丢失的副作用却是任何日志都无法追回的。一条铁律调度系统可以重复但不能丢。重复可以用幂等补偿丢了就只能祈祷下游有对账机制。6.2 状态机越简单越好调度系统的核心是一张实例状态表pending到running到success或者failed外加一个retry分支。我给ax加过更复杂的状态比如suspended、canceled、timeout后来发现大多数场景下一个pending被重新认领就已经能表达超时重派了一个手动暂停也无非是把任务的下次触发时间无限推后。状态越多恢复逻辑就越复杂排查也越慢。简单状态机反而是生产稳定性最好的状态机。6.3 重试配置要按任务模型定制MaxRetry和TimeoutSec这两个配置不同任务差异很大。有的任务是对账任务重试多跑几次没有坏处有的任务是外部通知任务重复通知会造成客诉重试就得谨慎。统一配置会导致要么重试不足、要么重试过度。正确做法是在任务注册的时候就把这两个参数打通让每个任务自己决定自己的可靠性策略。6.4 调度精度和扫描成本是一对宿敌从DB轮询到ZSet本质上是把“扫描”从全表同步变成了索引范围查询。但即使如此调度器也不能无限制地提高扫描频率。过高的扫描频率会让Redis的QPS持续打满最后影响触发链路上其他组件。根据我的经验调度扫描间隔设置在100到200毫秒之间实时性和开销比较平衡。6.5 调度系统是为业务服务的不是给自己加戏的做ax这一年里我反复问团队一个问题这个组件是业务真实需要的还是只是为了技术上的合理性凡是回答不清楚的组件最后都被删或者被重构了。调度的价值在于让业务按时、可靠地发生不在于技术栈的复杂度高低。时刻记住这一点设计才不会跑偏。这些经验都是运行了一年多以后花了真金白银踩出来的。如果你也在规划自己的“ax调度”我的建议是先别急着定技术栈把触发类型、任务状态、幂等和重试四条规则写在纸上跑一个最小闭环再决定要不要加分布式锁、要不要引入消息队列。技术选型永远只是水到渠成真正难的是想清楚调度系统到底要为业务承担什么职责。还有一点容易被忽略的是调度系统一旦上线运维体验和使用体验同样重要。给任务起一个清晰的名字、在同一条链路上保持一致的日志格式、为每一次失败留下可查询的状态记录这些看似琐碎的小事会在某个凌晨三点成为支撑你排查故障的关键。愿你的任务都能按时跑完失败都能及时告警。
