☰
从定时任务到分布式调度:自研轻量级调度引擎AX的设计与实践
2026/9/26 9:24:49 网站建设 项目流程

先说一个可能很多人都有过的经历:业务量还小的时候,接到需求就是写个定时任务,扔到服务器上让它跑。等任务数量从几十个涨到几千个、上万个,执行时间、失败重试、集群环境下的重复调度、任务堆积这些乱七八糟的问题全冒出来之后,我才意识到手里的"万能工具"早就不够用了。这个项目的主角,就是我自己从零写的一个轻量级分布式调度引擎,代号叫 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 Map<String, 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 名和反射调用)、HttpExecutor(POST 请求到目标 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 { Future<AxResult> 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 );

选举逻辑:

  1. 每个节点启动时尝试插入自己的 node_id 作为 Leader;如果插入成功,它就是 Leader。
  2. Leader 每隔 3 秒更新一次心跳时间。
  3. 非 Leader 节点持续读这张表,如果心跳时间超过 10 秒未更新,视为 Leader 失联,尝试删除旧记录并插入自己,抢锁成为新 Leader。
  4. 原 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 ms
P99 调度延迟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 出发,设计一个支持前置依赖、分支选择和失败重跑的任务编排执行器。

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

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

立即咨询