从0搭建ax调度:分布式调度平台架构设计与高可用实践
有些项目一开始并不起眼等真正把它做起来你才发现它重新定义了整个团队对“稳定性”的理解。“ax调度”就是这样一套系统。刚开始它只是我手上几个跑批任务收拢在一起的小工具后来慢慢长成了覆盖全公司上百条业务线、日均触发几十万次的分布式调度平台。这篇文章我准备把从0搭建“ax调度”的过程翻出来讲透设计动机、架构选型、任务从注册到执行发生了什么、高可用怎么做、依赖和重试哪些地方容易踩坑还有我们实际碰过的疑难杂症。适合正在自建调度平台的人参考也适合和定时任务、批处理链路打交道、想知道背后机制的人阅读。1. 为什么需要一套调度系统ax调度的设计起点1.1 从Cron到调度平台的必然演进说句实话业务规模没上来之前crontab完全是够用的。早期我们每个服务都是单体跑批就是往一台机器上挂几条crontab凌晨2点对账、凌晨3点清缓存、早上6点出报表。这个阶段维护成本低问题也直观哪条任务挂了你去看那台机器就行。但业务线从2条变成20条、50条之后crontab的维护就成了一场灾难。首先是分布问题不同服务的跑批散落在各自的机器上运维和开发互相不知道对方挂了什么任务其次是执行能力问题几十台机器没法统一查看执行情况任务撞车之后谁慢谁快全靠猜最后是可靠性问题一台机器偶然重启当天的任务直接漏跑没有任何补偿机制。我记得有一次早上业务方反馈报表没出来查了半天原因就是承载任务的那台机器3点被自动更新重启了任务根本没执行。所以当任务数量突破几百个且出现了“任务成功后触发另一个任务”的诉求之后我决定不再往crontab里塞了而是认真搭一套调度平台。这个平台的核心思路就是把“什么时候触发”和“怎么执行”拆开。触发逻辑集中管理执行资源分散部署所有调度事件全部落库。ax调度这个名字就是顺手起的内部代号结果一直叫到了今天。1.2 ax调度解决的四大核心问题搭建初期我给自己定了四个硬指标后来几乎所有设计决策都围绕它们展开。第一是可见性。调度系统必须能回答“谁在跑、跑了多久、成功没有、下一步是谁”。这意味着任务实例必须有完整生命周期状态不能只扔出一个执行结果。第二是可靠性。调度中心不可用时不能丢任务执行节点故障时任务要能转移重跑。第三是伸缩性。业务和任务量翻倍时加机器就能扛住不能重构系统。第四是可控性。上线后你要能随时暂停、手动触发、跳过、模拟执行不能一改就跑数据库。这四个指标听着朴素但每一项都会把系统往工程化方向推。可见性逼着你设计实例状态机可靠性逼着你做持久化和补偿伸缩性逼着你把调度器和执行器横向拆开可控性逼着你做一套完整的运维操作接口。ax调度的整体架构就是围绕这四个问题生长出来的。2. 核心架构与关键设计取舍2.1 控制面与执行面分离ax调度的第一版架构我没有采用调度器直接调用执行器那种简单的RPC模式而是把参与角色拆成了四类Scheduler调度器、Broker任务中转、Worker执行器、Admin管理端。调度器只负责一件事根据任务配置的cron表达式、依赖关系、触发条件决定“此刻需要产生一个任务实例”。它不关心这个任务有多耗时、跑在哪台机器上。Worker才是真正干活的人它从Broker里竞争拿任务、执行、上报状态。Broker用的是一套基于Redis Stream的队列系统连接调度器和Worker。为什么中间一定要加一层Broker而不是让调度器直接告诉某个Worker执行核心原因是削峰和解耦。跑批任务往往在整点集中爆发比如每天凌晨0点会有上千个任务同时触发。如果调度器直接压给WorkerWorker瞬时负载可能飙升而且任务在途状态很难持久化。有了Broker之后调度器只负责快速入队Worker根据自己的消费能力去拉任务天然就实现了流量整形。曾有人问调度器直接写MySQLWorker轮询MySQL不行吗行是行但MySQL在大量短轮询下的IO压力很难看用Redis Stream做缓冲之后同样的任务量数据库压力降了一个数量级。2.2 为什么选用Redis Stream和数据库双写“入队用Redis但Redis可能丢数据”这件事是很多团队不敢让Redis扛调度可靠性的原因。ax调度这里做了一个设计双写。任务触发瞬间调度器先往MySQL的任务实例表写一行记录状态是PENDING拿到自增ID。然后才往Redis Stream的pending队列里写一个实例ID。Worker消费到实例ID之后会把MySQL里对应记录改成DISPATCHED或RUNNING。一个任务算真正开始执行。如果Worker执行完再更新为SUCCESS或FAILED。如果Redis整个崩掉所有pending状态还在MySQL里躺着恢复时可以把PENDING记录重新放回队列。这种设计牺牲了一点点入队延迟但换来了“不会因为缓存丢失而漏任务”的底线。Redis Stream本身的xadd和重平衡机制我们已经用得比较熟它比List结构好在哪里主要是消费者组概念成熟。多个Worker可以作为同一个消费者组的不同消费者互相不会重复消费同一条任务。它还有pending entries list可以配合xack机制做消费确认。任务被Worker拉走后如果Worker挂了没回ack消息会重新进入待处理状态这套机制天然适配我们的超时重试需求。2.3 分片与负载均衡策略调度面临的一个经典场景是一个任务需要在几万台设备上跑或者一个大表要拆成几百个分片处理。如果调度系统只支持“一个任务对应一个实例”这种粗粒度模型根本玩不转。ax调度在配置层就支持分片声明一个Job可以定义shardCount100调度器触发时会一次性生成100个子实例每个子实例带自己的shardIndex和总片数Worker执行时就知道自己处理哪一段数据。选用分片策略时我们一开始是平均分即每个Worker按拉取顺序消费子实例通过Broker天然竞争来实现负载均衡。后面发现热量倾斜很严重有些任务是一天一千次的小任务有些是一天一百万次的大任务默认策略明显不合理。后来在Worker心跳上报里增加了“当前执行队列长度”和“CPU负载”两个字段调度器入队时会参考这两个值做权重分配。虽然精确性谈不上但大任务的倾斜明显缓解了。实测下来一百个分片的任务均匀铺开之后总耗时从原来所有分片挤在一个Worker上的45分钟压缩到3分钟出头。3. 从触发到执行一次完整调度任务的流转拆解3.1 任务注册与Cron解析在ax调度里任何一个任务要先注册到Admin配置中心才能被调度器识别。任务配置大概是这样的结构{ name: billing-daily, type: cron, cron: 10 2 * * *, shardCount: 50, executor: billing-worker, params: { bizDate: yesterday }, timeout: 3600, retry: { maxAttempts: 3, backoff: exponential } }cron解析这块我们直接采用了标准5位字段秒级任务默认不用秒级任务消耗大一般建议消息机制处理。解析时需要注意时区问题服务器部署在UTC环境很常见如果不对cron做时区转换任务早跑出一个小时对数据报表来说可是重事故。ax调度在配置里显式声明时区默认Asia/Shanghai解析器会先把cron转成UTC再计算下一次触发时间。注册完成后调度器会把任务元数据加载到本地内存并开启监听。为什么加载到内存而不是每次触发都查库因为调度器的触发路径对延迟敏感每次触发都要查一把MySQL在高频触发下会有严重的压力。内存里放一份最新的job表配置变更通过info推送让内存热更新即可。3.2 触发、入队、分发、执行四步走一次任务从“时间到”到“执行成功”整体链路由四步构成。第一步触发。调度器在一个秒级循环里遍历内存中的活跃任务计算每个任务在当前时刻是否达到cron表达式要求的触发点。触发命中后生成一个全局唯一的实例ID并把实例基础信息插入MySQL的状态表。第二步入队。把实例ID写入Redis Stream等待Worker消费。这一步我故意做得快整个触发逻辑必须控制在10毫秒以内否则调度器自身就会成为瓶颈。第三步分发。Worker进程启动后会订阅Redis Stream对应消费者组。每来一批消息组内的某个Worker就会收到实例ID。Worker收到之后先更新MySQL状态为DISPATCHED再真正拉取任务的详细参数。这样可以避免一个任务被分发后因为参数拉取异常而丢失。第四步执行。Worker执行任务方法退出码或返回值会回传状态更新接口。执行成功置为SUCCESS失败置为FAILED并进入重试判断。这四步的伪代码逻辑大概是def schedule_loop(): for job in active_jobs: if should_fire(job, now): inst_id create_instance(job) stream_add(inst_id) def worker_consume(): for msg in stream_consumer(): mark_dispatched(msg.inst_id) result execute_task(msg.inst_id) update_result(msg.inst_id, result)3.3 一次跑批任务的现场实录举个例子每日对账单任务原来是业务服务里自己起的线程每到凌晨2点拉全量订单经常跑到早上5点。接入ax调度后我们把它拆成50个分片按userId哈希段分片每个分片各处理一批订单。凌晨2点10分调度器触发50个实例入队2点14分全部由Worker消费执行2点17分最后一个分片执行完成整个生成周期从3小时缩到7分钟。这中间唯一的代价就是我们要开发一个分片参数解析逻辑但这是值得的。如果你也要处理类似批处理任务我建议先算一笔账单分片处理耗时乘以总片数如果超过业务SLA就必须分片。分片粒度还要避免单分片内数据量过大一般控制在1-10分钟能处理完的量级太久的话某一分片失败后的重跑成本很高。4. 分布式高可用与一致性保障4.1 Leader Election与脑裂处理调度器如果只有单节点那它本身就是单点故障源。ax调度的调度器一开始就是多副本部署但同一时间只允许一个Leader节点对外做触发决策。选主我们用了Redis的setnx过期锁。每个调度器启动时都尝试获取锁拿到锁的是Leader每5秒续期一次。其他副本作为Follower监听锁的过期情况一旦发现Leader失联就立即竞争选主。脑裂问题是这种主备模式最容易踩的坑网络抖动导致Leader和Follower之间的锁过期Follower成为新Leader但旧Leader实际上还在运行这时候如果两个节点同时触发同一任务就会造成重复实例。我们的解法是引入fencing token机制每个Leader任期内触发任务时生成的实例ID都会带上任期号。比如Leader A的任期号为7它触发的实例ID是“7-时间戳-序号”。新Leader B的任期号是8。Worker执行时会把实例ID写入执行记录如果发现同一任务同一周期已经有了其他任期的执行记录就拒绝重复执行。有了这层遮挡脑裂下最坏情况是多一个被丢弃的重复实例而不是两份数据污染。4.2 任务状态的幂等与补偿模型任务状态机我们收敛得比较严格PENDING - DISPATCHED - RUNNING - SUCCESSFAILED则根据重试配置回到PENDING或直接终态。这五个状态之间任何跳转都必须满足前置状态条件我们用SQL条件更新来实现比如UPDATE task_instance SET status RUNNING WHERE id ? AND status DISPATCHED;如果影响行数是0说明前置状态不对直接返回冲突。这套机制防止了多个Worker并发更新同一个实例导致的状态错乱。此外每个实例本身要具备幂等性尤其是“手动补录”和“依赖续跑”场景下同一个业务日可能被触发两次。ax调度给每个实例定义了唯一键由任务名加调度周期加轮次组成数据库加唯一索引。第二次插入直接失败从源头杜绝重复跑数。4.3 时钟漂移与跨时区问题调度系统对时间敏感但分布式的时钟问题很容易被忽略。我们踩过最早的坑是调度器用的物理机NTP没配好某天时钟慢了40秒任务全部晚出发对账链路险些报警。现在所有调度器节点都强制开启chronyd做时间同步并且监控脚本会检查每台调度器与基准时间的偏差超过200毫秒就告警。跨时区问题主要体现在“天级任务”的语义上。一个只配置了cron的任务在东京和上海执行其“凌晨1点”的UTC时刻完全不同。我们的做法是在任务配置中显式声明业务时区调度器在计算触发时间时先拿业务时区的时间解析cron再转成UTC统一处理。所有实例落库的调度时间也统一记录为UTC时间戳另存一个业务日期字段供业务方消费。5. 任务编排与重试补偿写调度方案时最容易忽略的细节5.1 用DAG编排上下游依赖跑批不止是单个任务按时跑那么简单更多的是“A完成之后才能跑BB和C都成功后才能跑D”这种依赖关系。ax调度支持配置DAG链路每个任务节点可以声明dependsOn列表。调度器生成实例后不会立即触发下流而是等待上游实例进入终态然后由依赖调度器判断是否满足触发条件——上游全部成功才触发下游上游有失败则给下游打BLOCKED并等待人工处理。DAG的实现我们是用MySQL表存边关系每条边包括上游实例ID和下游Job ID。每次状态变更后系统查一遍出边看是否所有上游条件满足。如果DAG特别大这个查询要优化不能每次都全表扫。我们做了一层依赖缓存内存里只保留还没跑完的依赖关系已经终态的实例直接从缓存删除。5.2 重试到底怎么设置才合理很多人配置重试就是“重试3次”这是偷懒做法。重试的前提是“值得重试”而“无法通过重试解决”的错误重试再多次都没意义。ax调度允许把失败错误分类网络超时、依赖服务暂时不可用这类推荐配置重试参数错误、数据校验失败这类直接终态。重试策略默认指数退避比如第一次失败后1分钟重试第二次2分钟第三次4分钟最大间隔不超过30分钟。还要设置最大重试次数避免失败任务无限挂起重试队列。我这里有一个真实教训某个下游数据同步任务连不上数据库我们配了最大重试10次但没配这个任务的SLA监控结果10次重试跨了4个小时等我们发现时下游报表数据已经晚了。后来所有重试任务都加了“预计完成时间”这个虚拟字段超过预计时间的告警直接拉群。重试的本质是增加系统鲁棒性但它不应该变成掩盖故障的工具。5.3 暂停、跳过、补录三类运维操作做调度系统这三类操作一定要同时做。暂停分两级一级是暂停整个Job调度器不再生成新实例但在途实例继续跑另一级是暂停某个具体实例适用于已经触发但发现数据源异常的紧急拦截。跳过则是把某个实例直接置为SKIPPED终态常用于节假日不需要跑批的场景或者某个任务已经手动处理过了。补录是调度系统的救命稻草业务方发现某天数据漏跑时选定业务日期系统会把这一天的实例重置为PENDING并重新入队。补录时唯一键要防重复不能让同一天的任务生成两个实例同时跑。6. 常见问题排查与调优实录6.1 问题速查表ax调度上线到现在我把常见问题的排查沉淀成了一张速查表建议新接入的团队先存一下。现象可能原因排查步骤任务整体延迟Worker消费慢、队列积压看Redis Stream lag、Worker心跳负载偶发重复执行脑裂、ack超时重投看实例ID任期号、唯一索引冲突日志任务卡在RUNNINGWorker被杀、执行器僵死看Worker租约、心跳最后时间、强制失败Cron触发时间不准调度器时钟漂移ntpcheck、调度器GC耗时监控依赖任务一直不触发上游实例状态不是终态查依赖边表、上游实例状态手动触发无效任务处于暂停状态看Job配置状态、实例唯一键冲突这几项里最容易被忽略的其实是“依赖任务一直不触发”。有一次排查花了一下午结果是因为上游任务有一个分片卡在RUNNING状态而RUNNING的实例心跳早已超时但系统又没有自动判定失活。后来我们给运行中的实例加了一根“心跳租约”线Worker每30秒续一次超过2分钟没续就自动判定异常并触发重试。6.2 调优心得如何把调度吞吐撑到每秒3000次调度器最容易成为瓶颈的地方是触发循环的“单条处理”习惯。我们一开始也是遍历所有活跃任务、逐个判断是否触发结果任务量到几千之后每秒只能处理约300次触发明显拖后腿。优化之后调度器改成批量预取模式每秒只做一次全局扫描把未来1秒内所有该触发的任务一次性算出来然后批量生成实例、批量写入MySQL、批量写入Redis Stream管道。吞吐量直接提升了十倍能够稳定跑到每秒3000次触发。另外必要的监控指标非常关键。调度系统至少要看调度器扫描耗时、触发循环延迟、Redis Stream积压数、Worker消费速率、实例终态分布、重试队列深度。这些指标不光要看图还要配上阈值告警。有一次Redis Cluster跨分片网络抖动pending队列积压了三万条图表上肉眼就能看到断崖式下跌这个监控帮我们节省了至少半小时定位时间。6.3 一个真实的“重复调度”事故复盘最后分享一次我们经历过的严重事故希望对你有参考价值。当时az调度器部署了三副本用Redis锁选主。某一天Redis主从切换导致Leader的锁提前过期新Leader选主成功开始触发任务。但旧Leader所在节点网络并完全没挂只是和Redis之间断了一会儿恢复后它以为锁还在自己手里继续触发。这场“双主”持续了6分钟产生了大量重复实例。核心原因就是只靠Redis锁做选主没有配合fencing机制。事故复盘后我们做的第一件事就是给每个实例唯一键加入Leader任期号第二件事是在任务触发入口增加MySQL防重表所有入队任务必须先insert ignore谁插入成功谁负责跑第三件事加强了对Redis锁过期事件的监控。这之后即便再发生类似选主抖动系统至多会产生几个会被丢弃的无效实例不会影响实际数据。7. 写在最后调度系统的工程感悟ax调度从立项到今天我最大的一点体会是调度系统的复杂度不是来自某个花哨的算法而是来自对“状态一致性”和“可观测性”的执念。状态不一致任务就会要么丢要么重复可观测性做不好出了问题你就只能对着数据库手动排查。如果你也在自建调度平台我的建议是第一步不要急着追求极致的性能先把任务实例状态机、唯一键、审计日志这三件套做扎实。它们不性感但都是救命的。性能反而是后面通过加Broker、优化扫描循环就能解决的。调度系统是为业务兜底的系统它稳定运行的时候没人会注意到但一个调度事故造成的经济损失往往比业务应用本身的故障更大。希望这篇内容对你有实际帮助我后面打算把实例级别的执行链路血缘也做进去让每一次跑批都能清楚看到它消费了哪些上游数据生成了哪些下游结果。