☰
ax调度:轻量级分布式定时任务引擎的设计与实践
2026/9/26 8:53:06 网站建设 项目流程

项目代号定为ax的时候,我一度觉得这名字随意得像随手敲的。后来有同事问起,我解释为 Auto eXecution 的缩写——一个只负责自动触发、自动调度的小引擎。业务方把越来越多的定时任务、延时任务、批处理任务丢过来之后,"ax调度"反而成了团队里出现频次最高的词,比它正式的名字好记太多。这篇文章就是把 ax调度 从设计、实现到上线踩坑的过程完整拆一遍,适合正在做内部调度系统、或者准备接手类似项目的同学参考。

像很多内部系统一样,ax 一开始并不是什么大工程。它只花了两周写出来,本意是替代散落在各服务里的定时器。结果半年后,它开始承担大促场景下上万级任务的编排,每天处理几百万次触发。回头看,它既不是最全能的分布式调度器,也没有炫酷的可视化界面,但它解决了一个很具体的问题:在一个以业务交付为主的技术团队里,用最小的维护成本把"到点该做的事"统一管起来。写这篇文章,目的是把这类轻量调度系统里真正有效的部分讲清楚,也希望你能少踩几个我已经踩过的坑。

1. 一个叫 ax 的项目,为什么值得单独写一篇文章

1.1 散落在各处的定时器,才是最初要解决的敌人

催生 ax 的根本原因,不是什么高深的技术选型,而是"定时器到处长"。当时各个业务服务里都有调度逻辑:有人用 goroutine 加 select 做超时控制,有人用一个跑在常驻进程里的 cron 表达式定时执行,还有人干脆在用户的请求路径里判断当前时间戳,时间到了顺手执行一段逻辑。任务少的时候,这种做法还行得通,顶多是在代码仓库里多搜几个关键字。

可一旦任务量起来,问题就非常现实:没人能说清楚当前一共有多少定时任务在跑,每个任务上次什么时候执行的、失败了多少次、有没有重试、日志去了哪里。更难受的是,某个服务发版重启后,内存里的定时器全部归零,业务方隔天来问"为什么昨天凌晨的推送没发",排查要大半天。所以 ax 要解决的第一件事,不是"调度效率更高",而是"让所有需要到点执行的东西,有一个统一出口,并且状态可查、失败可追"。

1.2 为什么没直接用现成的调度框架

动手之前,我把市面常见的方案过了一遍,也说一下为什么没有直接拿来用。

  • xxl-job:调度能力和管理界面确实成熟,但需要额外部署调度中心,权限、报警体系也更偏向 Java 团队。我们的执行器主要是 Python 服务,接入成本不算低。
  • Airflow:DAG 编排是一等公民,适合数据管道场景,但对短延时和秒级触发不太友好,调度器本身偏重,装一套最小集群也不轻松。
  • APScheduler:嵌入应用内很舒服,但任务状态不持久化、分布式支持弱。多实例部署同一个服务时,每个进程都会独立调度,任务就会被重复触发。

ax 的设计目标因此定得非常克制:调度引擎和业务逻辑分离,管理者只需要注册任务和触发器,执行器挂到引擎上,由引擎决定在什么时间点回调谁。它不做工作流审批,不做数据血缘,只做"到点触发"这一件事。明确边界之后,自研的成本一下子就降下来了。

1.3 先做一个"五分钟"级别的可行性验证

当时我先写了一个最小 Demo,只支持两个任务:一个每分钟执行一次,另一个延时 30 秒后执行。存储用 Redis,触发用最简单的时间轮。这么做不是为了炫技,而是为了快速验证两个问题:一是这种模型能不能覆盖现有业务里的绝大多数场景,二是运维起来麻不麻烦。跑通之后才决定正式立项。现在回头看,这个"先做最小路径"的习惯很管用——如果一开始就想去设计一个完美的分布式调度器,大概率两周内做不出来。

2. AX调度的任务模型:把"触发"和"执行"解耦是第一步

2.1 任务、触发器、执行器:三张表撑起整个模型

ax 的数据模型没有走复杂的微服务拆分,核心就三类对象:任务、触发器、执行器。

任务(Task)描述"做什么",包含任务名、执行器 ID、超时时间、重试策略、当前状态等字段。触发器(Trigger)描述"什么时候做",包含类型、类型相关的参数、是否启用等字段。执行器(Executor)描述"由谁来做",是暴露给调度器的回调地址或本地函数,需要实现一个注册协议。

这三层分离之后,很多需求就变成了简单的数据操作。比如想每天凌晨跑一次数据统计,只需要给同一个任务挂一个 cron 类型触发器;想临时改成每小时跑一次,改触发器即可,任务本身不动。反过来,同一个触发器也可以挂在多个任务上,只要把任务列表塞进去就行。这种模型在代码里体现为三张表,接口层面的成本很低,但给后续扩展留下了非常大的空间。

2.2 触发器的核心接口:next_time 和 is_due

为了让调度循环不关心具体触发器类型,所有触发器都实现了同一个接口。用 Python 写大概像这样:

class BaseTrigger: def next_time(self, after: datetime) -> datetime: """返回 after 之后的下一次触发时间。""" raise NotImplementedError def is_due(self, now: datetime) -> bool: """判断当前时刻是否应该触发。""" next_run = self.next_time(now) return next_run is not None and next_run <= now

cron 类型就解析 cron 表达式算出下一个时刻,interval 类型就基于上次执行时间加固定间隔,once 类型则直接比较预设时间点是不是已到。调度循环不再需要判断"这个任务到底是日切、周切还是延时任务",只需要统一调用next_time,把算出的时间点放回调度索引里。

这里有个细节值得注意:next_time必须基于"上次执行后的语义"来算,而不是每次都用当前时间重新对齐。否则一个周期任务如果执行耗时太长,下一次触发会被当前时间带着漂移,最终整个周期全部错乱。这也是很多用while True: sleep(interval)写定时任务的人会踩的坑。

2.3 依赖编排:任务之间不是越多越好,而是要有先后

很多调度场景不是孤立的,任务之间有上下游关系。比如先同步数据,再生成报表,最后发推送。ax 里我用了一张依赖表和一张事件队列来实现。

一个主任务执行完成后,会把完成事件写入事件队列。依赖服务消费这个事件,判断它下游任务的"前置条件"是否全部满足,满足才把下游任务标记为可调度。如果前置任务失败,默认策略是直接标记下游任务为 failed,不盲目往下走。依赖图在注册阶段做过拓扑排序校验,一旦发现有环,直接拒绝创建,避免出现两个任务互相等待导致死锁。

这个模型很简单,但没有引入重量级的 DAG 引擎。实际用下来,普通团队 90% 的编排需求就是"先 A 后 B"或"多个 A 完成后做 B"这两种,一张依赖表完全够用,没必要把复杂度引进来。

3. 调度器内部的时间轮与状态机,决定了稳定性上限

3.1 时间轮不是越精细越好:tick 与 wheelSize 的取舍

调度引擎内部我选用了时间轮(HashedWheelTimer)。它的原理很简单:一个环形数组,每个槽位挂一个待触发任务链表,指针按照固定 tick 间隔向前移动,落到哪个槽就把哪个槽的任务拿出来执行。用生活里的例子类比,就像一个钟表,秒针每走一格,就把这一格抽屉里到期的便签全部拿出来处理。

ax 的参数最终定为 tick=1 秒、wheelSize=3600,也就是说时间轮可以覆盖未来一个小时内的任务。内存占用就是 3600 个槽位,非常小。为什么不用 100ms 甚至更细的 tick?因为 ax 面向的业务是推送、批处理、数据同步,秒级延迟已经完全满足需求。tick 越细,指针空转频率越高,CPU 浪费越大。如果真有毫秒级触发需求,那应该考虑消息中间件的延时队列,而不是在调度器里死磕精度。

对于延时超过一小时的任务,时间轮本身存不下,我在任务结构里加了一个"圈数"概念:任务实际到期时间除以时间轮覆盖范围,得到的商就是还需转几圈。指针每轮回到对应槽位时,把圈数减一,减到零才真正执行。这种处理方式避免了链表无限膨胀,代价是极端场景下大延时任务的精度会随轮转略有偏移,但对秒级业务来说可忽略。

3.2 任务状态机:pending、running、success、failed、retry

状态机不需要花哨,但状态转换规则必须严格。

  • pending:任务注册成功、等待触发,这是所有任务的起点。
  • running:调度器把任务投递给执行器之后立即进入,而不是等执行器反馈后再进入。
  • success:执行器回调上报成功。
  • failed:执行器返回失败,或回调超时。
  • retry:failed 后,如果重试次数没超过上限,就进入 retry,重新计算下次触发时间。

这里有一个容易被忽略的规则:状态持久化必须在转换发生时同步落库,而不是等整个任务流程结束时统一保存。如果只在任务结束后保存,崩溃恢复时系统就分不清这个任务到底是"还没触发"还是"已经触发正在运行"。这两种情况的恢复策略完全不同。调度器恢复时,只要那些还处于 pending 和 retry 状态的任务,running 状态的任务交给执行器侧的幂等逻辑去兜底。

3.3 崩溃恢复:先选"至少一次",再用幂等兜底

分布式调度绕不开语义选择。ax 默认采用 at least once(至少一次),也就是任务可能被重复触发,但绝不能因为崩溃而丢失。配合执行器侧的幂等设计,重复触发不会产生重复的副作用。

具体恢复机制分两层。内存里,时间轮负责快速判断"现在该触发谁";数据库里,任务表负责记录所有任务的持久化状态。系统运行期间,除了状态转换写库之外,我还会每隔一段时间把时间轮里即将到期的任务快照写一次。崩溃重启后,调度器扫描数据库中所有 pending 和 retry 任务,重新计算各自的 next_time,再放回时间轮。

这个过程有一个默认策略很关键:已经错过多次触发窗口的任务,恢复后只补偿一次,不追着补跑。因为一旦系统停机超过几个小时,补跑所有错过的批次会让下游系统瞬间被打爆。把"补偿一次"作为默认值,业务方如果有特殊需求再按任务单独放开。

3.4 执行超时、线程池与调度漂移

执行器回调很容易出现超时。ax 的默认回调超时是 30 秒,超时后任务标记为 failed 并走重试策略。但真正要小心的不是单次超时,而是线程池被打满。

调度器向执行器发起请求时用的是有界线程池,队列长度默认 5000。如果等待队列已经满了,新任务直接拒绝并进入 retry,而不是把队列做成无界队列,否则内存会先炸。这个取舍非常重要:宁可让任务重试,也不能让调度器进程因为堆积而崩溃。

另一个经验是"调度漂移"。比如一个每 5 分钟执行一次的任务,某次执行花了 15 分钟,等到结束时,按固定间隔算法算出的下一个触发时间点已经过去了。ax 的默认行为是:无论推迟了多久,只补偿一次并立即触发,然后回到正常节奏。如果不做这个限制,一个慢任务会把后续所有触发时间全部向后挤,整个调度序列彻底乱掉。

4. 上线前我踩过的四个坑:每个都差点导致线上事故

4.1 第一个坑:分布式锁的自动过期,导致任务被重复执行

最开始为了保证双节点不重复触发任务,我直接用 Redis setnx 抢锁:谁的锁过期时间更长谁就拥有触发权。测试环境一切正常,直到一次线上压测时发现,同一个任务在极短的时间内被触发了两次。

排查过程很痛苦。日志显示两件事几乎同时发生:节点 A 拿到锁,但进程发生了一次接近锁过期时间的 GC 停顿;节点 B 在锁自动过期后抢到了同一把锁,于是两个节点同时执行触发逻辑。这个问题的本质是:分布式锁的过期时间太"静态"了,没法感知持有者的健康状态。

修复方案分两层。第一层,给锁加续期机制,也就是租约续期——持有锁的节点每隔一段时间延长锁的过期时间,节点失联后锁才会真正释放。第二层,也是最关键的兜底:执行请求里带上全局唯一的 requestId,执行器侧用 requestId 做幂等。就算两个节点真的同时发起了触发,执行器也只会接受第一个 requestId 对应的请求,后面的直接丢弃。这比任何聪明的锁方案都可靠。

4.2 第二个坑:数据库扫表成为了新的瓶颈

最初版本里,调度器每秒执行一次SELECT * FROM tasks WHERE next_time <= now AND status IN ('pending', 'retry'),把到点任务捞出来。任务量小的时候没问题,到几万条时数据库 CPU 开始报警,慢查询日志里全是这张表的全表扫描。

后来做了两个改动。第一,把"每秒全表扫描"改成"时间轮到期槽位批量取任务 ID",再用WHERE id IN (...)把任务捞出来,查询量立刻下降了几个量级。第二,引入 Redis ZSet 作为二级索引,score 就是任务的秒级到期时间。调度器只查ZRANGEBYSCORE拿到到期任务 ID,数据库只负责最终的状态持久化。经过这一轮改造,高频扫描压力彻底转移到了 Redis 上,数据库回归到它擅长的持久化角色。

4.3 第三个坑:手动触发和自动触发抢同一个任务

管理后台加了一个"立即执行一次"的按钮,让业务方可以手动触发。当时实现得很偷懒,直接调用了和自动触发同一个接口。结果上线第二天就出了事故:一个数据同步任务既被手动触发了一次,又被自动调度触发了一次,两边同时跑,把下游表写重了。

问题根源在于任务模型里没有区分触发来源。后来的修复方案是:任务模型增加trigger_source字段,手动触发走单独的通道并带上force标记;同时约定,同一个任务在同一时刻只能存在一个 running 实例,靠 Redis 原子递增和请求 ID 双重校验来实现。如果手动触发时任务已经在 running,系统会明确提示"任务正在执行中",而不是静默地再发一次。

4.4 第四个坑:时钟回拨让调度直接乱掉

这个问题排查时间最长,因为开发环境几乎无法复现。现象是:某次线上任务出现了重复触发,但去查锁、查幂等都没有问题;后来发现同一批任务在某些时刻被反复扫描,而且时间越往前,扫描次数越多。

最终定位到是云主机发生了一次 NTP 时间回拨。时间轮指针基于墙上时钟,系统时间往回跳了一秒,指针也跟着倒退,于是已处理过的槽位又被扫了一遍。修复方案很明确:调度器内部计时全部改用单调时钟。在 Python 里用time.monotonic,它只保证两次调用之间的间隔是递增的,不受系统时间调整影响;墙上时钟只用于生成任务的实际触发时间点。分布式环境下,多节点之间的对时以数据库时间或 leader 节点时间为准,不使用本地墙钟。这个改动修复之后,时钟回拨类问题再没出现过。

5. AX调度实测:数据说话,顺便说点它的边界

5.1 我们环境里的压测结果

压测环境是一台 4C8G 的虚拟机,Redis 6.0 部署在同机房。任务类型是 50 万条一次性延时任务,调度节点为单节点。数据如下:

任务总量每秒触发量P99 调度延迟错误率
5 万835 次/秒45ms0.02%
50 万8200 次/秒180ms0.05%

这里的"调度延迟"指任务到点时间到调度器真正发出回调之间的时间差。可以看到,触发量上升后延迟并没有线性恶化,说明时间轮加 Redis ZSet 的组合在单机容量内表现稳定。CPU 峰值在 75% 左右,内存约 500MB。如果任务量再上一个量级,单节点会顶不住,需要引入分片。

5.2 真实业务里的表现:每天几百万次触发

线上以 cron 和 fixed_delay 类型任务为主。高峰时期每天约 300 万次触发,调度成功率约 99.97%,失败大头在下游服务超时和依赖接口抖动。回调平均耗时大约 3.2 秒,整个调度服务的部署成本只有两个 4C8G 节点,其中一个还是备节点。

ax 真正有价值的地方,不是它每秒能触发多少任务,而是它把所有散落在业务代码里的时间判断全部收了口。团队从"某个服务里可能有个定时器"变成了"所有定时任务都在 ax 里,能看到状态、历史、日志"。这种可观测性,才是内部工具最大的收益。

5.3 边界:哪些任务别交给 AX

ax 不是万能调度器,有几类场景它明确不做。

第一,亚秒级高频触发,比如每秒几百次以上的行情推送。这种场景更适合消息中间件或流处理系统,调度器做不了精细时间控制。第二,有状态的流式计算任务,需要窗口、聚合、状态恢复这些能力,应该交给专用计算引擎。第三,复杂 DAG 且节点数以万计的工作流,时间轮这种模型不是为大规模 DAG 设计的,硬塞进来只会让依赖表膨胀到难以维护。

如果业务是真的遇到这些边界,与其改造 ax,不如在 ax 前面加一层分流:让 ax 只负责"什么时候开始",具体怎么做交给专业系统。这样边界清晰,两边都不会被拖垮。

最后分享一下我个人的体会。如果让我重新写一次 ax,我会把调度器和执行器彻底拆成两个独立进程再上线。第一版为了省事,直接在调度进程里 import 了业务执行函数,导致每次发布调度器都要连带把业务代码一起发一遍,踩过好几次头。另外,我不会再急着优化时间轮这类内部结构,而是先把一套完整的任务血缘日志做好。调度系统最难的地方,从来不是把任务触发出去,而是事后能让人准确回答"这个任务为什么在那个时间跑、到底跑得成不成功"。ax 调度能走到今天,靠的不是算法有多漂亮,而是这些朴素的边界被一次又一次踩实了。

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

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

立即咨询