资讯动态

自研分布式调度引擎AX:从0到1构建高可用任务调度系统

发布时间:2026/9/26 13:07:49 来源:尧图企业网站定制
凌晨两点半手机被告警短信震醒。打开电脑一看订单同步任务堆了3.7万条没跑下游报表全空运营群已经炸了。排查发现根因很讽刺集群凌晨发布时重启了一台机器由于调度器把任务集中分发到了单个节点新节点还没注册完成其他节点又误以为任务已经被领取于是全部跳过。那一刻我就决定必须把公司这套用了三年的调度体系推倒重来搭一套自己能完全掌控的调度引擎内部代号就叫AX调度。这篇文章不聊虚的把我从0到1设计AX调度的全过程拆开来讲为什么没直接搬开源框架、核心模型怎么设计、稳定性怎么打延迟/重试/幂等/限流、怎么保证任务不丢、以及集群高可用怎么演进。如果你也在维护一套定时任务系统或者正准备自研调度平台这篇应该能帮你省掉不少弯路。1. 为什么我不直接拿开源框架而是动手搭了AX调度先说结论不是开源框架不行而是它们解决的是通用调度问题我们需要的却是贴合业务语义的调度。这两个诉求的差距比想象中大得多。1.1 Quartz的硬伤调度和执行搅在一起公司最早用的是Quartz。单机场景下Quartz其实很好用API简单、文档全、Trigger灵活。但一旦上集群问题就出来了Quartz集群模式是靠数据库行锁QRTZ_LOCKS表来协调多个节点抢任务每个调度周期都要经历一次加锁→查询→执行→释放的过程。我们压测过三节点Quartz集群单表千万级任务实例时触发延迟从秒级飙到十几秒。更麻烦的是Quartz支持分片吗不支持。支持DAG编排吗不支持。任务执行失败后想按业务维度重试你得自己在外面包一层。这还只是功能层面。运维层面更头疼Quartz没有管理界面看不到任务状态、下次触发时间、执行历史。出了问题只能去翻日志开发查一次问题要半天这谁受得了。1.2 XXL-JOB的取舍调度能力强但编排磕磕绊绊XXL-JOB确实好用管理台开箱即用调度和执行分离的设计思路也正确。我们当时试了一周发现日常的简单任务调度完全没有问题但到了复杂场景就力不从心了分片路由策略说白了就是轮询、随机、一致性哈希那几种想按业务语义自定义分片得改源码。任务编排靠任务链实现本质上是一条链想表达订单创建成功后才去扣库存库存扣完再发通知这种分支逻辑没法原生描述。执行器通过HTTP回调注册任务量上来后调度中心就成了瓶颈单机调度中心支撑上万任务时页面操作明显卡顿。XXL-JOB适合中小团队快速落地但对我们这种要接上百个下游系统、任务量每天千万级的场景定制成本反而比自研还高。1.3 Airflow太重为数据管道设计的调度器Airflow的DAG编排能力确实强但它是为数据Pipeline设计的跑Python脚本没问题想调度Java服务得包一层API。我们团队是Java栈为了调度任务再去维护一套Python基础设施多一个中间件就多一份凌晨三点被叫起来的风险。直接否掉。1.4 一个对比表格选型一目了然维度QuartzXXL-JOBAirflowAX调度自研秒级触发支持有延迟支持不支持支持时间轮任务分片不支持有限支持不支持三种分片模式DAG编排不支持链式强支持原生支持含条件分支集群高可用数据库锁调度中心集群多节点Leader选举防双跑管理界面无有有自研Web端定制成本低中中完全可控所以最终的路线定了基于成熟的时间轮组件做核心触发自研任务模型、编排语义、运维闭环。不重复造轮子但把业务要的那层调度语义牢牢抓在自己手里。2. AX调度的核心模型任务定义、调度规则与DAG编排调度系统最重要的不是调度本身而是建模。模型设计得好不好直接决定了后面所有的稳定性、可观测性、扩展性。这一节是我整个AX调度里最满意也最值得展开的部分。2.1 三个抽象分开TaskDef、TaskInstance、ScheduleRule很多自研调度系统上来就建一张task表把所有东西混在一起。这张表既要描述任务怎么定义又要记录每次执行情况还要保存调度规则。一旦字段多了状态一乱后面就彻底没法维护了。AX调度从第一天就拆成三个核心实体TaskDef任务定义描述做什么执行类型、执行参数、所属业务域、超时时间、重试策略。ScheduleRule调度规则描述什么时候触发Cron表达式或FixedDelay固定延迟、生效时间范围、是否启用。TaskInstance任务实例描述某一次触发后发生了什么执行状态、开始结束时间、执行结果、错误信息。为什么要拆开因为规则和定义的生命周期不一样。一个任务可能白天每小时跑一次大促期间改成每5分钟一次。如果把规则挂在任务定义上改规则就得停任务拆开之后直接新增一条ScheduleRule指向同一个TaskDef就行历史实例的归属也不受影响。更关键的是同一个TaskDef可以挂多条ScheduleRule天然支持灰度发布。比如先让10%的流量走新逻辑通过调整实例参数实现而不是复制一整个任务。2.2 调度规则的实现细节Cron和FixedDelay的取舍AX调度支持的调度规则有两类Cron表达式和FixedDelay固定延时。Cron用于绝对时间点触发的场景比如每天凌晨2点跑财务日切。我们内部只支持6位标准Cron秒 分 时 日 月 周不支持L、W这些特殊字符理由很简单——特殊字符的语义在不同解析库里实现不一致线上出过每月最后一个工作日在不同年份少跑一次的bug干脆禁用。FixedDelay用于相对时间循环触发的场景强调上一次执行结束后再间隔固定时长。比如订单同步任务我们不希望它每秒准点触发而是希望上一次跑完后等5秒再跑下一次。这是避免任务堆积的关键手段。给一段任务定义的Java代码示例public class OrderSyncTaskDef { // 任务标识 private String taskDefId order_sync_daily; // 执行器名称 private String executor order-center-executor; // 执行参数JSON字符串 private String params {\batchSize\:2000, \syncType\:\INCREMENT\}; // 超时时间毫秒 private int timeoutMs 30 * 60 * 1000; // 失败重试策略指数退避最多3次 private RetryStrategy retryStrategy new RetryStrategy() .maxAttempts(3) .backoffBaseMs(1000) .backoffMultiplier(2.0); }对应的时间轮调度逻辑核心就是拿到一批到期的ScheduleRule根据规则类型生成TaskInstance然后推给执行器。这里有一个被我反复强调的点执行器只认TaskInstance不认ScheduleRule。调度器负责什么时候该生成实例执行器负责一个实例该怎么执行完两者通过数据库解耦互不干扰。2.3 DAG编排从各自为战到流程化单任务调度做完了真正的难点在编排。我们的业务场景里订单对账是一个典型的DAG订单数据同步完成后要先做数据校验校验通过才做记账记账成功才发通知校验失败直接走人工补偿通道。AX调度的DAG模型不复杂TaskDef之间通过边Edge建立依赖关系每条边带一个前置条件默认是SUCCESS也可以配成FAILED失败后触发告警任务或者COMPLETED无论成功失败都要执行下游。JSON结构大概是这样的{ edges: [ { from: order_sync, to: order_validate, on: SUCCESS }, { from: order_validate, to: order_post, on: SUCCESS }, { from: order_validate, to: order_compensate, on: FAILED }, { from: order_post, to: order_notify, on: SUCCESS } ] }调度器里跑的是拓扑排序。每个TaskInstance执行完成时更新自己的状态然后沿着出边找到所有下游检查前置依赖是否全部满足满足就把下游的TaskInstance创建出来状态PENDING交给执行队列。这样天然支持分支、多入度、多出度不用为每种流程单独写代码。2.4 表结构设计一张实例表扛住了千万级写入再分享一组核心DDL这些字段都是经过线上打磨的CREATE TABLE t_task_instance ( id BIGINT AUTO_INCREMENT PRIMARY KEY, task_def_id VARCHAR(64) NOT NULL, schedule_rule_id VARCHAR(64) NOT NULL, biz_key VARCHAR(128) DEFAULT NULL, status TINYINT NOT NULL COMMENT 0:待执行 1:执行中 2:成功 3:失败 4:取消, trigger_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, start_time DATETIME NULL, end_time DATETIME NULL, retry_count INT DEFAULT 0, executor_ip VARCHAR(32) DEFAULT NULL, error_msg TEXT, UNIQUE KEY uk_def_trigger (task_def_id, trigger_time), KEY idx_status_trigger (status, trigger_time) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;特别注意uk_def_trigger这条唯一索引这是防重复调度的第一道防线。同一任务的同一触发时间在数据库层面就不可能存在两条实例。这是后面防双跑、防重试的关键基础。3. 稳定性实战延迟、重试、幂等、限流这条链路怎么打通调度系统上线后才是真正噩梦的开始。前三个月我们踩了无数坑下面这几条是最典型的也是我认为所有自研调度平台都得过的坎。3.1 一次调优把派发和执行彻底拆开初期版本里时间轮线程触发任务后直接在当前线程里同步执行。听起来没毛病但线上跑起来就出事了一个跑批任务卡了十分钟后续几百个任务的触发全部阻塞延迟越来越严重最后连心跳都过期了被误判为节点宕机。复盘后立刻改架构把整个链路拆成两条线程池Dispatch线程池只负责从时间轮取到期的触发事件、生成TaskInstance、写入数据库。要求快进快出绝对不能在这里执行任务逻辑。Execute线程池消费PENDING状态的实例真正跑任务代码执行完把结果写回实例。拆开之后调度延迟直接从秒级降到毫秒级。哪怕某个执行任务卡死了最多消耗Execute线程池的线程Dispatch线程完全不受影响。这也是AX调度后来能支撑千万级日活调度的基石。3.2 重试不能无脑重试指数退避才是正解有一次下游订单系统发版短时间不可用结果我们的任务全部在1秒内连续重试直接把下游数据库连接池打满了。那次事故之后我定了一条铁律默认不做立即重试一律指数退避。现在的重试策略是最大重试次数默认3次可配置。退避基数初始1秒每次乘2也就是1s、2s、4s。超过最大次数后实例进入FAILED状态同时触发告警链。每次重试都会生成新的retry_count写回实例表方便事后审计。还有一个容易被忽略的点重试前一定要检查幂等键。为什么因为第一次执行可能是业务成功但响应超时如果不检查幂等键直接重试就会造成重复扣款、重复入账。AX调度在每个执行器里强制要求实现idempotentCheck方法重试前先查一下biz_key是否已经成功。3.3 幂等设计状态机是调度系统的灵魂说到幂等就必须讲AX调度的状态机设计。TaskInstance的一生只允许走下面这条路径PENDING - RUNNING - SUCCESS/FAILED/CANCELED不允许从SUCCESS再回到PENDING不允许重复从PENDING推进到RUNNING——所有状态变更都通过SQL的UPDATE ... WHERE status ?条件更新实现。如果影响行数为0说明状态已被其他线程变更当前操作直接放弃。举个例子UPDATE t_task_instance SET status 1, start_time NOW(), executor_ip ? WHERE id ? AND status 0;这样即使两个节点同时消费同一条实例也只有一个节点能抢到另一个更新0行直接跳过。这是防重复执行的基础也是性能优化的关键——我们在DB层面解决并发问题不去依赖复杂的分布式锁。3.4 限流不让任务打垮下游系统任务多了以后下游系统开始抱怨你们订单同步任务一跑我们数据库CPU直接100%。于是AX调度加了双层限流执行器本地限流每个执行器最多同时跑N个实例超出就排队等待不新建线程。全局信号量限流跨节点的总并发限制比如订单同步任务全局最多20个并发。全局信号量用Redis Lua脚本实现核心逻辑是这样一个原子操作local current redis.call(INCR, KEYS[1]) if current tonumber(ARGV[1]) then redis.call(DECR, KEYS[1]) return 0 end redis.call(EXPIRE, KEYS[1], ARGV[2]) return 1任务完成后执行DECR归还信号量。这套方案比数据库行锁轻量得多也比纯本地限流准确得多。3.5 告警降噪不是每个失败都要半夜叫醒你告警轰炸是我们踩过的另一个坑。最初版本只要任务失败就发钉钉结果三天后大家直接把机器人静音了。后来的方案是连续失败3次才产生告警事件。同一个任务在15分钟内的告警只推送一次后续失败合并统计。告警带上任务名称、执行器IP、失败次数、最近一次错误信息、任务管理页链接。从此以后半夜被叫醒的次数少了一个数量级而且每次被叫醒都是真需要处理的事。4. 任务不丢的兜底机制时间轮数据库扫描的双保险调度系统最可怕的事不是跑得慢而是该跑的任务没跑。由于时间轮的数据在内存里JVM一崩溃所有在途的触发点全部丢失。所以AX调度设计了双保险机制。4.1 第一道保险数据库扫描兜底每个调度节点上除了时间轮主线程在跑还会启动一个兜底扫描线程。每隔60秒扫描一次t_task_instance里所有状态为PENDING、当前时间比计划时间晚5分钟以上的实例重新推入执行队列。为什么是5分钟因为正常任务的触发到执行启动有延迟三五秒很正常但如果超过5分钟还没动静大概率是调度线程卡死或者节点假死了。这个阈值可以根据业务调但我们的经验是别太激进否则会出现重复消费。扫描脚本长这样SELECT id, task_def_id, executor_ip, params FROM t_task_instance WHERE status 0 AND trigger_time DATE_SUB(NOW(), INTERVAL 5 MINUTE) AND lock_time IS NULL LIMIT 500;拿到这批实例后尝试用UPDATE ... WHERE status 0 AND id ?去抢占抢到就执行。整个过程对主调度线程零侵入相当于一个安全网。4.2 第二道保险节点重启恢复节点重启后第一件事不是直接开始调度而是先执行一个StartupScanner。这个Scanner负责找出所有RUNNING状态但超过心跳阈值的实例按策略处理如果上游业务接口支持幂等直接把实例重新入队置为PENDING。如果不支持幂等置为FAILED触发人工补偿流程。为什么必须这么做因为节点重启时内存时间轮里的触发点全没了如果不扫描数据库把未完成任务捞回来这些任务就会永远消失直到业务对不上账才发现。4.3 时钟漂移这个隐雷调度分布式系统里时钟漂移是一个经常被忽略的大坑。我们遇到过一台调度节点的系统时间慢了1分半结果任务延迟触发下游等了半天没数据我们排查了两个小时才发现是NTP没同步。后来的做法是所有调度节点强制NTP同步部署脚本里检查ntpq -p的偏移量超过500ms直接拒绝该节点注册。代码里做偏移窗口容忍。比如判断实例是否超时时允许2秒的误差。避免因为毫秒级的时钟抖动导致误判断。时间轮本身不依赖数据库事务但触发计划时间戳一定要用统一的数据库时间或NTP校准后的本地时间不能两台机器各玩各的。5. 集群模式下AX调度的高可用演进选举、分片与防双跑单节点调度始终有单点风险AX调度上线第二个月就开始了集群化改造。这一节的内容很多是从线上故障里逼出来的。5.1 Leader选举用etcd还是ZooKeeper我们最终选了etcd。原因很简单etcd的Lease机制天然适合做Leader选举和自动续期而且Watch接口能实时感知Leader失联比ZooKeeper需要维护会话的复杂度低不少。另外在Kubernetes环境里部署etcd更顺滑。基本原理所有调度节点启动时尝试创建/ax-scheduler/leader这个Key带上自己的节点IDTTL10秒。创建成功的节点成为Leader启动时间轮。其他节点进入Standby模式Watch这个Key一旦Key过期Leader宕机Standby节点立即竞争创建新Leader接管时间轮。Leader失联到新Leader接管我们的实测数据是2~3秒左右。这段时间内可能会有任务延迟触发但因为有数据库扫描兜底不会丢任务。给一段etcd选举伪代码// 尝试创建 leader key try { etcdClient.put(/ax-scheduler/leader, nodeId, leaseClient.newLease(10)); isLeader true; } catch (KeyAlreadyExistsException e) { isLeader false; watchLeader(); }5.2 防双跑只靠分布式锁远远不够这是整个AX调度里最经典的一个坑值得单独拿出来讲。当时我们用Redis分布式锁保证同一时间只有一个Leader在调度。结果凌晨发生了一次故障旧Leader因为GC停顿超过锁过期时间Redis把锁释放了新Leader抢到锁开始调度此时旧Leader GC恢复但它的线程还在继续往下执行——两个Leader同时跑任务大量重复执行下游被重复入账搞了一整天。要彻底解决这个问题光有分布式锁是不够的得加Fencing Token栅栏令牌。etcd天然支持这个机制每次创建Key时拿到当时的Revision号这个Revision就是一个单调递增的令牌。旧Leader在每次执行任务前先检查自己持有的Revision是否还是当前Key的Revision。如果不是说明自己已经被踢出局立即终止调度。// 每次调度循环开始前检查 long currentRevision etcdClient.get(/ax-scheduler/leader).getRevision(); if (currentRevision ! myRevision) { // 我已被替换退出调度 isLeader false; return; }这才是一个闭环的高可用设计分布式锁解决的是谁有资格Fencing Token解决的是谁真的有资格。两个缺一不可。5.3 分片从单Leader调度到多Leader分片调度单Leader虽然高可用但所有触发事件都压在Leader节点上到了大促高峰还是会有压力。AX调度的第三个版本引入了分片机制。三种分片模式手动分片按租户或业务域指定某些任务只在特定节点上调度。适合敏感业务或低频任务。一致性哈希分片按TaskDefId哈希均匀分布到所有节点。默认模式适合大部分场景。动态分片节点负载过高时自动把部分任务迁移到低负载节点。实现复杂我们目前还在打磨。一致性哈希分片下每个节点的时间轮只承载部分TaskDefId的触发事件。这样既保留了Leader选举的高可用底子又避免了单点性能瓶颈。5.4 性能表现这组数据让我松了口气最后一个版本上线后我们做了压测结果比预期好指标数值单节点时间轮触发能力约3000个触发事件/秒实例表写入单实例约1.2万TPS数据库瓶颈Leader故障切换耗时2~3秒任务延迟P50小于100ms任务延迟P99小于500ms瓶颈最终落在数据库写入上而不在调度本身。如果未来实例量再上几个量级可以按业务域拆库把实例表分片。最后再说一点实操体会AX调度从立项到稳定运行前后迭代了四个大版本。如果让我重新做一次我会在第一天就把任务血缘视图TaskDef之间的依赖关系图做出来——排查问题时看着一张图比翻十个表效率高十倍。另一个深刻的体会是调度系统的核心价值不在触发机制有多华丽而在运维体验。一个任务失败了你能不能3分钟内定位到原因一个任务重复跑了你能不能自动止损。这些才是每天真刀真枪要面对的。希望这篇AX调度的拆解能给你带来一些启发。如果你也在折腾调度系统欢迎在评论区聊聊你们是怎么解决防双跑和任务编排的——这两块我始终觉得还有很大的优化空间。

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑