ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

自研分布式任务调度中间件AX:架构设计、幂等保障与事故复盘

自研分布式任务调度中间件AX:架构设计、幂等保障与事故复盘 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大小设为4MBMaxGCPauseMillis设为200同时在调度器端给时间轮驱动线程单独设置了守护线程优先级确保GC期间它不会被业务线程池完全抢占。这里有一条经验值得记录调度器这种组件GC调优的影响比业务系统大得多。业务系统GC停顿几百毫秒可能无所谓但调度器停顿几百毫秒可能导致数百个任务错过触发窗口。所以AX从架构上把时间轮驱动线程和业务执行线程做了严格的隔离任何情况下时间轮驱动线程不允许被阻塞。6. 最后再聊聊调度系统的运维观从工具思维到产品思维AX上线运营到现在我最大的体会是调度系统的价值不在于它多快多稳而在于出故障时它能让你在几分钟内搞清楚发生了什么。这个认知是几次线上事故后慢慢建立的。以前我们看调度系统就像看一个工具配置好任务、能触发、能回调就认为万事大吉。现在我们把AX当产品来运营重点盯的是三类数据触发成功率、触发延迟分布、任务生命周期状态机流转异常数。任何一个指标出现拐点都比单看业务成功率更能提前暴露问题。结合AX的落地经验给想要自研调度系统或者正在改造调度架构的团队三条建议。第一先埋点再上线不要等出了问题再补监控。AX第一版上线时监控只有任务成功率导致第一次事故排查花了三个小时。后来补上了调度器到执行器每一跳的耗时、队列积压量、续期成功率、锁获取失败次数再排查问题基本都能定位到具体环节。第二把幂等当成一等公民设计而不是事后补救。无论你的锁做得多完美无论你的调度算法多精妙分布式环境总有你想不到的场景让你的任务执行两次。任务实例ID加唯一约束这一条能让你的系统在无数次故障演练中都站得住。第三分片参数要留出动态调整的接口别把分片配置写死。我们去年的压测数据只能说明当时的业务模型。半年后业务字段分布变了静态的分片参数就会变成瓶颈。AX后来把分片策略做成了可插拔的SPI线上调整分片参数不用重启这个改动间接解决了很多数据倾斜问题。如果你也是在业务规模到了一定程度后被定时任务搞到头疼的人我建议你先别急着把开源框架拿过来优化而是花一个周末把以下三件事想清楚你的任务分布是什么样、你的失败容错边界在哪、你调试一个任务从触发到结果要几步。这三件事想清楚了哪怕最终不写一行框架代码你的调度系统也会比现在稳定一个档次。
返回列表