☰
自研分布式任务调度中间件AX:架构设计、幂等保障与事故复盘
2026/9/26 19:24:22 网站建设 项目流程

1. 为什么放着现成的调度框架不用,非要在团队里自研一个AX

两年前我们团队接到一个比较头疼的需求:现有业务系统里的定时任务、延迟任务、数据对账任务加起来超过两万个,日触发量到了千万级别,而且很多任务要求分钟级甚至秒级准时触发。当时线上用的是单机Quartz加手动分环境部署,每逢大促前都要靠人工去核对任务是否有重叠,凌晨两三点被任务抖动告警叫醒是家常便饭。也因为这个背景,我们决定评估是不是要引入一套分布式调度中间件,结果评估来评估去,最后走了自研路线——内部代号叫AX。

先说清楚AX到底是什么:它是一套面向Java后端的分布式任务调度中间件,核心解决三类问题——定时任务准时触发、延迟任务可靠投递、大批量任务动态分片执行。它跟市面上开源的XXL-Job、ElasticJob不同之处在于,我们围绕自己业务的高峰场景做了不少针对性的设计,尤其是秒级任务的抖动控制、任务实例的幂等去重、以及调度链路的全链路埋点。这篇文章会把AX从架构设计到上线压测再到线上事故复盘的全过程整理出来,适合正在调研分布式调度方案的后端工程师、架构师,也适合那种项目里已经被定时任务搞到头大的同学参考。

为什么不是直接用开源框架?这不是我们矫情,是真的算过账。XXL-Job的核心调度模型是调度中心轮询任务表,拿到待触发任务后分发到执行器。这个模型在万级任务、分钟级触发的场景下很稳,但当我们把触发粒度压缩到秒级,再叠加每秒钟上千个待执行任务时,调度中心对任务表的轮询压力和数据库连接占用会变得非常难看。而ElasticJob的分片模型很漂亮,但它的强依赖是ZooKeeper,团队里当时没人愿意为了一个调度系统再去维护一套ZK集群。

还有一个问题在于动态分片。我们有很多数据对账类任务,需要对几千万的数据按客户维度拆成几百片跑,执行器节点又在不停上下线,开源框架的静态分片在这种场景下要么手动干预太重,要么重平衡时会打断正在执行的任务。所以最终决定自研一套轻量级的调度器,目标就三个:靠谱触发、智能分片、全程可观测。别的花活一概不做。

这个决定背后还有一个现实考量——调度系统是典型的"平时没存在感、出事全责在你"的组件。一旦用了开源框架,出了线上问题你只能去翻它的源码,而自研至少每个设计决策我们都清楚当初为什么这么做,排查问题时的思路路径是完整的。

2. AX调度器的整体架构:注册中心、时间轮和分片器是如何协作的

AX的架构没有搞得很复杂,核心由四个角色组成:调度器(Scheduler)、执行器(Executor)、注册中心(Registry)和控制台(Console)。

调度器负责生成触发计划、驱动任务到期、完成分片计算。执行器部署在业务侧,负责接收调度指令并真正执行业务逻辑。注册中心用的是Etcd,没选ZooKeeper也没选Consul,原因是Etcd的租约机制和Watch推送配合得很好,任务执行器上下线的感知延迟能控制在几百毫秒内。控制台是纯前端项目,主要用来查看任务状态、手动触发、调整参数。

2.1 任务从注册到触发的完整链路

一条任务从创建到真正被执行,要经过下面这些环节:

首先是任务注册。业务方在控制台创建一个任务,需要指定任务的类型:定时任务(cron表达式)还是延迟任务(延迟秒数),以及执行器名称、回调接口、超时时间、重试次数和分片策略。注册信息写入Etcd,调度器通过Watch监听任务列表的变化。

然后是触发计划的生成。调度器内部为每个任务维护一个触发器,定时任务会解析cron表达式计算出下一次触发的时间点,延迟任务则是在提交时直接算出计划触发时间。触发时间到点后,调度器会从任务队列里取出任务,做前置校验,再进入分片逻辑。

接着是分片计算。AX的分片结果是一个包含目标执行器ID和分片参数的消息体。举个例子,一个数据对账任务配置了10个分片,当前有4个执行器在线,调度器会根据分片策略把1、5、9分配到一个执行器节点上,而不是愚蠢地按顺序切。这里的设计逻辑是:尽量让每个执行器拿到的任务数量均衡,同时降低重平衡时迁移的任务量。

最后是任务下发和结果回传。调度器把任务消息投递到执行器对应的消息通道,执行器执行完毕后回传状态、耗时和业务附加信息。调度器端维护一张进行中的任务表,只有收到终端状态后才算闭环。

2.2 时间轮的选择:为什么不用延迟队列

任务触发的核心难点在于:海量任务挂在内存里,如何高效判断哪些任务已经到期。第一版我们用的是Java的DelayQueue,按触发时间排序,到期就弹出。说实话功能没问题,但有个尴尬的场景:当任务数量超过五万个时,DelayQueue的插入和删除操作是O(log n)复杂度,高并发下锁竞争很明显,调度线程经常空转。

后来改成了多级时间轮。AX的时间轮参数是这样设置的:tick时长1ms,第一层wheelSize是512个slot,第二层是64个slot,每个slot代表一层的时间跨度。第一层覆盖512毫秒内的触发任务,超过的进入第二层,第二层每个slot代表512毫秒,总共覆盖约32秒。为什么做两层而不是直接一个大时间轮?因为大部分延迟任务集中在分钟级以内,两层时间轮的内存占用很低,每个slot只需要维护一个环形链表。

时间轮的驱动线程是单线程的,每秒转动一圈,把到期的任务拿出来做分发。这里有个细节容易忽略,就是时间轮的推进频率和调度线程池的解耦。时间轮只负责"到点把任务标记为可执行",真正的执行是丢给后面的线程池去处理,这样即使某个任务执行回调很慢,也不会阻塞后续任务的到期判定。

2.3 执行器注册与心跳保活机制

执行器启动时会向Etcd注册一个临时节点,节点带有执行器ID、IP、端口和健康检查路径。注册之后,执行器与Etcd之间保持租约,默认租约时间是30秒,执行器需要每10秒续约一次。如果调度器连续三次没有收到执行器的心跳,调度器就判定该执行器下线,触发分片重平衡。

这里我想强调一个设计取舍:心跳和租约是分开的。心跳用来感知存活,租约用来控制节点的生命周期。两者不能混为一谈,否则会出现这样的情况:执行器业务线程卡死但心跳正常,调度器依然往里使劲派任务,最后阻塞在业务侧。

实际上AX最终把健康检查做成了两级:一级是基础心跳,执行器进程活着就算正常;二级是任务队列水位,当执行器内部的任务积压超过阈值时,执行器会主动向调度器上报"繁忙"状态,调度器收到该状态后会将新任务分片到其他空闲节点。这个机制带来了非常好的效果:大促期间扩缩容时,任务的倾斜程度明显下降。

3. 调度正确性的命门:时钟、锁和幂等,这三件事环环相扣

说到分布式调度的正确性,大多数人第一反应是"不能让任务重复执行"。但真正动手做会发现,重复只是表象,导致重复的原因往往有三个:时钟不同步、锁失效、幂等没有兜底。这三个问题互相纠缠,你只解决其中一个,另外两个会在某个深夜准时出现。

3.1 时钟漂移:这是最容易被忽视的坑

调度系统里有两个时钟概念,一个是墙上时钟(Wall Clock),比如2025年某月某日14点30分;另一个是单调时钟(Monotonic Clock),表示从某个时间点开始流逝的秒数,只增不减。判断一个任务是否到期,只能依赖单调时钟,因为它不受NTP校时回拨的影响。

但业务上创建延迟任务时,用户指定的往往是墙上时间。这里就需要一个转换:任务创建时,把墙上时钟转为单调时钟的基准偏移量存入任务元数据,后续时间轮判断全部基于单调时钟。这套逻辑落地后,我们遇到的"任务提前几十秒触发""任务触发时间忽前忽后"的现象基本绝迹。

真实的线上案例是:某次机房服务器NTP配置错误,导致一整批机器的墙上时钟比标准时间快了4分钟。如果触发判断直接用系统当前时间,那一批延迟任务将全部提前4分钟执行,上游还没准备好数据,下游直接拉取,结果是批量报错。换成单调时钟之后,这类问题被彻底隔离在时间轮之外。

3.2 分布式锁的二段式设计

调度器集群模式下,同一个任务可能同时被多个调度器节点感知到期,所以必须在触发前加一把分布式锁。AX最初的实现很天真,直接用Redis的SETNX命令,锁的key是任务ID,获取到锁的调度器才允许触发。

这套实现上线后遇到了一个典型问题:任务执行时间超过了锁的过期时间,调度器A执行到一半锁过期了,调度器B获取到锁又把任务触发了一次,下游重复扣款,客户投诉电话直接打到CTO那里。

后续我们改成了带版本号的锁,也叫fencing token方案。调度器在获取锁时,Redis返回一个自增的token,任务执行时每次写状态都要带上这个token,业务侧校验token是否是最新版本,如果不是就拒绝写入。同时锁的过期时间改成了动态续期,而不是固定值:每次续期脚本检查当前任务是否还在执行,如果还在执行就延长过期时间。这套机制让大部分重复问题在源头被拦截。

3.3 幂等兜底:最终防线必须靠任务实例ID

分布式环境里,任何锁都不是绝对可靠的。所以AX在设计时就把"任务必须幂等"作为硬性要求,不是建议,是必须。

具体做法是:每触发一次任务,调度器生成一个全局唯一的任务实例ID,格式是任务ID+时间戳+随机数。这个实例ID会随调度消息一起发给执行器,执行器在处理任务前先检查本地是否已经处理过这个实例ID,如果是就丢弃。下游业务系统接收任务回调时,也要求以实例ID作为唯一键做去重。

数据库层面我们建了一张唯一的去重表,表里只有一列主键,就是任务实例ID。触发前插入,成功了才继续执行后续逻辑。这张表的作用是防止调度器端重复分发、执行器端重复执行、下游回调重复接收。三层防护下来,重复率从最初的千分之几降到了百万分之几。这个设计才是AX调度正确性的真正地基,锁只是在前面挡子弹的。

4. 线上事故复盘:一次重复调度引发的完整排查链路

讲讲我们上线后遇到的最严重一次事故,整个过程非常典型,对理解分布式调度非常有价值。

4.1 故障现象:凌晨的重复扣款

某个业务方接入AX后第一次跑月度结算。凌晨3点,数量约8000条的结算任务按计划触发,分片到6个执行器节点。执行到一半,下游账务系统突然告警,同一笔款项被重复扣了两次,涉及约200个用户。

当时第一反应是查执行器日志,发现部分任务确实在同一时间点被两个不同的执行器节点各执行了一次。更诡异的是,被重复执行的任务分片基本集中在某两个节点上,而其他节点的任务执行完全正常。

4.2 第一轮排查:锁与幂等日志全部正常

我们拉取了调度器端的分发日志,显示每个任务实例ID在调度器侧只生成了一次,锁的获取和释放记录也完整,没有异常覆盖。执行器端去重日志同样显示实例ID没有重复处理。这意味着问题不是锁失效,也不是重复分发,而是同一任务被生成了两个不同的实例ID。

顺着这个方向查,我让运维把两个被执行节点的任务接收时间都打出来,发现时间戳相差只有1.2秒,都在各自节点认为的"合理触发窗口"内。这说明两个节点都有自己的理由认为自己应该触发这个任务。

4.3 根因定位:Full GC引发的租约续期失败

继续深挖,找到了一条被大家忽略的链路——执行器节点上线前,注册信息会同步到调度器端,而调度器端维护了一份在线执行器列表。正常情况下,执行器节点A下线后,任务分片才会迁移到节点B。但这次事故中,节点A和节点B是并存的,两个都在调度器看来是"在线"状态。

为什么两个节点同时在列表里?看监控数据发现,节点A在凌晨3点前后发生了两次Full GC,每次停顿时间超过4秒。而执行器与Etcd之间的租约续期间隔是10秒,两次GC导致节点A错过了租约续期,Etcd端判定节点A过期,触发了节点下线事件。

节点A的下线事件传到调度器后,调度器把节点A持有的分片重新分配给了节点B。此时节点B开始执行原本属于节点A的任务。可是节点A在GC结束后并没有真正退出,它从Etcd租约恢复后,又通过心跳重新上线,此时节点A认为自己是新节点,继续消费队列里积压的旧任务。节点A执行了一份,节点B执行了一份,虽然实例ID不同,但对应的业务处理逻辑是同一个定时任务,于是重复扣款发生了。

整个过程可以用一句话总结:节点假死被踢出,GC恢复后又重新上线,新旧实例同时跑同一批任务。这跟注册中心的数据一致性没关系,而是任务队列里的消息没有做归属标记,节点恢复后不知道自己已经被顶替。

4.4 修复方案与验证

这个事故让我们复盘出了四条整改措施,每条都很具体。

第一,执行器注册时增加启动纪元(Epoch)。纪元是一个自增的单调递增数字,节点每次上下线都会变化。任务消息投递时携带目标纪元,执行器收到消息后先对比自己的纪元,不一致就丢弃。这个机制从根源上防止了旧节点恢复后继续消费新任务。

第二,租约续期增加重试缓冲。执行器的续期操作在底层加了独立线程池和重试队列,即使业务线程发生长时间GC,只要JVM进程还活着,续期线程仍然能撑住,默认重试5次,每次间隔500毫秒。

第三,调度器端把任务从"已分配"到"已完成"的最长生命周期做了硬性限制。超过限制还没返回结果的任务,不允许被重新分片。这条规则能拦掉大部分因节点假死导致的重复执行场景,宁可让任务挂起待人工确认,也不要贸然派给别的节点。

第四,执行器端新增了实例消费记录的快照表。每消费一条任务消息,先把消息内容和消费时间写入本地磁盘,重启后加载去重。这个兜底解决了消息队列在节点长时间不可用期间积压后重新放量的场景。

整改之后,我们专门压了一轮故障演练:手动触发节点的Full GC,对比整改前后的重复触发次数。整改前重复率约为千分之二,整改后重复次数降为0,连续三轮演练都没有复现。

5. 压测后的调参记录:线程数、队列容量和分片数的平衡

AX上线前我们在测试环境做了几轮压测,压测过程中暴露出来的问题比功能测试时的Bug有意思得多,主要集中在参数调的平衡,因为这些参数互相牵扯,不是拍脑袋能定下来的。

5.1 压测环境与基础数据

压测环境是三台调度器节点、五台执行器节点,每台机器配置是8核16G。模拟的任务模型参考线上真实分布:40%定时任务,60%延迟任务,任务平均执行耗时约800毫秒,目标吞吐是每分钟10万次触发。

压测过程分三档:每秒500触发、每秒1500触发、每秒3000触发。每一档跑30分钟,看调度器端的触发延迟、执行器端队列堆积、以及任务成功率的分布。

5.2 线程池参数对吞吐的影响对比

第一轮压测,调度器端触发线程池配置的是固定线程数32,队列容量用默认的10万。结果是每秒1500触发时,触发延迟出现了明显的锯齿形波动,最高延迟到了400毫秒,而P99只有80毫秒。这说明线程池在任务瞬时时长暴涨时产生排队,造成触发高峰期延迟异常。

我们把触发线程池从固定线程数改成了核心线程数16、最大线程数64、队列容量2万的弹性配置。同时把任务下发改成批量模式,不是一条一条地发,而是将时间窗口内到期的任务聚合成一个批量消息,减少网络开销和线程切换次数。调整后,每秒3000触发时P99稳定在120毫秒以内,锯齿波动消失。

执行器端的调参也很有讲究。Worker线程数最初设置的是CPU核数的两倍,也就是16,后来发现大部分任务都在等待下游接口返回,IO密集型场景下线程数明显偏低。改成32后吞吐提升了约35%,但再往上加到48,吞吐不仅没提升,任务超时率反而增加了。原因是线程切换成本和下游系统承受的连接压力上来了。最终停在32这个值。

5.3 分片不均:一个带着惯性思维踩进去的坑

排查分片问题时,我们一开始的习惯是看各执行器CPU和内存负载,发现负载很均匀,就认为分片没问题。直到有一天某个执行器节点突然报OOM,查看任务分配详情才发现,它实际处理的任务数量是其他节点的3倍。

根因出在分片策略选择上。分片键我们最初用的是用户ID模10,而业务数据的分布天然偏好某些尾号。比如尾号0和5的用户数是尾号3和7的三倍,那分片自然不均。这个问题只靠"分布式"解决不了,得回到数据分布本身。

后续我们改成了基于统计的分段分片。具体做法:在每个执行器节点上按天采集任务处理量,调度器定期汇总成一张分片量分布表。分片时不是拿任务ID做均匀切分,而是根据这张表把"预计耗时相近"的任务段合并成一个分片组。虽然计算复杂了一些,但效果非常明显,各节点任务耗时方差降了70%以上。

5.4 GC与触发延迟的联动调优

压测后期我们注意到一个现象:调度器端每两个小时发生一次Full GC,每次停顿约600毫秒,这期间触发延迟飙升到秒级。

用G1替换CMS后,Full GC变为Mixed GC,单次停顿降到150毫秒以内。但G1的Region大小和期望停顿时间需要调,我们最终把Region大小设为4MB,MaxGCPauseMillis设为200,同时在调度器端给时间轮驱动线程单独设置了守护线程优先级,确保GC期间它不会被业务线程池完全抢占。

这里有一条经验值得记录:调度器这种组件,GC调优的影响比业务系统大得多。业务系统GC停顿几百毫秒可能无所谓,但调度器停顿几百毫秒,可能导致数百个任务错过触发窗口。所以AX从架构上把时间轮驱动线程和业务执行线程做了严格的隔离,任何情况下时间轮驱动线程不允许被阻塞。

6. 最后再聊聊调度系统的运维观:从工具思维到产品思维

AX上线运营到现在,我最大的体会是:调度系统的价值不在于它多快多稳,而在于出故障时它能让你在几分钟内搞清楚发生了什么。这个认知是几次线上事故后慢慢建立的。

以前我们看调度系统,就像看一个工具:配置好任务、能触发、能回调就认为万事大吉。现在我们把AX当产品来运营,重点盯的是三类数据:触发成功率、触发延迟分布、任务生命周期状态机流转异常数。任何一个指标出现拐点,都比单看业务成功率更能提前暴露问题。

结合AX的落地经验,给想要自研调度系统或者正在改造调度架构的团队三条建议。

第一,先埋点再上线,不要等出了问题再补监控。AX第一版上线时监控只有任务成功率,导致第一次事故排查花了三个小时。后来补上了调度器到执行器每一跳的耗时、队列积压量、续期成功率、锁获取失败次数,再排查问题基本都能定位到具体环节。

第二,把幂等当成一等公民设计,而不是事后补救。无论你的锁做得多完美,无论你的调度算法多精妙,分布式环境总有你想不到的场景让你的任务执行两次。任务实例ID加唯一约束,这一条能让你的系统在无数次故障演练中都站得住。

第三,分片参数要留出动态调整的接口,别把分片配置写死。我们去年的压测数据只能说明当时的业务模型。半年后业务字段分布变了,静态的分片参数就会变成瓶颈。AX后来把分片策略做成了可插拔的SPI,线上调整分片参数不用重启,这个改动间接解决了很多数据倾斜问题。

如果你也是在业务规模到了一定程度后被定时任务搞到头疼的人,我建议你先别急着把开源框架拿过来优化,而是花一个周末把以下三件事想清楚:你的任务分布是什么样、你的失败容错边界在哪、你调试一个任务从触发到结果要几步。这三件事想清楚了,哪怕最终不写一行框架代码,你的调度系统也会比现在稳定一个档次。

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

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

立即咨询