有些项目一开始并不起眼,等真正把它做起来,你才发现它重新定义了整个团队对“稳定性”的理解。“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点会有上千个任务同时触发。如果调度器直接压给Worker,Worker瞬时负载可能飙升,而且任务在途状态很难持久化。有了Broker之后,调度器只负责快速入队,Worker根据自己的消费能力去拉任务,天然就实现了流量整形。曾有人问,调度器直接写MySQL,Worker轮询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可以定义shardCount=100,调度器触发时会一次性生成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 -> SUCCESS,FAILED则根据重试配置回到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完成之后才能跑B,B和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任期号、唯一索引冲突日志 |
| 任务卡在RUNNING | Worker被杀、执行器僵死 | 看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、优化扫描循环就能解决的。调度系统是为业务兜底的系统,它稳定运行的时候没人会注意到,但一个调度事故造成的经济损失,往往比业务应用本身的故障更大。希望这篇内容对你有实际帮助,我后面打算把实例级别的执行链路血缘也做进去,让每一次跑批都能清楚看到它消费了哪些上游数据,生成了哪些下游结果。