ARTICLE DETAIL

资讯详情

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

轻量级异步任务调度器设计:从单机延迟队列到分布式演进

轻量级异步任务调度器设计:从单机延迟队列到分布式演进 凌晨两点半被电话叫醒线上告警提示订单超时任务大量积压。当时我负责的业务系统里有一个很狼狈的模块订单15分钟未支付要自动关闭用户签到后的优惠券要延迟生效运营要批量推送活动消息。这些功能各有各的执行时间但实现方式几乎是同一个套路——起一个定时线程每隔几秒扫一次数据库把到期数据捞出来再同步执行。结果就是设计粗糙、资源浪费一到高峰期所有任务全挤在一起把数据库搞得半死。折腾了半个月我干脆自己写了一个轻量级的异步调度模块代号就叫“ax”全称是 Async eXecution也就是异步执行。这篇文章会把 ax 调度器的设计思路、核心代码、踩坑过程、压测数据以及从单机走向集群的演进路线完整记录下来。适合那些需要在业务系统里引入异步任务调度但又不愿意直接上 Quartz、XXL-Job 这类重量级框架的团队参考。如果你也在为“任务到点该执行却没人执行”或者“一压测就线程爆炸”这类问题头疼这篇应该对你有实际帮助。1. 为什么要自己造一个“ax”——不是轮子是被现有方案卡住了1.1 业务场景里的真实痛点当时的业务里主要有三类典型任务。第一类是超时订单关闭用户下单后15分钟未支付就要自动置为已关闭量级跟着订单量走高峰期每秒能产生上千个待关闭订单。第二类是延迟生效类任务比如优惠券领取后要等30分钟才可用、会员体验包到期后自动降级。第三类是批量的补偿对账任务一般在凌晨统一执行要求任务之间有顺序保障失败之后还要能原样重跑。这三类任务有个共同特征执行时间不确定不能阻塞主流程而且失败了不能静默丢弃。用“定时线程 轮询数据库”的老办法做会出现两个很直接的恶果。第一数据库被高频扫描明明只有几千个到期任务却要让全表索引每秒被扫几十次数据库负载完全被无意义的 IO 拖垮。第二任务执行之间没有任何隔离一个慢 SQL 能把整批任务全堵住晚点的订单跟着晚点用户投诉一排排地来。1.2 为什么不用现成的调度框架很多人第一反应是“用现成框架不就行了”。我当时也认真评估过最后还是决定自己写。先看 JDK 自带的 ScheduledExecutorService它只能做简单的延迟执行和周期执行不支持优先级、没有任务状态、没有失败重试更别谈执行超时的控制。再看 Quartz功能确实很完善分布式集群、持久化、Cron 表达式都有但代价是配置重、依赖多、概念多一个小团队维护它本身就要花不少精力。XXL-Job 这种中心化调度平台更不用说了至少得部署调度中心和执行器两个端还要配数据库和前端控制台。对一个只有十几人维护的中小型业务系统来说这完全是杀鸡用牛刀。我真正需要的其实是一个足够轻量、能像库一样直接嵌进业务代码里的调度器。所以 ax 的目标从一开始就很明确核心代码控制在几百行以内零外部依赖延迟触发、优先级、超时、重试、状态追踪全都要有。这也是后来团队里很多同学愿意直接接它的原因——不用部署任何服务加一个类进来就能用。1.3 ax 的设计目标与原则总结下来ax 要满足六个基本目标轻量无框架依赖核心代码不超过一千行可以嵌进任何 Java 模块。异步非阻塞提交任务后调用方立刻返回业务主流程不受影响。支持延迟与优先级任务可以指定什么时候执行同一时刻多个任务时可以按优先级竞争 worker。状态可观测每个任务从提交到结束有哪些状态能被追踪方便排查问题。优雅关闭应用重启或发版时能销毁正在执行的任务不丢任务或者至少能拿到剩余任务清单。单机支撑万级排队内存队列要能支撑单实例同时排队一万个以上任务且延迟抖动可控。这里有一个核心设计原则很大程度影响了后面的代码结构调度线程只负责把任务“催熟”绝不在调度线程里直接执行业务逻辑。也就是说判断任务是否到期、到期后交给谁执行这两件事必须由不同角色负责。否则一旦某个任务执行了 5 秒钟后续所有到期任务都会被堵住调度器就退化成一根串行管道了。2. ax 调度的核心语义任务抽象与执行模型2.1 任务对象怎么设计才够用设计任务对象是第一步也是容易被低估的一步。很多开发者在写自己的任务调度器时只知道往队列里塞一个 Runnable但 Runnable 信息量太少后面做状态追踪和超时控制会非常痛苦。我在 ax 里用的是自定义的 Task 对象一个泛型任务类核心字段包括 taskId、bizType、payload 业务数据、priority 优先级、executeAt 计划执行时间戳、timeoutMs 超时时间、maxRetry 最大重试次数以及当前状态 status。public class TaskT { private final String taskId; private final String bizType; private final T payload; private final long executeAt; // 计划执行的时间戳毫秒 private final int priority; // 数值越小优先级越高 private final long timeoutMs; // 单次执行超时时间 private final int maxRetry; // 最大重试次数 private volatile int status; // 任务状态 public Task(String taskId, String bizType, T payload, long executeAt, int priority, long timeoutMs, int maxRetry) { this.taskId taskId; this.bizType bizType; this.payload payload; this.executeAt executeAt; this.priority priority; this.timeoutMs timeoutMs; this.maxRetry maxRetry; this.status STATUS_SCHEDULED; } }这里有一点要特别说executeAt 用的到底是墙上时钟还是单调时钟。如果任务只在本机内存队列里跑用 System.nanoTime 计算相对延迟更安全因为墙上时钟可能会被 NTP 校准乱跳导致任务提前触发或者迟迟不触发。但我最终用的是 System.currentTimeMillis核心原因是任务要支持持久化和集群扩展一旦跨了进程就必须用墙上时间来对齐。为了弥补时钟跳变的问题后面会在调度循环里做更稳的延迟判断这一点放到踩坑章节细讲。2.2 队列选型为什么是 PriorityBlockingQueue 而不是 DelayQueue任务对象定了之后下一个关键决定是选什么数据结构做排队。大部分人的第一反应是 JDK 自带的 DelayQueue毕竟名字里就带着延迟。我也试过 DelayQueue它的入队出队天然支持按剩余延迟时间排序队头永远是最先到期的任务。但它的数据结构其实不适合做复合优先级排序因为 DelayQueue 的排序依据是 getDelay 方法的返回值也就是说“到期时间”几乎锁死了顺序很难再把业务优先级作为一个独立维度塞进去。我最后选择的是 PriorityBlockingQueue 加自研调度循环。PriorityBlockingQueue 本身是一个支持并发访问的优先队列可以自定义 Comparator这样就能把执行时间和优先级结合起来排序。但注意不要想当然地以为“优先级越高越靠前就行”如果队列纯粹按优先级排序高优先级任务会不断插队低优先级任务会像慢性饥饿一样永远得不到执行。ax 里实际采用的排序口径是执行时间 executeAt 为主排序键优先级 priority 只作为同一到期时间内的次级排序键。也就是说一个低优先级但马上就要到期的任务一定会排在高优先级但还差很久才到期的任务前面。这样排序逻辑最直观也符合业务直觉——调度调度首先是“准点”其次才是“抢资源”。private static final ComparatorTask? TASK_COMPARATOR (a, b) - { int cmp Long.compare(a.executeAt, b.executeAt); if (cmp ! 0) return cmp; return Integer.compare(a.priority, b.priority); };如果你业务里真的需要“高优先级必须插队”正确做法不是改排序器而是在提交任务时把高优任务的 executeAt 提前让它自然排到队头。把时间语义和优先级语义分开来设计整个系统行为会好预测很多。2.3 worker 线程池模型固定线程数比动态伸缩更稳任务到期后要交给 worker 线程执行。这里的第一个版本我用的是 ThreadPoolExecutor 默认策略任务一多线程数就自动往上扩最大线程数设得还挺大结果一压测就发现性能不升反降。原因很典型线程扩张会带来上下文切换开销尤其是在任务以 I/O 密集为主、线程大多数时间在等待外部响应的情况下线程越多数次轮转浪费越严重。到后来我直接改成固定线程池线程数按经典公式估算线程数 CPU核数 * (1 平均等待时间 / 平均计算时间)。对 I/O 密集任务线程数通常设成核数的 2 到 4 倍对纯 CPU 密集型任务设成核数加一就够。调度线程和 worker 线程也做了严格区分。调度线程只有一个专职负责从队列取任务、判断任务是否到期、把到期任务交给线程池。worker 线程池独立执行业务逻辑。这样的好处是调度延迟不会因为任务执行耗时变长而劣化哪怕是队列里有一堆 800 毫秒的慢任务新到期的任务依然能被及时拿出来执行。3. 从零实现 ax 调度器关键代码与决策点3.1 调度器骨架与核心循环直接上核心代码。AxScheduler 有几个部分一个 PriorityBlockingQueue 放任务队列一个 ThreadPoolExecutor 做 worker 执行池一个 dispatchThread 做调度循环。这里为了讲解清晰我省略了部分状态管理代码但核心调度逻辑是完整的。public class AxScheduler { private final PriorityBlockingQueueTask? queue new PriorityBlockingQueue(1024, TASK_COMPARATOR); private final ThreadPoolExecutor executor; private final Thread dispatchThread; public AxScheduler(int coreWorkers) { this.executor new ThreadPoolExecutor( coreWorkers, coreWorkers, 0L, TimeUnit.MILLISECONDS, new LinkedBlockingQueue(4096)); this.dispatchThread new Thread(this::dispatchLoop, ax-dispatch); this.dispatchThread.start(); } public void submit(Task? task) { queue.put(task); } private void dispatchLoop() { while (!Thread.currentThread().isInterrupted()) { try { Task? task queue.take(); long remaining task.executeAt - System.currentTimeMillis(); if (remaining 0) { // 还没到期不能让出队后的空档要轮询等待剩余时间。 // poll(remaining) 在有新任务入队时会提前返回不会卡死。 Task? next queue.poll(remaining, TimeUnit.MILLISECONDS); if (next ! null) { queue.put(task); // 把没到期的放回去处理新的队头 continue; } } executor.execute(() - executeTask(task)); } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } } } }这个循环看起来简单但里面有个细节容易被忽略。我用queue.take()取出队头之后发现任务还没到期不能直接 sleep因为如果 sleep 期间来了一个更紧急的新任务队列里已经没人处理了。所以这里用的是queue.poll(remaining, TimeUnit.MILLISECONDS)它有两个语义如果等待期间有新任务入队会立刻返回新队头如果一直没新任务就等够 remaining 时间返回 null。返回 null 说明原来的任务到期了继续往下走执行它返回了新任务就把旧任务放回队列重新排序。这个设计兼顾了延迟正确性和新任务的紧急抢占。3.2 调度线程被频繁唤醒怎么办有的同学可能会问每次到期任务都要 poll 等待如果队列里同时有几百个相差几秒的任务是不是调度线程会被反复唤醒然后把任务放回去再取出来效率很低是的按上面的写法确实会这样但实际测试下来问题不大。原因有两个。第一PriorityBlockingQueue 的 put 和 poll 操作是有锁的锁竞争在任务数几千到几万时没有想象中严重因为大多数时间只有一个调度线程在操作队列。第二更关键的是这种“取出来发现没到期再放回去”的情况不会频繁发生。因为在高峰时队列里通常有大量已到期的任务take 出来就立刻执行了不需要等待只有队列“空转”的时候才会反复做这个动作。如果你实在介意可以把没到期的任务挪到一个专门的“待命集合”里由单独一个 timer 负责唤醒但这就增加了代码复杂度对 ax 的定位来说并不划算。3.3 优雅关闭别一梭子打断正在跑的任务调度器最难写的地方其实是关闭。很多自己写的调度器在应用停机时直接把线程池 shutdownNow结果正在执行的任务被强行打断数据库事务状态卡在半路。ax 里单独实现了 shutdown 方法并且支持两种关闭模式。public ListTask? shutdown(boolean cancelPending) { dispatchThread.interrupt(); // 1. 调度线程先停 executor.shutdown(); // 2. worker 停止接收新任务 if (cancelPending) { ListTask? remaining new ArrayList(); queue.drainTo(remaining); // 剩余任务全部捞出来返回 executor.shutdownNow(); // 同时中断正在执行的任务 return remaining; } try { executor.awaitTermination(30, TimeUnit.SECONDS); // 等待执行中任务结束 } catch (InterruptedException e) { Thread.currentThread().interrupt(); executor.shutdownNow(); } return List.of(); }这里要提醒一下cancelPending 为 true 时虽然返回了剩余任务清单但 shutdownNow 只能发中断信号不能保证任务真的停掉。如果业务代码没有正确响应中断任务可能还在后台继续跑。所以 ax 的团队使用规范里写得很清楚shutdown 只是调度层停止不代表业务逻辑层停止需要真正确保任务不再继续必须业务侧自己配合幂等和事务回滚。4. 调度的可靠性设计状态机、异常恢复与幂等4.1 任务状态机与全链路可观测只把任务丢到队列里跑完就结束是不足以支撑生产环境的。ax 里每一个任务都有一整套状态流转这也是后续排查超时和重复执行问题的关键。任务从创建开始会经过 CREATED、SCHEDULED、RUNNING最终进入 SUCCESS、FAILED、TIMEOUT 或者 CANCELED。初始提交时是 SCHEDULED表示已经在队列里等待worker 取出准备执行时变成 RUNNING执行成功变 SUCCESS抛异常或者到达最大重试次数后变 FAILED执行超时被取消变 TIMEOUT被上层主动取消变 CANCELED。每个状态之间的流转都会记录时间点方便后面定位“这个任务到底卡在哪一步”。状态不只是用来展示的它直接决定调度器的行为。比如 RUNNING 阶段的任务如果被超时扫描线发现超时调度器会把它标记为 TIMEOUT并从 executor 里提交 cancel如果 FAILED 且还有重试次数调度器会重新计算 executeAt 再放回队列。每个状态行为必须清晰否则一个任务可能同时被超时扫描和重试机制处理造成重复执行。4.2 失败重试与退避策略别把退避写成一个固定值失败重试是最容易被写错的。第一版 ax 的重试间隔是固定 3 秒结果线上出现过一次“重试风暴”一大批任务同时失败又同时 3 秒后重放数据库压力瞬间翻倍反而把服务打得更不可用。后来我把重试策略改成指数退避加抖动。公式很简单retryDelayMs baseDelayMs * 2^(attempt - 1) Random(0, 500)比如基础延迟 1000 毫秒第一次失败后等 1 到 1.5 秒第二次等 2 到 2.5 秒第三次等 4 到 4.5 秒。抖动的作用是让同一批失败的任务不会整齐地“踩点”重试而是散开分布。maxRetry 也必须设上限ax 里默认最多重试 3 次防止一个永远失败的坏任务无限循环消耗队列。举个例子订单关闭任务如果第一次执行时数据库连接池满了导致失败不能立刻又让一万个同类任务同时去重试。加入退避后每个任务的重试时间点被打散数据库连接池才有机会恢复。4.3 超时中断future.cancel(true) 不一定能停掉任务执行超时是 ax 里最容易踩坑的地方简单说就是future.cancel(true)只是发了个中断信号业务线程如果不响应中断任务实际上还在跑。ax 里对超时的处理是这样实现的worker 拿到任务后将业务逻辑提交给 executor 得到一个 Future并记录一个 timeout 扫描任务扫描时间等于任务的 timeoutMs。private void executeTask(Task? task) { Future? future executor.submit(() - doBiz(task)); timeoutScanner.schedule(() - { if (!future.isDone()) { task.status STATUS_TIMEOUT; future.cancel(true); } }, task.timeoutMs, TimeUnit.MILLISECONDS); }这里我单独起了一个单线程的 timeoutScanner 专门负责扫描超时避免超时逻辑占用 worker 线程。注意timeoutScanner 只负责标记状态和发送中断信号不负责业务补偿。如果业务代码里没有处理 InterruptedException任务可能照常执行到一半甚至执行完只是状态已经被标记成超时。这个情况在打日志时特别有迷惑性明明日志显示任务超时了但订单还是被关闭了因为它真的跑完了。所以 ax 的使用规范里强调任何被超时标记的任务必须通过业务侧幂等逻辑来兜底。调度器不能因为超时了就盲目重跑否则可能出现两个线程同时操作同一个订单的状态造成数据错乱。4.4 任务幂等调度系统不会替你做这件事既然调度器做不到完美取消业务侧就必须在上层自己做好幂等。这一点无论单机还是分布式都是铁律。ax 侧的幂等逻辑很简单每个任务自带唯一 taskId调度器本身保证同一个任务不会在同一时刻被两个 worker 从队列中取走执行。但真正落到业务侧就要靠状态查询和乐观锁。比如关闭订单任务里执行的不应该是UPDATE orders SET status closed WHERE order_id ?而应该是UPDATE orders SET status closed WHERE order_id ? AND status paid这样即使任务被重复执行第二次也会因为 status 不是 paid 而不再生效。如果单机版 ax 的队列换成分布式幂等就更加重要了。任务可能在 worker 宕机后被重新投递如果没有唯一索引或者业务状态机兜底同一个任务跑两遍几乎是必然发生的。5. 实测数据与调参复盘一个低优先级任务的苏醒故事5.1 压测方法与场景设计ax 在正式上线前我在测试环境做了一轮压测。压测环境是 8 核 16G 的容器实例JVM 堆设 4G。模拟的真实业务任务共 10 万个其中 20% 是耗时 50 毫秒的轻量 I/O30% 耗时 200 毫秒50% 耗时 800 毫秒。这些任务被随机设置了 0 到 30 秒不等的延迟时间尽量贴近生产环境分布。我重点测的是两个指标整体吞吐量也就是每秒能完成多少任务以及 P99 延迟也就是任务到达计划执行时间后实际被 worker 拉起来执行的延迟。后者特别重要它反映的是调度器的“守时能力”。5.2 核心数据表格线程数与队列容量怎么配下面这组数据是从多轮压测里挑出来的每一行代表一个固定线程数和固定队列容量的组合。线程数队列容量吞吐量任务/秒P99 延迟ms420483121240820486058301620481175460322048138652016102410426301640961288390数据有两个值得注意的现象。第一线程数从 16 加到 32吞吐量只涨了约 17%但 P99 延迟反而从 460 毫秒涨到了 520 毫秒说明线程数到 32 之后上下文切换已经开始拖后腿收益明显递减。第二同样线程数下队列容量从 1024 加到 4096吞吐量和延迟都变好了因为队列变长之后突发提交时 worker 不容易因为队列满而被拒绝任务积压也少了。但队列容量也不是越大越好。队列太长任务的“在队等待时间”会变高极端情况下一个任务可能排队几分钟还没执行而且一个任务对象包含业务 payload队列里十万个任务可能吃掉几百兆内存。ax 里我给团队的参考建议是队列容量设为每秒新增任务数的 3 到 5 倍比如每秒大概新增 2000 个任务队列设 4096 到 8192 比较合适。5.3 血泪教训低优先级任务真的会饿死压测过程中我遇到一个特别典型的坑。当时为了测试优先级功能我故意塞了一批 priority10 的低优先级任务和大量 priority1 的高优先级任务。按我最初的排序器低优先级任务在 executeAt 相同的情况下永远排在后面结果那批低优先级任务压测结束后一个都没执行。后来加了老化策略才解决。ax 里每隔 30 秒会用一次 O(n) 的线性扫描找出所有等待时间已超过 40 秒并且优先级低于某个阈值的任务把它们的 priority 临时上调一级。这里做了个 O(n) 扫描对一万级任务的队列来说开销完全可接受因为优先级队列的 remove 走的是数组扫描低频操作没必要优化到极致。老化策略保证了低优先级任务不会被彻底饿死最多就是延迟执行不会永不执行。这轮压测让我意识到优先级调度不只是 Comparator 里一行字段排序的问题它背后还牵涉到公平性和饥饿控制。你在自己的调度器里如果也要做优先级一定要先把“低优任务会不会被饿死”这个问题想清楚。6. 从单机到集群ax 调度的分布式演进6.1 单机 ax 的边界在哪里ax 单机版上线后稳定跑了几个月但任务量涨到日均百万级之后问题开始暴露。首先任务只存在内存队列里进程一重启还没执行的任务全部丢失。其次单台实例的线程池和内存有限高峰期已经需要多副本部署但内存队列在不同实例之间彼此不感知同一个任务可能被多个实例重复提交。我当时明确了一个判断标准如果你的任务量只是每天几千条延迟等级在秒级到分钟级单机 ax 加一层数据库持久化完全够用没必要上分布式。但如果任务量到了每天几十万条或者服务要做到多副本无状态那 ax 就必须演进成“调度中心 worker 执行器”的架构。6.2 轻量级分布式方案Redis ZSet 消息队列ax 的分布式演进没有引入复杂框架核心思路是拆分把“任务排队和到点触发”交给 Redis ZSet把“延迟消息投递和执行”交给消息队列worker 只负责消费执行。Redis ZSet 的做法非常直接score 就是任务的计划执行时间戳。入队时执行 ZADD调度循环每隔 500 毫秒用 ZRANGEBYSCORE 取当前时间之前的所有任务投递到消息队列。# 任务入队 ZADD ax:delay:orders 1710000000000 taskId # 调度循环每 500ms 执行一次 ZRANGEBYSCORE ax:delay:orders -inf 1710000001000 WITHSCORES LIMIT 0 100 # 把到期任务投递给 worker 消费 PUBLISH ax:ready:orders taskId这个方案最大的好处是思路直观、组件都是现成的Redis 负责延迟排序MQ 负责削峰和缓冲worker 完全无状态天然支持水平扩容。但有一个细节必须处理多个调度实例同时执行 ZRANGEBYSCORE 时同一个 taskId 可能被多个实例投递到 MQ造成重复消息。所以调度中心必须保证只有一个实例在跑最简单的方式是抢一把分布式锁也可以用 Redis 的 Lua 脚本把“扫描 移除”做成原子操作。我在生产里是用分布式锁做降级保护确保同一时刻只有一个 ax-scheduler 实例在调度。6.3 分片、去重与故障恢复集群化之后三个问题必须解决。第一个是去重worker 消费 MQ 消息时必须按 taskId 做幂等最简单的是数据库唯一索引插入失败说明消息重复直接 ACK 丢弃。第二个是分片如果 worker 数量多可以让每个 worker 按 taskId 哈希取模归属固定分片这样同一个订单的任务只会被同一组 worker 处理减少锁冲突和跨节点协调。第三个是故障恢复worker 宕机时它正在执行的任务不会被任何人接管这时候需要引入心跳检测调度中心定期检查每个 worker 的心跳超过阈值就认为 worker 挂了把这些 worker 的未完成任务重新投递。重投递的操作必须慎之又慎因为上一节说过超时和故障并不能保证任务真的没在执行。所有重新投递的任务都要靠业务侧的状态机和乐观锁兜底否则就是雪上加霜的双重执行。这一轮演进之后ax 的调度语义其实没变任务对象、状态机、重试、超时还是同一套逻辑变的只是队列和执行器从内存换成了 Redis 和 MQ。这也是我后来最深的体会写调度器的时候千万不要把“调度语义”和“存储执行”焊死在一起。只要抽象好任务和状态换存储、换执行器都只是适配层的事上层业务代码可以完全不动。我在实际使用中的体会是调度器这种东西最难的不是“到点执行”这四个字而是把“任务什么时候执行、执行了没有、失败了怎么办”这三个问题回答得明明白白。ax 这套东西做到后面真正值钱的部分并不是那个优先级队列而是它给业务侧提供了一套清晰可靠的任务生命周期管理。如果你也想写一个轻量调度器我的建议是第一步先把单机版的延迟触发、优先级、重试、超时全部打磨透再考虑分布式第二步在任何自动化重试和补偿之前先确保业务侧幂等否则调度器越智能数据错乱越严重。调度器永远只是工具业务逻辑的正确性还得靠自己的状态机来保证。
返回列表