1. 从“任务炸了”到动手写 ax:一个调度内核的诞生背景
先说个真实场景。去年年中,我们整个微服务集群的定时任务和延迟任务开始失控:订单超时自动关闭、积分过期提醒、异步对账补偿、消息重试投递,这些逻辑散落在十几个服务里,各自用time.Sleep、Redis 过期回调、数据库轮询硬扛。高峰期一到,任务一多,日志就开始出现“任务积压”“执行超时”“重复触发”。更头疼的是,每次排查都要翻遍不同服务确认“到底哪个模块在管这个任务”,改一个延迟策略还得发版重启。
我翻遍了市面上所有调度产品:Quartz 太重,架构老,集群方案繁琐;xxl-job 调度定期任务很顺手,但延迟事件的支撑偏弱;Airflow 定位批处理 DAG,天然不适配高吞吐事件流。至于直接用消息队列做延迟队列,确实能解决一部分问题,可要做到任务编排、重试、优先级、分布式一致性,还是得自己补一大圈代码。
抱着“与其到处打补丁,不如做一个统一调度内核”的想法,我花了几个月时间,从零写了一套命名为 ax 的调度系统。它不是什么颠覆性创新,也没有炫酷的界面,核心目标就一个:把散落在业务代码里的任务调度逻辑收拢到一个内核里,业务方像调用本地函数一样提交任务,调度、分派、重试、幂等、监控全都交给 ax 处理。这篇就聊聊 ax 的整体设计、调度核心、一致性方案,以及压测过程中踩到的几个隐藏深坑。如果你也在为任务散落、触发不准时、积压追查困难发愁,这篇应该能给你一些可以落地的思路。
2. 为什么没直接选现成框架:ax 的定位与整体架构
2.1 现成框架的“最后一个短板”在哪
先说清楚,不是现成框架不能用,而是它们解决的核心问题跟 ax 不太一样。Quartz 是单机调度,靠数据库锁做集群,节点多了锁竞争会很酸爽;xxl-job 的强项是 cron 定期任务,对“延迟多久触发一次”这种海量事件型任务,需要把每个事件都映射成一条调度记录,存储和扫描开销都不小;Airflow 这类工作流引擎擅长的是有向无环图的批处理编排,一个任务实例的调度周期通常秒级以上,不适合我们这种秒级以内要响应几千次的场景。
ax 想补的短板,一句话概括就是“事件驱动的轻量级延迟调度”。业务方不关心你的任务是被定时器触发还是被消息触发,不关心它被派发到哪台机器,只关心“我提交了一个任务,希望在 3 秒后执行,如果失败了自动重试,最多三次”。为了支撑这个模型,我一开始就定下了三层架构:ax-core 负责接收任务并维护调度队列,ax-broker 负责把就绪任务派发给 worker,ax-worker 负责真正执行并回传状态。三层职责单一,每一层都可以独立扩展。
2.2 核心组件划分与各层职责
下面这张表是我在项目初期定下的组件边界,后续所有功能都围绕这几条线展开:
| 组件 | 核心职责 | 关键接口 |
|---|---|---|
| ax-core | 接收任务提交、维护延迟队列与优先级分桶、触发就绪任务 | SubmitTask(task)、CancelTask(id)、AckLease() |
| ax-broker | 拉取就绪任务、按 worker 负载做派发、背压控制 | PullReady()、Dispatch(task, worker)、ReportBacklog() |
| ax-worker | 接收任务并执行、上报心跳与执行结果 | Execute(task)、Heartbeat()、Complete(result) |
ax-core 不做任务执行,它只关心“时间到没到”。ax-broker 不做任务存储,它只关心“任务该发给谁”。ax-worker 不关心全局状态,它只关心“当前这一条任务怎么跑完”。三层通过内部 RPC 通信,任务状态统一落到一个存储里,方便对账和排障。
从业务方的视角看,接入 ax 的体验非常轻。比如在 Go 服务里,只需要引入 ax-client,然后client.Submit(ctx, &Task{Payload: orderID, Delay: 3 * time.Second, Topic: "order.close"}),剩下的等待、重试、超时都由 ax 管理。这比每个服务自己塞一个延迟队列要清爽得多,也让“任务治理”第一次有了统一的抓手。
3. 任务模型设计:从“定时触发”升级成“状态机驱动”
3.1 任务数据结构与状态流转
任务模型是整个调度内核的地基,地基不稳,后面全白搭。我在设计任务数据结构时,没有走“定时任务表 + 扫描线程”的老路,而是把每条任务当成一个有限状态机。任务有两个时间字段:ready_time(最早可执行时间)和deadline(最晚必须开始执行的时间)。ready_time由提交参数里的 Delay 决定,deadline则用来兜底,防止任务在队列里待太久导致业务上过期失效。
任务状态流转我用 JSON 描述一下:
{ "id": "task_8f3a2c", "topic": "order.close", "payload": "{\"orderId\":\"20250101\"}", "status": "pending", "priority": 5, "ready_time": 1735689600, "deadline": 1735689900, "retry_count": 0, "max_retry": 3, "version": 1 }状态机一共六个状态:pending(已提交未到期)、ready(时间已到待派发)、running(已派发给 worker 正在执行)、succeeded(执行成功)、failed(执行失败,等待重试)、canceled(手动取消)。流转规则只有五种:pending -> ready由时间轮驱动;ready -> running由 broker 派发驱动;running -> succeeded / failed由 worker 上报结果驱动;failed -> pending表示进入重试,重新计算 ready_time;任何状态下都可以被用户主动取消进入canceled。
这里有一个容易犯迷糊的点:为什么ready和running不合并成一个状态?因为任务从“时间到了”到“真正被 worker 拉走”之间是有延迟的,如果合并,就说不清楚任务到底卡在调度环节还是执行环节。拆成两个状态后,监控指标就非常清晰:pending -> ready的耗时反映触发准时度,ready -> running的耗时反映派发效率,running -> succeeded的耗时反映执行质量。
3.2 乐观锁与状态更新
任务状态并发更新是调度系统最容易出 bug 的地方。一开始我图省事,用数据库行锁,结果任务一多,锁等待直接把数据库拖垮。后来改成乐观锁:每次更新都带version字段,UPDATE task SET status = 'ready', version = version + 1 WHERE id = ? AND version = ?,更新影响行数为 0 就说明冲突,重新读取再判断。
举个例子,一个任务正在被 worker 执行,这时候用户发来取消请求。worker 执行完上报succeeded,取消请求想把状态改成canceled,两者同时操作,乐观锁保证只有一个能成功。如果取消请求先成功,worker 上报时 version 对不上,会被拒绝,然后 worker 侧做“执行结果已被忽略”的处理。这个机制看似简单,但它保证了状态流转不会出现“既成功又取消”的脏状态。
优先级是另一个绕不开的设计点。我把优先级分成 0 到 9 十个档位,9 最高。调度时不是简单的“高优先级永远先跑”,那样会让低优先级任务饿死。ax 采用带权轮转:每个优先级桶里的任务都有权重,调度机会按权重比例分配。比如优先级 9 和优先级 5 的任务权重比是 3:1,那么高优先级桶最多连续弹出 3 条,就得让低优先级桶弹 1 条。这个设计在业务上很实用,比如“订单超时关闭”这类时效敏感任务给优先级 9,而“每日统计报表”这类容忍延迟的任务给优先级 2,两边都能在合理时间内被执行。
4. 调度核心原理:时间轮、双缓冲队列与精确触发
4.1 时间轮为什么比“扫表”快了不止一个量级
很多团队做延迟任务的第一反应是用数据库轮询:SELECT * FROM task WHERE ready_time <= NOW(),每秒扫一次。数据量小的时候没问题,几十万条延迟任务之后,扫描开销和索引命中率都会变差,更别说同一秒内大量任务集中触发时的毛刺。
ax 用的是分层时间轮(Hierarchical Timing Wheel)来管理延迟任务。时间轮的基本思想是:把时间分成一个个 tick 槽位,每个槽位挂一个任务链表,指针每 tick 移动一格,指向的槽位里所有任务就是当前到期的任务。单层时间轮的问题是精度和槽位数彼此制约——精度到毫秒、范围到一天,需要的槽位数是 86,400,000,内存根本扛不住。分层时间轮则用多轮嵌套解决:第一层每 tick 1 毫秒,共 1024 个槽位,代表约 1 秒;第二层每 tick 1 秒,共 1024 个槽位,代表约 17 分钟;第三层每 tick 17 分钟,容纳更长时间范围的任务。任务插入时先放进最底层能容纳它的那层,随着时间推移逐层降级,最终落到第一层触发。
这里贴一段简化后的代码逻辑,帮助理解:
class TimingWheel: def __init__(self, tick_ms=1, slots=1024): self.tick_ms = tick_ms self.slots = [deque() for _ in range(slots)] self.current = 0 def add(self, task, delay_ms): # 计算应该放入哪个槽位 ticks = delay_ms // self.tick_ms idx = (self.current + ticks) % len(self.slots) task.remaining_ticks = ticks self.slots[idx].append(task) def advance(self): self.current = (self.current + 1) % len(self.slots) due_tasks = self.slots[self.current] self.slots[self.current] = deque() return due_tasks实际工程里 ax 的分层时间轮还包含了任务取消、层级间级联等逻辑,但核心思路就是这样。用时间轮管理触发的效果很直接:单机每秒可以稳定触发数万甚至数十万条到期任务,而 CPU 开销远低于数据库轮询。
4.2 双缓冲队列:把“入队”和“派发”的锁竞争降到最低
任务触发只是第一步,触发后如果直接把任务塞给 broker 派发,高并发下锁竞争立刻会成为瓶颈。ax 的解法是双缓冲队列:每个优先级桶里维护两个切片——incoming和ready。定时器 tick 到来时,先把本 tick 触发的任务追加到incoming,并不直接暴露给消费者。当ready切片被消费到一定程度,或者一个 tick 周期结束时,把incoming整体切到ready,同时新开一个空切片作为新的incoming。
这个设计和 Go 的sync.Pool思想有点像:写操作永远只碰incoming,读操作永远只碰ready,中间靠一次“交换”完成数据交接,而不是每个元素加锁。实测下来,单机一万 TPS 提交时,调度模块的锁等待时间几乎可以忽略不计。如果你也想优化你的任务队列,我强烈建议先想想“能不能把锁粒度变成整个队列级别的交换”,而不是每个元素都去抢锁。
4.3 精确触发与误差控制
最后聊精度。ax 承诺的调度误差是“P99 触发延迟不超过 20 毫秒”,这意味着 tick 周期不能太长。我最终把第一层时间轮 tick 设为 2 毫秒,这是因为系统时钟本身有波动,太精细的 tick 反而会因为系统调用开销引入更多误差。实测中,tick 为 2 毫秒时,P99 误差在 8 到 16 毫秒之间,完全满足大多数业务场景。
有一点要提醒:如果你要做的是“秒级精确”的金融级任务,只看 P99 是不够的,还需要关注 P999 和最大误差,并且对系统时钟做单调性保护。这个问题我在后面“压测深坑”部分会详细展开,很多调度系统被 NTP 坑过,ax 也未能幸免。
5. 一致性保障:多副本选主、幂等执行与任务对账
5.1 多副本选主:不是每个节点都在跑时间轮
调度系统的高可用很容易被人误解成“多部署几个实例就行”。如果多个 ax-core 实例同时推进时间轮,同一个延迟任务会被触发多次,然后被多个 broker 重复派发,业务端就会收到重复执行。ax 的做法是:ax-core 做多副本,但同一时刻只有 leader 节点在推进时间轮和触发任务,follower 节点只接收任务写入并同步状态,不参与调度。
选主机制我权衡过几种方案:Raft 实现起来太重,首批版本等不起;最终选了基于存储租约(lease)的简化方案。所有 ax-core 实例启动后尝试在存储里写入一条租约记录,内容包括节点 ID 和过期时间。持有租约的节点是 leader,每 5 秒续租一次。如果 leader 宕机,租约过期后其他节点竞争续租,谁先成功谁成为新 leader。这个方案比 Raft 简单很多,代价是租约切换期间存在最多 5 秒的调度空窗。对于延迟任务来说,5 秒空窗通常可以接受,因为任务只是“晚触发”,不是“不触发”。
选主切换的关键点在于:新 leader 上任后,必须重新扫描所有状态为pending的任务,按ready_time重建时间轮。这个重放过程要快,我用的优化手段是维持一张“延迟任务索引”,按 ready_time 排序,重建时间轮时只需要扫描索引头部的一小段,而不是全量扫描。
5.2 执行幂等:at-least-once 与业务去重
调度系统通常保证的是“至少一次”投递,而不是“恰好一次”。原因很简单:网络超时、worker 崩溃、broker 重试,都可能导致同一条任务被派发两次。要真正做到 exactly-once,需要业务端配合做幂等,调度系统只能尽量降低重复概率。
ax 的方案是给每次任务派发生成一个全局唯一的execution_id。worker 执行完毕后,把execution_id和结果写入结果表,结果表对这个字段建唯一索引。如果一条任务被重复派发,第二个 execution_id 对应的结果写入会失败,ax 就能识别到“这条执行是重复的”,直接丢弃,不再触发补偿。这个方案不能完全避免重复执行——两个派发如果同时进行,worker 确实会跑两遍——但至少能保证最终状态一致,不会出现“显示成功但实际失败”的混乱。
5.3 对账任务:兜底最后一公里
不管机制多完善,分布式系统里总会有漏网之鱼:broker 派发出去的消息丢了、worker 执行成功了但上报结果超时、网络分区导致状态卡在running。ax 每天凌晨跑一个对账任务,把所有长时间停留在running状态超过 10 分钟的任务捞出来,重新检查 worker 心跳,如果 worker 已经失联,就把任务状态重置为ready,重新派发。
这个对账机制被很多人忽略,但它恰恰是生产环境救命的最后一道防线。没有对账,一次网络抖动造成的“僵尸任务”就可能永久卡在running状态,不触发、不失败、不告警,业务上表现为“订单一直不关闭”,排查起来极其痛苦。我建议任何做任务系统的团队,都给自己留一个对账任务,频率不用高,但一定要有。
6. 压测实录:ax 踩过的三个隐藏深坑,以及修复过程
6.1 时钟跳跃:NTP 一键送走所有延迟任务
第一个大坑发生在压测环境切到生产环境的第二天。早上 7 点整,监控告警突然炸了:本该在未来几小时内逐步触发的延迟任务,在同一秒内几乎全部触发了。后台日志显示,大量任务的ready_time明明在未来,却被判定为到期。
排查过程很有代表性。先怀疑数据库时间字段有问题,查了半天没有异常;再怀疑代码里时间比较的逻辑写反了,review 了几遍也没有问题。直到我随手执行了date命令,才发现系统时间被 NTP 校准往前拨了整整 5 分钟。这 5 分钟的跳变,让时间轮里所有任务的“剩余时间”变成了负数,于是全部被当成到期任务触发。
修复方案基于一个常识:程序内部的“剩余时间”计算应该用单调时钟,而不是墙上时钟。单调时钟保证只增不减,不受 NTP 调整影响。我把调度模块里所有“比较当前时间和任务 ready_time”的逻辑都改成基于time.monotonic()计算剩余时间,墙上时钟只用于任务创建时设置初始延迟。同时给调度线程加了一个“校准窗口”:如果检测到墙上时钟发生超过 1 秒的跳变,就触发一次全量重建时间轮,而不是让旧的时间轮带病运行。
6.2 broker 派发太快,worker 本地队列积压导致的超时
第二个坑是在压测 5000 TPS 时暴露的。增加 worker 节点数量之后,系统吞吐量并没有线性提升,反而出现大量“任务派发成功但执行超时”。一开始怀疑是 worker 执行能力不足,但单看每个 worker 的 CPU 都在 20% 以下。
翻了两天日志才找到根因:broker 派发任务用的是 push 模式,不管 worker 能不能吃得下,派发了就先记录成功。结果就是 worker 本地队列越堆越长,最早派发的任务在队列里蹲了几分钟才被执行,而执行超时是从“派发时间”开始算的,任务还没被业务函数执行就已经超时了。典型的生产者速度远超消费者速度,中间没做流控。
修复方案是把 push 改成 push-with-credit:worker 启动后向 broker 注册,并上报自己的“可接收任务额度”,比如 30。broker 每派发一条任务,就把对应 worker 的额度减 1;worker 每执行完一条任务,就向 broker 汇报并返还额度。额度为 0 时,broker 不再给这个 worker 派发新任务。这个机制和 TCP 的滑动窗口本质上是同一个道理——消费能力决定发送速度。改造之后,worker 本地队列长度稳定在个位数,任务积压问题直接蒸发。
6.3 GC 停顿:200 毫秒的延迟尖刺从哪来
第三个坑比较隐蔽。压测持续跑了 40 分钟后,监控图上出现规律的“锯齿”:每 3 到 5 分钟,调度延迟就会出现一次 200ms 左右的尖峰。这种周期性的毛刺没有任何业务流量波动,一开始我根本没往 GC 上想,直到有一次顺手导出了 JVM GC 日志,才发现每次尖峰都对应一次 Full GC。
根因是 broker 模块在派发任务时,把每个任务的 payload 都复制了一份放进派发上下文,高峰期大量对象在新生代迅速填满,触发频繁的 Young GC 和周期性 Full GC。调度线程虽然没有休眠,但在 GC 停顿时所有线程都会暂停,所以延迟尖刺直接反映在任务触发上。
修复方案分两步:一是压测时默认启用低停顿垃圾回收器,减少 Full GC 频率;二是重构派发上下文,把“复制 payload”改成“引用传递”,并尽最大可能避免在调度热路径上分配大对象。具体来说,就是把 payload 的字节数组统一放到一个内存池里,派发时只传递引用和偏移量。优化之后,GC 停顿对调度延迟的影响从 200ms 直降到 10ms 以内。如果你也在做低延迟系统,建议把“热路径零分配”写进开发规范,而不是事后补救。
7. 调参与落地建议:从压测数据到生产参数
7.1 关键参数推荐范围
ax 部署时有一组参数直接影响调度效果,下面是我压测多轮后得出的推荐范围,你可以根据自己的场景调整:
| 参数 | 推荐值 | 备注 |
|---|---|---|
| 时间轮 tick | 2ms | 调小会增加 CPU 开销,调大牺牲精度 |
| 优先级桶数量 | 10 | 0 到 9,与业务档位一一对应 |
| broker 派发额度 | 30 | 压测中 20 到 50 是均衡区间 |
| 任务最大重试次数 | 3 | 超过后进入死信队列,人工处理 |
| 租约续租间隔 | 5s | 短了增加存储压力,长了增加空窗 |
| 对账扫描间隔 | 60s | 只捞超过 10 分钟的僵尸任务 |
这里特别说一下“任务最大重试次数”。很多团队把重试次数设成 5 甚至 10,觉得多点机会总是好的。实际上重试过多会造成“坏任务霸占 worker”:一条调用第三方接口超时的任务,重试 10 次,每次跑 2 分钟,一个 worker 的 20 分钟就被吃掉了。ax 的做法是重试 3 次后进入死信队列,由人工排查业务逻辑问题,而不是让机器盲目重试。
7.2 监控指标体系与告警阈值
调度系统上线后,第一步不是加功能,而是把监控建全。ax 在 Prometheus 上暴露的指标里,我认为最需要盯的是这五组:
ax_task_trigger_latency_ms:任务触发延迟分布,P99 超过 50ms 就要排查ax_task_execution_time_ms:任务执行耗时分布,用来发现业务逻辑变慢ax_queue_backlog:各优先级桶的积压数量,积压持续上涨说明消费能力不足ax_worker_heartbeat_stale:worker 心跳过期数量,非 0 说明有 worker 失联ax_task_retry_distribution:重试次数分布,重点看重试 3 次的占比
结合这些监控,我发现一个很有用的经验:不要只看平均值,要看 P999。很多调度问题在平均值上毫无波澜,但 P999 已经悄悄飙到几百毫秒,等真实用户感受到才去排查就晚了。把 P999 的告警阈值设成平均值的 5 到 10 倍,能提前暴露大部分隐患。
7.3 接入新业务的建议流程
最后分享一点接入流程上的经验。新业务接入 ax 时,我强烈建议先从那个业务模块的“最不重要的任务”开始,比如每日一次的数据清洗,跑两周观察运行情况,再逐步接入线上核心链路。第一次接入就上“订单超时关闭”这种高优业务,一旦出了问题,业务方对调度系统的信任感就没了,后面再推任何改造都会很费力。
接入过程中还有一个小技巧:让业务方提交任务时把payload里带上一个全局唯一的业务追踪 ID,比如订单号、用户 ID。这样任务出问题时,可以从 ax 的日志系统直接反查到业务记录,省去两个系统之间来回比对的时间。我自己踩过这个坑——早期接入的任务没有追踪 ID,线上一条任务执行失败,业务方问“是哪个订单”,我只能说“不知道”,那种感觉真的非常尴尬。
8. 写在最后的实战体会
ax 这套调度内核从设计到落地,前前后后改了无数轮,最大的收获不是学会了时间轮和选主算法,而是搞明白了一条朴素的道理:调度系统的难点从来不在调度本身,而在边界情况。时钟会跳变、worker 会失联、消息会重复、网络会抖动,这些“异常中的异常”才是系统设计真正比拼的地方。
如果让我重新设计一次,我会在第一时间就把“单调时钟”和“push-with-credit”这两块做进去,而不是等压测踩坑之后再加。前者花两个小时就能改完,后者也只需要一条“额度上报”的链路,但它们在关键时刻能省下整整一周的排障时间。
最后分享一个我调试 ax 时的小技巧:全链路压测时,不要把每个任务的延迟都设成固定值。把延迟设成随机数(比如random(0, 2000ms)均匀分布),更容易暴露时间轮、队列、派发链路上的共振问题——固定延迟会让所有任务在同一时刻触发,队列瞬间被冲垮,掩盖掉很多真实分布下的性能特征。这个技巧也适用于任何任务系统的压测场景,实测非常有效。