先说一个可能很多人都有过的经历业务量还小的时候接到需求就是写个定时任务扔到服务器上让它跑。等任务数量从几十个涨到几千个、上万个执行时间、失败重试、集群环境下的重复调度、任务堆积这些乱七八糟的问题全冒出来之后我才意识到手里的万能工具早就不够用了。这个项目的主角就是我自己从零写的一个轻量级分布式调度引擎代号叫 AX。它不是什么千万级流量的明星框架但确实把我从定时任务的泥潭里拉了出来。这篇文章我会从需求拆解、核心建模、调度算法、可靠性保障到集群一致性把整个设计和实现过程完整拆开适合正在用传统定时任务框架、或者准备自研调度系统的同学参考。1. 为什么我要从零写一个调度系统这个问题我几乎被每个同事问过现成的 Quartz、XXL-JOB 不香吗为什么非要自己造一个轮子说实话不是想造轮子是被实际业务场景逼的。1.1 传统定时任务框架的瓶颈我当时所在的业务线核心是大量数据同步和报表生成任务。最开始用 Quartz 做单机定时调度配上org.quartz.properties调参跑一两年也还凑合。直到两个问题集中爆发第一单机处理能力有天花板。任务总数过万之后Quartz 默认的线程池默认10个线程经常被长任务占满短任务排队时间越来越不可控。调大线程池数据库连接池、下游服务又扛不住。第二调度和执行耦合。Quartz 把调度逻辑放在每个应用节点里集群部署时虽然可以用 JDBC JobStore 做持久化但负载不均、重复触发等细节需要处理得很小心。一旦某个节点假死又恢复那批任务到底该谁执行经常要靠人工查日志判断。1.2 市面上现成方案为什么不合适我认真对比过两类方案。一类是增强版定时框架比如 ElasticJob。它的分片思路很好但引入依赖较重而且它的 job 分片 模型更适合处理超大任务队列我们很多任务是短频率高消耗型分片反而带来额外复杂度。另一类是完整调度平台比如 XXL-JOB。功能确实全部署也简单但它有一整套管理员 UI、权限模型和执行器通讯约定。接入方必须按照它规定的模式改造自己的执行逻辑——用 HTTP 暴露执行器、注册到调度中心等。我们这个团队当时人少不想被一个平台绑死希望调度核心保持轻量执行方式能灵活扩展。1.3 AX 的设计目标所以我给 AX 定了几个目标它们后来也成为所有设计的评判标准调度核心不依赖任何重量级框架核心代码可以独立嵌入业务进程。支持多种触发方式CRON、固定延时、一次性延迟任务、依赖触发。调度逻辑与执行逻辑完全解耦执行器可以是本地方法、HTTP 调用、Shell 脚本。集群环境下同一个任务在同一时刻只能在一个节点执行一次严格幂等。单机调度器需要承受每秒至少 2000 次触发判定。所有状态变更必须可审计。这个项目代号 AX 其实就是 Action Executor 的缩写意思是动作执行引擎。名字是临时起的后来用得顺手了也就没改。2. 调度系统的核心建模任务、触发器与执行器的抽象调度系统的核心不是定时这两个字而是把业务动作抽象成可以被统一管理、触发、追踪的状态机。这一章讲清楚我如何给 AX 定义这三个最基础的概念。2.1 任务Task调度的基本单位在 AX 中一个任务代表一段具备明确业务含义的操作逻辑。我设计的数据结构如下public class AxTask { private String taskId; // 全局唯一任务ID private String name; // 任务名称 private TaskType taskType; // 本地方法、HTTP、Shell private String target; // 执行目标Spring Bean名称 / URL / 脚本路径 private MapString, String params; // 执行参数 private int retryTimes; // 失败重试次数 private int timeoutSeconds; // 超时时间 private boolean concurrentEnabled; // 同任务是否允许并发执行 }这里的重点是taskType和target的设计。我不希望调度器知道任务具体怎么跑它只需要负责在正确的时间把任务标记为可执行然后丢给对应的执行器。所以任务定义里不含执行逻辑只有执行方式和目标地址。有个小细节容易被忽略concurrentEnabled。默认情况下如果上一个任务实例还没跑完下一次触发到达时会被丢弃默认策略这样能防止数据同步任务意外叠加。但如果某个任务的执行时间可能超过调度周期且业务上允许并行就需要显式打开开关。这个字段在后续集群场景中也很重要我在可靠性章节会展开。2.2 触发器Trigger让任务动起来的信号触发器负责产生该执行了这个信号。AX 支持的触发器类型如下类型适用场景精度备注CronTrigger报表生成、定时同步秒级标准 Cron 表达式扩展支持年字段IntervalTrigger轮询型任务、心跳检查毫秒级固定时间间隔可跟具体时间对齐OnceTrigger延迟队列、临时任务毫秒级注册后仅执行一次DependencyTrigger跨任务依赖编排取决于上游上游完成后触发这里有个教训Cron 表达式虽然好用但很多人对秒级 Cron 有刻板印象觉得Cron 最小粒度就是分钟。其实 Quartz 的 Cron 支持秒级字段AX 也沿用这个约定。比如每5秒执行一次的表达式是0/5 * * * * ?。比较难处理的是 DependencyTrigger。上游任务完成后触发下游听起来简单实际要解决上游成功才算完成还是上游结束就算完成的问题。我的方案是把触发源事件任务完成事件抽象成带状态的事件流DependencyTrigger 只监听状态SUCCESS的事件失败或超时事件不会触发下游。2.3 执行器Executor任务的最终归宿执行器的抽象决定了调度系统的可扩展性。我定义了四个核心接口方法public interface AxExecutor { AxResult execute(AxTask task, String instanceId); default boolean matchTaskType(String taskType) { return false; } default void onSuccess(AxTask task, String instanceId) {}; default void onFailure(AxTask task, String instanceId, Throwable t) {}; }内置实现有三种LocalMethodExecutor通过 Spring Bean 名和反射调用、HttpExecutorPOST 请求到目标 URL、ShellExecutor执行本地脚本。接入方想加一种新执行器只需要实现AxExecutor接口并注册到执行器路由器。一个真实例子有一个团队接入了数据仓库的 SQL 查询任务直接通过 LocalMethodExecutor 把 SQL 执行器注册进来。整个过程没改调度核心一行代码这验证了抽象边界的合理性。2.4 任务状态机与实例跟踪任务定义是静态的每次执行会产生一个任务实例。这个实例的状态流转是调度系统可靠性的基础。AX 的状态定义如下CREATED - SCHEDULED - RUNNING - SUCCESS |- FAILED |- TIMEOUT |- RETRYING我在这里就吃了不少亏。最初状态设计只有三种待执行、执行中、已完成一旦出现重试或超时根本没法区分是首次失败还是重试后的失败日志审计时非常被动。后来引入 RETRYING 和 TIMEOUT 两个独立状态才把整个链路理顺。状态变更走统一的 StateStore 记录每次变更都写入instance_log表。后来排查线上问题这张表帮了大忙可以直接回答谁在什么时间把任务状态改成了失败。3. 时间轮调度器从 Timer 到分层时间轮的选择如果说任务建模是骨架调度算法就是心脏。这一章会详细拆解为什么我用分层时间轮作为 AX 的核心调度器以及它在面临高触发频率时做了什么优化。3.1 JDK 内置定时器的硬伤很多人写定时任务用java.util.Timer或者ScheduledExecutorService简单场景完全没问题。但在 AX 设计时这两个方案我直接否了原因很实际Timer只有一个后台线程一个任务执行时间过长会阻塞后续所有任务这在调度系统里是不可接受的。ScheduledExecutorService把任务存在DelayedWorkQueue中插入是 O(log n) 的复杂度但删除、取消、重新排序都需要额外代价。当任务数达到一万、十万级时每次触发都要做一次优先级队列操作调度延迟会明显上升。3.2 时间轮原理一个所有人都能理解的比喻时间轮本身不是新东西Netty 的HashedWheelTimer、Kafka 的定时器都用它。原理可以拿钟表来理解表盘有60个刻度秒针每走一格对应一秒每一格上挂着该在这一秒触发的任务。当秒针转到某个刻度取出该刻度对应的一串任务一一执行。这种设计下任务的添加是 O(1)——只需要计算应该挂到哪个刻度而不是每次都做全局排序。但单层时间轮有个问题精度和内存的矛盾。如果刻度是1秒、一轮60格最多只能调度60秒内的任务超过一轮就会丢失。解决办法是把任务转几圈再醒来但这样长时间任务的效率并不高。3.3 分层时间轮精度与内存的平衡AX 使用的是两层时间轮结构public class LayeredTimeWheel { private WheelLevel secondLevel; // 刻度间隔 100ms共 10 格 private WheelLevel minuteLevel; // 刻度间隔 1s共 60 格 }插入一个延时5秒的任务优先放到秒级轮的合适刻度如果延时超过当前轮的跨度就放到更粗的一层。每次 tick 优先处理较细粒度轮上到期的任务并按需把粗粒度轮的任务降级到细粒度轮。这套设计与单时间轮最大区别在于上线后实测一万个任务的内存占用不到 10MB而使用ScheduledExecutorService时同样的任务规模要吃掉约 60MB 内存在优先级队列上。而且触发延迟从平均 5ms 降到 0.8ms 左右。3.4 触发判定的小优化合并同类任务在实际业务中很多任务有相同的 Cron 表达式比如每天凌晨2点同步可能有300个任务。如果每个任务单独触发、单独查找执行器压力虽不大但明显可以优化。AX 在时间轮之上加了一层TriggerGroup概念相同触发类型、相同触发表达式的任务会被归到同一组。时间轮上只挂一个组触发器到期后一次性拉出组内所有任务批量提交。这样既减少了时间轮上的节点数又便于统一扩缩容。这个优化前期没做后来压测发现调度器线程在逐任务触发上浪费了大量时间加上分组后 CPU 使用率直接降了40%。4. 任务执行的可靠性保障超时、重试与幂等一个调度引擎能跑起来很容易能稳定扛住业务不断跑才是真功夫。这章讲的是任务执行过程中出了问题怎么兜底而不是把任务丢给执行器就完事。4.1 超时控制每个任务都必须设超时时间我给 AX 设计了一个强制约定任务定义时不配超时时间默认 60 秒且日志给出告警。超时控制不是只靠 Future 的get(timeout)就完事关键在于超时后要标记当前执行实例已终止同时清理真正执行中的资源。用 HTTP 执行器举例如果一个 HTTP 请求超时后只是单纯把状态改成 TIMEOUT但实际请求还挂在底层的连接池里大量超时后连接池会被占满接口直接雪崩。所以 AX 的 HttpExecutor 超时之后要主动调用Future.cancel(true)同时携带线程中断信号让底层 HTTP 客户端能响应中断并释放连接public AxResult executeWithTimeout(AxCallable callable, int timeoutSeconds) { ExecutorService single Executors.newSingleThreadExecutor(r - { Thread t new Thread(r); t.setDaemon(true); t.setName(ax-executor); return t; }); try { FutureAxResult future single.submit(callable); return future.get(timeoutSeconds, TimeUnit.SECONDS); } catch (TimeoutException e) { future.cancel(true); throw new AxTimeoutException(); } finally { single.shutdownNow(); } }这里single.shutdownNow()非常关键它会给该执行线程发送中断信号。如果执行体内部对中断没有响应那超时控制就只是状态改了事还在跑早晚出问题。4.2 失败重试与指数退避不要一上来就疯狂重试重试是每个调度系统都会做的事但怎么重试大有讲究。AX 的重试配置如下retryTimes: 3 retryBackoff: EXPONENTIAL backoffBase: 1s maxBackoff: 60s指数退避的算法很简单第 n 次重试前等待min(base * 2^(n-1), maxBackoff)秒。为什么要这样下游服务如果出现瞬时故障往往一两秒就能恢复如果是长时间故障你每秒重试一次等于给下游火上浇油。重试间隔从小逐渐扩大才是对下游服务的基本尊重。重试还有一个隐蔽的大坑重试的实例 ID 需要保持不变。如果每次重试都生成新的 instanceId审计日志就无法串起第一次失败-第二次重试-最终结果这条链。我在设计时把 retry 定义成同一实例的状态变更而不是新实例的创建。4.3 幂等键防止重复执行造成业务事故分布式调度最怕什么同一任务在同一时刻被两个节点各执行了一次。下游如果是告警通知多一次可能只是打扰如果是扣减库存、同步数据主键冲突那就是事故了。AX 给每次执行实例生成一个instanceId格式如下[taskId]-[triggerTimestamp]-[nodeId]taskId任务唯一标识triggerTimestamp触发器判定该执行的时间戳毫秒级nodeId当前节点的唯一编号这个组合能保证即使两个节点在同一毫秒触发同一个任务instanceId 也不同。执行器接受到任务实例后可以通过 RedisSET NX EX抢占幂等键只有抢到的节点才真正执行。幂等键的核心思路很简单让重复执行的结果不产生重复影响要么在入口挡住要么在业务逻辑上天然幂等。调度系统能做的只是尽量保证同一时刻最多一个执行者真正业务幂等还是需要执行方配合。4.4 死信队列与人工介入即使有重试机制有些任务的错误始终无法自动恢复。比如下游数据库连接被删了、配置文件配错重试100次也是白搭。AX 的做法是超过重试次数上限后任务实例进入DEAD状态同时写一条死信记录并触发告警。死信记录会保存完整的任务参数、执行历史、最后一次错误堆栈。这样人工排查时不用再次凭空猜测现场直接能看到这个任务6月1日第一次跑失败索引翻倍了。很多调度框架没有死信概念任务失败就是失败日志刷过去了就没人管。实际运营中发现没有死信队列的任务系统等于把待人工处理这个状态弄丢了——每周都会漏掉几个重要任务。5. 集群模式下的一致性分布式锁与主从切换单机版本的 AX 跑顺利之后我把它部署到三台机器上问题随之而来怎么保证多个节点不会重复调度同一个任务怎么保证一个节点挂掉后任务不丢5.1 首先明确调度与执行应该分开考虑集群环境下最容易犯的错误是让所有节点同时抢任务执行。如果执行逻辑是幂等的倒还好但如果有状态修改抢任务就会导致资源浪费和潜在冲突。AX 的设计是所有节点一起接收任务定义但由同一个 Leader 节点统一负责触发判定。其他节点处于 Standby 状态只接收命令。这样调度这个写操作被收敛到一个节点不容易出冲突。而执行是允许分布式的Leader 触发后任务实例会按执行策略分散到不同节点运行更有利于利用集群资源。这个模式的好处是调度逻辑本身不用加大量分布式一致性算法只要保证 Leader 的选举和切换是安全的调度行为就稳定。5.2 Leader 选举基于数据库的简化方案我在 AX 中实现了一个基于 JDBC 的简化 Leader 选举机制。背后是一张ax_leader表只有一个字段记录 Leader 节点 ID 和心跳时间CREATE TABLE ax_leader ( node_id VARCHAR(64) PRIMARY KEY, heartbeat_time DATETIME NOT NULL, lease_seconds INT NOT NULL );选举逻辑每个节点启动时尝试插入自己的 node_id 作为 Leader如果插入成功它就是 Leader。Leader 每隔 3 秒更新一次心跳时间。非 Leader 节点持续读这张表如果心跳时间超过 10 秒未更新视为 Leader 失联尝试删除旧记录并插入自己抢锁成为新 Leader。原 Leader 恢复后发现自己不再是 Leader自动降级为 Standby。这种方式不引入 Zookeeper 等外部依赖对数据库只有一个主键约束的写操作性能可以接受。缺点是没有严格的 fencing 机制极端场景下可能出现旧 Leader 以为自己是 Leader新 Leader 也选出来了的情况。5.3 脑裂问题的实际处理上面提到的旧 Leader 以为自己是 Leader就是脑裂。为了缓解AX 在每个调度动作前增加了一个checkLease步骤调度器每次发布触发决定前重新校验 local node_id 是否仍是表里的 Leader且心跳在有效期内。伪代码如下public boolean checkStillLeader() { AxLeader leader leaderDao.findCurrentLeader(); if (!leader.getNodeId().equals(localNodeId)) { return false; } if (System.currentTimeMillis() - leader.getHeartbeatTime() leaseSeconds) { return false; } return true; }这个方案并不完美但很实用。它把脑裂窗口从一个任务必然重复执行缩小到Leader 心跳失效但节点还活着的短暂窗口再结合业务幂等实际线上事故大大减少。需要提醒的是如果你非常看重严格一致性还是需要引入共识算法Raft、ZAB或者使用提供租约语义的协调服务。AX 的简化方案适合中小规模集群3-10个节点容忍偶尔的重复触发但不适合对原子性要求极高的金融级任务。5.4 任务状态的集群可见性任务被 Leader 触发后执行在哪个节点状态如何更新需要全局可见。AX 将所有实例状态写入同一张 MySQL 表ax_instance各节点通过 Redis Pub/Sub 接收状态变更通知保证本地缓存和实际状态一致。这里有一个非常容易踩的坑状态更新直接用读取→修改→写回的方式容易覆盖别人的更新。比如两个节点同时收到同一个任务的执行结果一个成功一个失败后写的会把先写的覆盖。AX 的解法是状态写入采用条件更新UPDATE ... WHERE status ?先判断当前状态能否转移再执行更新。状态机本身也是防止覆盖更新的一种防线。6. 真实压测数据与踩过的坑最后这部分我把 AX 的上线压测数据和实测踩坑全程记录下来。每一条坑都有足够的背景和排查链路希望能帮你少走一段弯路。6.1 压测场景与数据压测环境3 台 8C16G 云主机MySQL 8.0单机Redis 5.0。共注册 20000 个任务其中 8000 个秒级触发任务、12000 个分钟/小时级任务。指标结果调度器每秒触发判定次数2170 次/秒平均调度延迟触发信号发出到执行器收到1.3 msP99 调度延迟6.8 ms单节点 CPU 使用率调度器线程21%集群执行任务吞吐HTTP 执行器860 次/秒这个数据说不上惊人但足够支撑当时的业务。压测过程中暴露的许多问题比数据本身更有参考价值。6.2 坑一时钟漂移导致任务提前执行某天凌晨同事报告凌晨2点的定时同步任务在1点59分55秒就跑了。查了很久发现多台机器系统时钟有轻微漂移Leader 节点时间比真实时间快了几秒调度判断提前触发了。这个问题的关键在于调度系统对时间敏感而服务器时钟并不是完全可信的。AX 后续在 Leader 选举和调度触发时都加入了时间校准机制——不直接用本机时间而是用 Redis 的TIME命令获取统一时间减少漂移。对于允许分钟级误差的任务这个优化完全够用。6.3 坑二GC 停顿导致的调度延迟尖刺压测时有段时间 P99 延迟从 6.8ms 飙到 800ms一开始怀疑是数据库问题查了慢查询发现 leader 表锁等待时间有点高但不是主因。后来用 JFR 采集 JVM 事件发现是调度器线程在 Full GC 时产生停顿。Full GC 主体来自调度器缓存任务定义时用了大量小对象老年代迅速膨胀。优化方式很粗暴调大年轻代大小减少对象晋升频率同时将部分高频访问的任务元数据缓存为扁平字节数组而非 Java 对象降低 GC 压力。改完后 P99 降到 12ms 左右。这个坑说明一个道理做高并发调度不能只看业务代码JVM 参数和数据结构也很关键。6.4 坑三任务阻塞在查询数据库上另一个性能问题时任务执行器线程池出现了积压大量任务排队等待执行。排查发现很多任务是数据库查询型任务SQL 优化不到位单个查询耗时 5 秒以上把执行线程都占住了。这里不是调度器的问题但 AX 提供了一个非常有效的兜底ExecutorPool采用分桶式线程池不同任务组绑定不同线程池某个组的慢查询阻塞了自己不会拖累其他组。这个设计和普通线程池的区别在于按业务维度隔离资源副作用是线程数增加了需要运维层面配合监控。6.5 坑四日志风暴引发的磁盘打满上线初期把每个调度动作、触发状态变更都打 INFO 日志当天凌晨就被磁盘警报吓醒。20000 个任务秒级触发任务一天产生的日志量非常惊人直接导致日志盘被写满。后来做了三层过滤正常状态变更只打 DEBUG任务失败记录 ERROR但不打堆栈堆栈单独存到死信表只有重试也用尽时才在 ERROR 中输出完整堆栈。为方便排查正常成功执行只输出一条摘要日志包含 instanceId、耗时和结果摘要。实测日志量降到原来的十分之一且排查效率更高。7. 个人体会与后续规划写到这里AX 调度系统这套实现思路基本就讲透了。我个人在实际操作中的体会是自研调度引擎的难点从来不是写个定时器而是把任务建模、时间轮调度、可靠性保障和集群一致性这些看似独立的部分串成一个整体每个环节都需要刻意设计。如果你也想复制这套方案我给三个具体建议先分析真实任务分布再决定技术选型。90% 的任务是小时级以上的低频任务时完全没必要为海量秒级任务提前优化过度设计比不设计更麻烦。幂等和状态机是实现可靠调度的基础底座先把这两块做扎实后面的分布式锁、重试、死信都好加。监控和审计从第一版就要有。调度系统是无人值守的系统没有状态监控和日志追踪出问题时往往已经造成业务损失。后续我计划给 AX 加上任务血缘追踪和依赖编排的 DAG 可视化让跨任务的工作流状态更直观。这个内容如果大家感兴趣我后面可以再单独写一篇聊聊如何从一张 DAG 出发设计一个支持前置依赖、分支选择和失败重跑的任务编排执行器。
