☰
AX调度实战:从零构建分布式任务调度与延迟队列内核
2026/9/26 8:21:36 网站建设 项目流程

最近圈子里聊得很热的一个新词,叫“ax调度”。坦白说,第一次听到这个词的时候,我以为是某个开源框架的简称,后来才发现它更像一类调度范式的代称——把海量异步任务、延迟消息、定时触达统一纳入一套可控的调度内核里,让业务方不用再各自为战地写定时循环和延迟队列。我在过去一年多里,因为项目需要,从零自研了一套任务调度器,内部代号就是“ax”,后来团队里干脆把这种调度方式统称为ax调度。这套东西解决了我这边最头疼的问题:任务量一大,要么堆积,要么重复执行,要么机器一宕就全盘乱套。这篇文章把我自己的设计思路、踩坑经历、和落地时的关键参数都摊开来说,适合正在搞任务调度、打算自研延迟队列或者维护高并发触发系统的后端同学参考。

1. 为什么需要AX调度:从业务痛点说起

1.1 裸写定时任务的问题

我最早参与的一个业务系统,要处理的动作很典型:用户下单后30分钟未支付自动关单、订单完成后48小时自动确认、每天的凌晨2点跑一次报表汇总、还有运营后台手动触发的一堆异步数据同步。最初大家图省事,直接在服务里起了几个Timer和ScheduledExecutorService,再配合数据库里的status字段轮询扫表。

扫表这个方案在任务量几百、几千的时候还能勉强跑,但一旦单表数据到千万级,问题就全冒出来了。最直接的是慢查询拖垮主库,每30秒扫一次create_time < now()-30min AND status='WAIT_PAY',索引稍微建得不对,锁和IO开销能把数据库CPU打到70%以上。更严重的是,如果同一个订单在多个实例上同时被扫到,两个节点都发起了关单操作,又没有做幂等控制,就会出现重复发短信甚至重复退款的风险。

还有一个容易被忽略的问题:ScheduledExecutorService的线程池是常驻的,每个定时任务至少一个线程,几十个任务就是几十个常驻线程,内存和上下文切换开销都不小。而且一旦目标系统响应变慢,线程池里的任务就会排队,后面的定时任务全部延迟。我们有一次线上事故就是因为外部物流接口超时,把执行线程全部占满,导致每天凌晨的汇总报表没跑出来,早上业务方在用数据的时候才发现是空的。

1.2 AX调度想解决的三个核心矛盾

正是这些经历,让我重新审视“调度”这件事。传统的定时任务框架通常解决了“什么时候跑”的问题,但很少解决“跑在哪台机器上”、“挂了怎么办”、“如何避免重复跑”和“任务积压如何消化”的问题。我给自己定的目标,是做出一个叫ax的调度内核,专门解决三个最核心的矛盾。

第一个矛盾是海量任务的存储成本。如果还是把所有任务写进一张MySQL表,那调度系统本身就成了瓶颈。我决定采用分层存储:短周期任务放在内存时间轮,长周期任务放在数据库但只存下次触发时间,配合批量预取。

第二个矛盾是分布式环境下的竞争条件。多个worker节点同时消费同一个任务,必须有租约机制或者强协调能力。我之前用过基于数据库悲观锁的方案,效果奇差,因为每个任务抢占都要SELECT ... FOR UPDATE,数据库扛不住。ax采用Redis分布式锁加本地缓存双保险,锁粒度精确到任务ID,抢锁失败的节点直接放弃该任务,不反复重试。

第三个矛盾是调度的准时性和可靠性。准时性要求时间精度,可靠性要求任务不能因为宕机而丢失。这两者天然有冲突。ax的做法是:把时间精度要求高的任务放到内存时间轮,把可靠性要求高的任务放到持久化队列,再通过状态补偿机制保证最终一致性。

2. 整体架构设计与调度模型

2.1 任务分片与一致性哈希

ax不是一个单机程序,而是一组进程构成的调度集群。为了避免所有任务集中在某一台机器上,我做了一层任务分片。每个任务有唯一的taskId,通过一致性哈希映射到固定的分片,每个分片归属一个worker节点。这样的话,任意一个节点收到调度信号,都可以根据taskId快速计算出它应该由谁执行,然后走RPC转发。

一致性哈希的环形空间我用了256个虚拟节点,每个物理节点映射32个虚拟节点。这样做的直接好处是集群扩缩容时,只有约 1/虚拟节点数 的任务发生迁移,而不是全局重新洗牌。以前用取模算法,加一台机器会有一大半任务换节点,那叫一个酸爽。换成一致性哈希后,扩容时只需要对移走的虚拟节点对应任务做一次重新分配,线上几乎无感。

分片的另一个作用是流量隔离。我可以把不同业务线的任务打上独立的bizTag,然后配置分片权重。比如A业务线的任务量占80%,就分给它更多的虚拟节点。实现上就是虚拟节点命名时加上bizTag前缀,路由时先按bizTag过滤,再哈希到对应的虚拟节点集合。

2.2 时间轮与延迟队列的设计取舍

ax调度的时间源头是时间轮。我用Netty的HashedWheelTimer做过原型,但发现它在任务量过万、单轮精度毫秒级时性能下降明显,所以后来自己用数组+双向链表实现了一个分层时间轮。底层是一个环形数组,长度默认512格,每格代表1秒,一圈8分多钟。超出当前圈的延迟任务会放入一个上层时间轮,每上一层,时间跨度乘以60。

这里有个关键取舍:时间轮适合“短延迟、高吞吐、允许少量重启丢失”的场景。比如短信验证码发送,延迟30秒,机器重启丢几个任务是可以接受的,因为用户会重新触发。但订单关停这种绝对不能丢。所以ax把任务分成两类:MEMORY类型只挂在时间轮上,PERSISTENT类型除了挂在时间轮上,还写一份到持久化日志。

时间轮触发后,任务并不会直接执行业务逻辑,而是进入一个内部延迟队列。这个队列不是JDK的DelayQueue,因为那个队列的take()只有一个消费者,吞吐上不去。我改成了一批有序桶,每个桶对应一个消费线程。任务触发后根据分片hash选择一个桶放入,桶内按执行时间排序,消费线程轮询桶头判断是否到期。这样可以把消费并发度提高到机器核数水平。

2.3 状态机与任务生命周期

调度系统最怕状态混乱,所以我给ax定义了一套严格的任务状态机。

  • CREATED:任务已创建,尚未加入调度环。
  • SCHEDULED:任务已写入时间轮或持久化日志,等待触发。
  • DISPATCHED:已下发到执行节点,正在等待ACK。
  • EXECUTING:执行节点确认接收并开始执行。
  • SUCCESS:执行成功,记录完成时间。
  • RETRY:执行失败,等待下一轮重试。
  • DEAD:重试次数超过阈值,进入死信队列。

任务状态流转只允许单向,不允许跳级。比如EXECUTING不能直接回到SCHEDULED,必须经过RETRY再重新调度。这样有一个好处:任何时刻我们都能回答“这个任务到底处于什么阶段”。排查问题时,只要看状态机就知道卡在哪一环。

状态存储我用的是Redis Hash,key是task:{taskId},field是state和updateTime。为什么不用数据库?因为调度路径上状态变化非常频繁,写数据库的IOPS太高。但Redis数据不能作为唯一事实源,所以周期性会把状态快照异步落库,供对账和审计使用。

3. 核心实现:关键代码与踩坑记录

3.1 调度器主循环

ax调度器的核心线程是一个跑在主节点上的循环,类似事件循环。它每一轮做三件事:收集当前时刻到期的任务、更新统计指标、执行一次持久化对账。

用伪代码描述主循环:

while (running) { long now = System.currentTimeMillis(); // 1. 从内存时间轮取出当前格子所有到期任务 List<Task> readyTasks = wheel.pop(now); // 2. 将到期任务写进持久化日志(PENDING_LOG) for (Task task : readyTasks) { if (task.isPersistent()) { writeLog(task, "READY"); } } // 3. 发布Ready事件到分发器 taskDispatcher.dispatchAll(readyTasks); // 4. 处理持久化日志中超过超时时间未ACK的任务 recoverTimeoutTasks(); // 5. 等待下一个tick Thread.sleep(wheel.tickDuration()); }

这里有一个非常容易踩的坑:Thread.sleep(tickDuration)是不准的,它只保证至少睡眠这么多时间,但JVM的GC暂停和系统负载都可能导致tick延迟。所以主循环里不要用sleep来算下一轮时间,而应该用System.nanoTime()计算硬实时截止时间,然后根据剩余时间决定是否sleep以及sleep多久。后来我干脆改成LockSupport.parkNanos(),情况好了很多。

另一个坑是主循环必须和任务分发解耦。之前我把任务执行也放在主循环里,结果一个任务的外部调用超时,整个时间轮都卡住。后来严格约束主循环只做触达和投递,真正执行逻辑放到Worker线程池中。

3.2 任务分发与ACK机制

任务从时间轮弹出之后,要交给正确的Worker节点。ax从主节点视角看,是一个分发器;从Worker视角看,是一个消费者。它们之间用gRPC长连接双向通信。

分发流程:

  1. 主节点根据taskId哈希找到所属分片和Worker节点。
  2. 主节点推送调度请求DispatchRequest给Worker,请求里带taskId, executeAt, bizTag, payload。
  3. Worker收到请求后,执行前先返回一个ACK_RECEIVED给主节点。
  4. Worker执行完毕,返回SUCCESS或FAILED,附带执行耗时和错误信息。
  5. 主节点收到ACK后更新任务状态。

这里最核心的是ACK超时机制。如果主节点发出Dispatch请求后,超过3秒没有收到ACK_RECEIVED,会认为节点可能不可用,把任务重新放入时间轮并标记为RETRY。如果收到ACK_RECEIVED但超过30秒没有收到执行结果,则视为执行超时,同样进入重试流程。

但是这里就会产生一个重复执行风险:Worker可能已经执行完业务逻辑,但结果回传时网络抖动,主节点判超时,重试后又执行了一次。所以ax强制要求业务方实现幂等,我在文档里专门强调:retryTimes和taskId会一并传给业务方法,业务方可以用taskId作为幂等键。我自己在做订单关单时,就是先通过taskId查一下order_action_log,存在相同taskId的记录就直接返回。

3.3 持久化与恢复:宕机不丢任务

持久化这一块,我尝试过写MySQL、写本地文件、写消息队列三种方案。最终方案是混合:每个节点把任务意图追加写到本地磁盘的WAL(Write-Ahead Log),同时周期性地把WAL摘要上报给协调节点。

WAL每行是一个JSON:{"taskId": "xxx", "executeAt": 1710000000000, "bizTag": "order", "payload": "{\"orderId\":123}"}。任务在创建时写入WAL,执行完成后写入一条KEEPALIVE标记,表示该任务已完成。恢复的时候扫描WAL,把没有KEEPALIVE且执行时间未过的任务重新加入时间轮。

这里有一个关键设计,WAL文件必须追加写且定期滚动。我用的是每64MB滚动一次,旧文件压缩后归档到对象存储。不然WAL无限膨胀,恢复扫描会越来越慢。

在实际测试中,一台节点被kill -9后,重启最多需要扫描最近64MB的WAL,换算下来大约几万条任务,恢复时间控制在5秒内。这个恢复机制曾经帮我避免过一次生产事故:当时某Worker节点磁盘满,进程被OOM Killer杀掉,重启后靠WAL把该节点负责的所有未完成任务重新拉了起来,业务方只观察到几分钟延迟,没有发生订单漏关。

4. 实操经验:AX调度在生产环境落地

4.1 部署拓扑与参数选型

ax集群我推荐按三层部署:调度主节点(Scheduler)、工作节点(Worker)、依赖存储(Redis + 对象存储)。

调度主节点至少2个,一主一备,用选主协议确认谁是当前主节点。主节点负责跑时间轮和分发,备节点只接流量不上岗。主节点建议4核8G起步,因为内存时间轮要保留大量任务。工作节点根据业务量横向扩展,我目前压测过8台Worker节点,日常任务量在每秒2000次触发时CPU利用率只有30%。

Redis的角色很微妙。它既要存任务状态,又要锁任务。建议单独为ax分配一个Redis实例,不要和业务缓存混在一起,否则一个大数据查询把Redis搞慢,调度锁的获取时间就会飙升。锁的超时时间我设置为10秒,任务租约续期是每3秒做一次。续期用Lua脚本保证原子性:

if redis.call("GET", KEYS[1]) == ARGV[1] then return redis.call("PEXPIRE", KEYS[1], ARGV[2]) else return 0 end

这个脚本解决了之前一个很经典的问题:锁到期而任务还在执行,另一个节点抢到锁后重复执行。续期机制保证只有在持有锁的节点失活时,锁才会真正过期被别人抢走。

4.2 容量评估与性能压测

关于容量,我先给个量级感觉:ax的调度能力瓶颈主要在Redis的写并发和Worker的执行能力,而不是时间轮本身。我压测的时候,时间轮每秒能吞吐5万次触发,但一旦加上持久化日志和Redis状态更新,单主节点稳定在2万/秒左右。如果你的业务需要超过这个量级,考虑分主节点按业务线隔离。

压测的时候要重点关注三个指标:

  • dispatch_success_rate:分发成功率,正常应接近100%。
  • ack_p99:从分发到确认ACK的耗时,我这边p99在40ms以内,超过100ms说明网络或线程池有瓶颈。
  • retry_rate:重试率,控制在1%以内。如果超过5%,说明超时时间配置不合理或Worker有热点。

还要注意一个容易被忽略的参数:时间轮的格数。格数太小会导致同一格内任务过多,弹出时瞬间产生大流量。我默认512格,每格1秒,经过估算足以应对每秒2000任务量。如果你的任务量更大,可以调大格数,但代价是内存上升。每格是一个双向链表头节点,512格的内存开销可以忽略。

4.3 监控告警与日常运维

调度系统的监控最好直接落Prometheus指标。我暴露了这几个核心指标:

指标名含义告警阈值
ax_task_scheduled_total累计调度次数无(趋势监控)
ax_task_dispatch_failure_total分发失败次数5分钟内增长超过10
ax_task_ack_latency_millisACK耗时分布p99 > 200ms
ax_worker_execute_time_secondsWorker执行耗时分布p99 > 30s
ax_task_dead_total死信任务数大于0即告警
ax_redis_lock_acquire_failure_total获取锁失败次数每分钟超过20

日常运维中最重要的一个命令是查看任务积压。ax提供了一套Admin HTTP接口,/tasks/pending可以查当前积压任务数。积压的原因通常是Worker执行太慢导致消费能力下降,或者某个业务方接口被拖慢。我这边有一次积压积了上百万,是因为下游订单服务故障,所有关单任务都在等待重试。排查思路是先看ax_worker_execute_time_seconds,如果执行耗时很高,就去查具体任务调用的接口;如果执行耗时正常但积压仍高,就扩容Worker线程池。

5. 常见问题与排查技巧

5.1 任务积压与消费倾斜

任务积压最常见的原因是消费倾斜。一致性哈希本来能均匀分配,但如果某个Worker节点宕机,它的虚拟节点被重新分配给其他节点,短时间内可能造成某个Worker接受的流量翻倍。此时测下来该Worker的线程池队列被打满,往外分发的响应变慢。

排查倾斜时,先看每台Worker的ax_worker_received_total指标曲线。正常情况下应该是平均分布,如果有单台明显高于其他节点,就检查它是不是新加入的节点,或者它的虚拟节点是否重复。我曾经在一次误操作里给同一个Worker重复注册了两次,导致它承担了一半以上的流量。修复后我在注册逻辑里加了机器ID字段,防止重复注册。

如果能确认倾斜来自节点故障,但不方便马上扩展节点,可以先临时调大倾斜节点的线程池,同时降低该节点锁续期的间隔,让它更快处理任务。但这只是应急,最终还是要均衡分片。

5.2 重复执行与幂等设计

ax调度没法保证只执行一次,这点必须说透。它的保证是至少执行一次,所以业务侧必须承担幂等。我在实际项目中用的幂等方案比较简单:每个任务在业务库里都有一条task_execution记录,字段包括task_id、status、update_time,唯一索引为task_id。

执行流程变成:

INSERT INTO task_execution (task_id, status, update_time) VALUES (?, 'RUNNING', NOW()) ON DUPLICATE KEY UPDATE status = CASE WHEN status = 'SUCCESS' THEN 'SUCCESS' ELSE 'RUNNING' END, update_time = NOW();

如果之前已经成功,这条SQL会直接返回,后续业务逻辑就不执行了。但这里还有一个并发间隙:两个请求同时执行这个INSERT,第一个成功,第二个会陷入阻塞,需要给SQL设置锁等待超时,不然会拖住任务线程。我实际用的连接池里lock_timeout设为1秒。

5.3 时钟跳跃问题

调度系统最怕的事之一就是时钟跳跃。无论是NTP同步导致的往回跳,还是手动改时间,只要系统时间往回跳了几秒,时间轮就会混乱。我之前踩过一个坑:运维在凌晨同步时间时往回跳了1秒,结果所有本应在下一格触发的任务全部提前了999毫秒触发,导致业务方收到一堆莫名其妙的消息。

解决办法是把系统时钟和单调时钟分离。时间轮的格子推进必须用System.nanoTime(),因为它是单调递增的,不受系统时间调整影响;但任务触发后的业务判断,比如“现在是否到了执行时间”,可以用System.currentTimeMillis()。我封装了一个ClockHolder,调度器内部全部用单调时钟,只有对外暴露时间戳时才转成系统时间。

同时,我给ax加了一个时间跳跃检测器:每秒计算当前系统时间与单调时钟导出的预测系统时间的偏差,如果偏差超过500ms,就暂停调度并报警。暂停时间持续2秒,等系统时间稳定后再恢复。这个方法救过我一次,不然那次的线上事故可能直接毁了整个业务日的对账数据。

最后再说一点关于AX调度的体会

自研调度器这件事,表面看是技术选型和代码实现,实际做起来最耗时的是对“不确定性”的处理。任务什么时候触发、哪个节点执行、执行到一半挂了怎么办、重复执行怎么兜底——这些才是调度的灵魂。我看到很多人一上来就追求高并发、低延迟,但忽略了恢复和幂等,结果后期天天在补bug。ax调度这套体系我打磨了挺久,真正让它稳定下来的是每次故障后的复盘和补丁。如果你也在做类似的调度系统,我建议你从“最少丢失、最多执行的边界”开始设计,先把失败路径走通,再考虑优化性能。这样系统上线后,你睡觉都能踏实一点。

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

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

立即咨询