
先说说 ax 是怎么来的。它不是拍脑袋造出来的轮子而是被一堆定时任务和失败重试的问题逼出来的内部项目。当时我手上有一个数据平台跑着大量同步、缓存刷新、报表生成、IoT 指令下发之类的活一开始全用 cron 顶着上了集群之后才发现cron 根本回答不了这几个问题同一时刻到底该哪台机器执行机器挂了任务谁来接管执行失败了要不要重试、怎么重试ax 调度这个内部代号指的就是为了解决这些问题而做的一套轻量调度内核。这篇文章把我从需求梳理、模型选型、核心实现到踩坑排查的完整过程写下来给正在考虑自研调度器、或者单纯想搞懂调度底层原理的朋友一个参考。1. 为什么会出现 ax 这个调度引擎先交代清楚项目从哪来这部分决定了后面所有设计决策不看背景直接看代码容易产生为什么要做这么复杂的误解。1.1 业务场景里定时任务到底缺什么在我们的业务里典型任务长这样每天凌晨 2 点同步一次订单数据到数仓每 5 分钟刷新一次热点商品的本地缓存每小时生成一份经营报表每分钟向一批智能设备下发控制指令。单机单实例的时候crontab完全够用写几行配置、配个日志轮转问题就结束了。但服务一上集群语义立刻变了。首先是重复执行两台机器同时从配置中心拉到了同一批任务到了触发时间各自执行一遍缓存刷新这种幂等操作还好数据同步、订单推送这种操作就会出现重复写入严重时直接把对端系统打挂。其次是无人接管某台机器宕机后它身上挂着的任务在它重启之前永远不会再被触发除非人工介入。再一个是失败无补偿任务执行到一个中间步骤抛异常cron 只能记一行日志然后等下个周期但很多任务需要马上重试、需要退避、需要告警。这些问题的本质是时间触发只是调度最表层的功能真正的调度是分布式环境下对任务归属、执行状态和失败补偿的管理。理解了这一点后面设计状态机、幂等、分布式锁的时候就不会觉得小题大做。1.2 为什么没有直接套用现成的开源调度框架不少朋友听到自研调度的第一反应是有病吧XX 不香吗。我承认成熟框架功能全、社区大、文档多但它解决不了我们的三个实际问题。第一是太重。我们有一部分服务要部署在边缘节点上设备内存很小平时只跑一个轻量 Agent塞进一个完整调度框架的客户端和依赖光依赖冲突就能折腾一整天。我们需要的是一个可以被裁剪的内核边缘节点只保留注册、触发、执行、回报四条链路中心节点才启用完整的重试、锁、状态机能力。第二是线程模型不透明。调度框架为了通用性线程池、队列策略、锁实现往往是黑盒。真到了线上任务积压、延迟飙升的时候你只能去论坛翻 issue没法直接通过代码判断瓶颈在哪。自研之后整个线程模型长什么样我们心里有数出问题能直接看源码定位。第三是学习成本。为了用好一个框架团队需要读一遍它的架构设计、配置项和故障排查手册这成本远比写一个够用的调度内核高。当然这里不是劝大家都去自研而是说当你的运行环境、资源约束和现成框架的假设不一致时自研内核是一个合理选项。ax 也不是一个完整平台它只做调度引擎不做工作流编排边界非常克制。1.3 命名由来与项目边界ax 原来的全称是 Agile eXecutor后来觉得这名字太正经就干脆当成一个内部代号沿用下来。项目边界从一开始就画得很清楚只做时间触发 状态管理 失败补偿不做 DAG 编排、不做分布式事务、不做 Web 控制台。这些能力留给上层业务按需搭建。这句话说起来轻巧实际上非常关键。很多自研项目死在顺便把 XX 也做了上边界一模糊就再也不可能轻量了。ax 的核心链路就五环注册任务、等待触发、投递执行、收集结果、失败重试。其他一切都要围绕这五环不许乱长功能。2. ax 的核心设计调度模型怎么选调度引擎的灵魂不在代码而在模型。模型选错了后面每修一个 bug 都是在给错误决策补窟窿。2.1 调度模型对比轮询、延迟队列、时间轮我最早考虑过三种方案挨个说下取舍。第一种是数据库轮询。每隔固定时间扫描一张任务表把到期任务捞出来执行。它的优点是真的简单任务定义、状态天然落库重启不丢数据。缺点也很明显扫描间隔决定了调度延迟下限间隔设置到 1 秒数据库压力就上来了而且每台机器都扫同一个表还需要额外做分布式抢锁把简单问题复杂化。这种方案适合任务量小、秒级延迟完全无所谓的场景。第二种是延迟队列。JDK 自带的DelayQueue或者用 Redis 的 zset 按执行时间戳排序再起一个线程不断取队首。延迟精度比数据库轮询高很多能做到毫秒级。但它有内存或存储开销问题而且取出到投递的过程仍然需要有人盯着队首本质上还是一个轮询只是粒度变细了。第三种就是时间轮。它用环形数组模拟表盘指针每走一个 tick就把当前槽位上的任务批量取出来。插入和移除的时间复杂度都是 O(1)不依赖数据库适合高频、短延迟的场景。我用一个餐厅类比来理解它不采用时间轮的做法是传菜员拿手机给每桌设一个倒计时菜越多越手忙脚乱时间轮相当于把厨房出菜口做成一个旋转转盘每个菜做好放到对应的格子转盘转一圈到口的菜自动流出来传菜员只要守着转盘出口就行。ax 最终选了时间轮作为核心调度结构但不是全盘取代延迟队列——长时间任务组合了一个辅助线程来处理后面 2.2 节细讲。2.2 时间轮的工作原理与选型原因时间轮的本质是一个环形数组每个槽位代表一个时间单位槽位上挂着一个任务链表。假设 tick 是 100 毫秒轮盘有 512 个槽那么转一圈的时间是 51.2 秒。指针每一 tick 前进一步把当前槽位链表里的所有任务取出来逐个投递给执行线程池。这里有几个点必须想清楚。第一任务延迟超过了轮盘一圈怎么办两种常见解决思路一是多级时间轮类似水表上的小数位低层转一圈高层走一格二是记录圈数任务进槽位时带上还需要转几圈每圈只做计数。ax 用的是更懒但更实用的方案51.2 秒以内的任务直接进时间轮超过 51.2 秒的定时任务放到一个独立的优先队列里由一个周期线程负责到点后把任务换算成时间轮内的绝对时间再注册进去。之所以这么设计是因为 90% 的业务任务周期都在分钟级到小时级用单层时间轮处理短延迟用辅助队列兜底长周期两套配合代码量最少心智负担最小。第二个需要想清楚的是槽位冲突。多个任务到期时间可能落到同一个 tick 里所以槽位上挂链表而不是单个节点。取出当前槽位链表后必须先把整个链表摘下来再逐个执行避免执行过程中新加入的任务污染当前批次。第三个是精度与性能的取舍。如果 tick 设到 10 毫秒理论上调度延迟更小但 CPU 成本明显上升因为每个 tick 都要做一次取链表、判空、加锁空转也是开销。ax 最终定了 100 毫秒 tick实测下来覆盖绝大多数业务没问题。谁要是跑到秒级以下实时调度那是实时任务系统的活不该让调度引擎硬扛。2.3 任务状态机设计调度系统最怕状态定义模糊。执行到一半算不算失败正在等待重试的任务被手动取消是哪种状态这些问题不提前定清楚后续写分布式锁和幂等一定会乱。ax 把任务实例的生命周期定义成下面这张表状态含义触发动作PENDING已注册等待触发注册时写入TRIGGERED触发时间到已投递给执行线程时间轮取出时写入RUNNING执行器正在执行执行器启动时写入SUCCESS执行成功执行器回报时写入FAILED执行失败不再重试重试次数耗尽时写入RETRY_WAITING执行失败等待退避重试计算退避时间后从 RUNNING 转入CANCELLED手动取消取消接口调用时写入每次状态变更都会发布一个事件告警模块订阅FAILED和RETRY_WAITING监控模块订阅RUNNING和SUCCESS做耗时统计。这里有个容易踩的坑RETRY_WAITING这个状态必须单独存在不能直接把任务丢回 PENDING否则你无法区分一个从没执行过的任务和一个失败后准备重跑的任务重试次数和退避时间的隔离也就无从谈起。3. 核心实现时间轮、分发与执行线程这一章直接拆源码级别的实现思路。我不会贴完整工程但会把关键数据结构和执行路径讲清楚照着写一个最小可用版本不难。3.1 时间轮的数据结构与推进逻辑时间轮可以用一个类来表达核心就三块数组槽位、当前指针、tick 间隔。简化版本长这样public class TimingWheel { private final long tickDurationMs; // 每格时间固定 100ms private final int ticksPerWheel; // 格子数量固定 512 private final AtomicInteger currentTick; // 当前指针位置 private final Node[] slots; // 环形数组 public TimingWheel(long tickDurationMs, int ticksPerWheel) { this.tickDurationMs tickDurationMs; this.ticksPerWheel ticksPerWheel; this.currentTick new AtomicInteger(0); this.slots new Node[ticksPerWheel]; } public void add(long deadlineMs, Task task) { long remainMs deadlineMs - System.currentTimeMillis(); if (remainMs 0) { // 已经过期立即投递 executor.submit(task); return; } int tickOffset (int) (remainMs / tickDurationMs); if (tickOffset ticksPerWheel) { // 超过一轮交给长周期队列处理 longQueue.offer(new LongTermTask(deadlineMs, task)); return; } int targetIdx (currentTick.get() tickOffset) % ticksPerWheel; slots[targetIdx].add(task); } public void advance() { int idx currentTick.getAndIncrement() % ticksPerWheel; Node node slots[idx].drain(); while (node ! null) { executor.submit(node.task); node node.next; } } }注意几个关键的取舍add里用的是System.currentTimeMillis()去算业务上的绝对到期时间但advance()的推进节奏不能依赖它否则系统时钟被 NTP 调整时可能出大问题这个坑我放到第 5 章细说。另外取出槽位链表用的是drain()不是getAndClear()含义是把当前槽位链表整体摘下来交给执行线程池同时立即允许新的任务重新挂到这个槽位。这样避免了遍历执行完再还回去过程中的锁持有可能阻塞调度线程的问题。3.2 调度线程与执行线程的协作方式调度线程是纯粹的时间驱动者它由一个单线程调度器驱动每隔 100ms 调用一次advance()。它的职责只有一个把到期的任务从槽位链表搬进执行线程池的队列里绝不亲自执行任务。这里的原则是调度线程永远不能阻塞。很多调度框架延迟增高的原因就是不小心在调度线程里做了 IO、打了日志、或者同步等待执行结果。ax 在这块画了死线调度线程里禁止抛出业务异常任务投递动作全部走 try-catch即使某个任务注册时因为参数非法崩了也不能影响同一槽位里其他任务的投递。执行线程池用的配置我建议这么设核心线程数按 CPU 核数乘 2最大线程数按核心线程数的 2 倍队列用有界队列容量根据任务量压测来定拒绝策略选CallerRunsPolicy。很多人会选AbortPolicy觉得任务太多宁可丢但调度场景下丢了意味着数据同步、缓存刷新、报表会缺一次后果比阻塞更严重。CallerRunsPolicy的意思是线程池满了之后让调度线程自己帮忙跑任务虽然调度线程会被拖住短暂时间但至少任务不会丢用微量的延迟换任务可靠性值得。3.3 重试、幂等与分布式锁任务接口我强烈建议返回Result(code, message)而不是裸抛异常。背后的原因很简单抛异常只能表达失败了表达不了这是可以重试的失败还是这是业务上确定失败的失败。比如第三方接口明确返回参数不合法你再重试一百次也没用返回 503 超时重试就有意义。所以任务结果至少要带两个字段retryable和message。重试退避我选了指数退避加噪声long backoff baseMs * (1L (attempt - 1)) RandomUtils.nextLong(0, 1000);为什么加噪声因为如果几十个失败任务刚好同时进入重试队列第一次退避后它们又会同时醒来产生的不是重试而是二次流量高峰。加一个 0 到 1 秒的随机量可以有效打散这批请求。分布式锁是自研调度绕不开的一道坎。ax 用的是 Redis setnx 实现的租约锁拿到锁的节点才允许执行任务。最核心的注意点是锁必须带租约时间租约期间执行者必须主动续期。否则出现任务执行耗时超过租约锁过期被其他节点拿走两个节点同时执行同一个任务的场景你根本查不出是哪台机器干了两遍。这个坑的完整排查过程我放到 5.3 节。此外每个任务实例生成一个全局唯一executionId执行器写数据时把这个 ID 落到数据库唯一索引。即便锁机制真的出了某种极端问题数据库也会兜住最后一层防线。分布式系统没有单点保障必须层层设防。4. 实操搭一个最小可用的 ax 调度核心讲完原理给一套可以直接参考落地的骨架代码。目标是让读者看完能跑起一个最小调度器并挂上一个真实任务。4.1 核心接口与代码骨架三个核心接口缺一不可public interface Job { Result execute(JobContext ctx); } public class JobContext { private String jobName; private String executionId; private long triggerAt; // 以及需要的配置项、日志句柄等 }任务注册请求的定义public class ScheduleRequest { private String jobName; private long triggerAtMs; // 首次触发时间 private long intervalMs; // 周期0 表示单次任务 private int maxRetries; // 最大重试次数 private long retryBaseMs; // 退避基准时间 private Job job; // 实际执行逻辑 }这里我把 Job 直接塞进请求里是为了演示方便生产环境建议改为jobName - Job 实例的注册表调度请求只带名字和参数。4.2 注册、启动、关闭流程启动流程三步走初始化时间轮、初始化执行线程池、启动调度线程。public class AxScheduler { private TimingWheel wheel; private ExecutorService executor; private ScheduledExecutorService scheduler; private MapString, Job jobRegistry; public void start() { wheel new TimingWheel(100, 512); executor new ThreadPoolExecutor( coreSize, maxSize, 60, TimeUnit.SECONDS, new ArrayBlockingQueue(queueSize), new CallerRunsPolicy() ); scheduler.scheduleAtFixedRate(wheel::advance, 0, 100, TimeUnit.MILLISECONDS); } public void register(String name, Job job) { jobRegistry.put(name, job); } public void schedule(ScheduleRequest req) { wheel.add(req.triggerAtMs, () - { JobContext ctx new JobContext(req.jobName, genExecutionId()); // 执行前的分布式锁检查、状态流转触发 jobRegistry.get(req.jobName).execute(ctx); }); } public void shutdown() { // 第一步停止调度线程不再产生新触发 // 第二步等待执行中任务完成带超时 // 第三步关闭执行线程池 } }关闭流程是很多人忽视的重灾区。直接调executor.shutdownNow()虽然快但会把执行到一半的数据同步任务硬生生掐断。ax 的优雅关闭顺序是先停调度线程避免新任务进入再给执行线程池一个宽限期让它把已投递的任务执行完宽限期过了还赖着不走的才强制关闭。同样的思路要应用到长周期辅助线程和状态上报链路上。4.3 完整示例缓存预热任务每天凌晨 2 点把热点商品信息从 MySQL 加载到本地缓存这个任务非常适合拿来示范jobRegistry.register(hot-cache-warmup, (ctx) - { ListString hotIds queryHotIdsFromDB(); for (String id : hotIds) { ProductInfo info queryProduct(id); cache.put(product: id, info); } log.info(cache warmup done, count{}, hotIds.size()); return Result.success(); });注册调度请求时可以这样写ScheduleRequest req new ScheduleRequest(); req.jobName hot-cache-warmup; req.triggerAtMs nextTriggerAt(2, 0); // 下一个凌晨 2 点 req.intervalMs 24 * 60 * 60 * 1000L; // 每天一次 req.maxRetries 3; req.retryBaseMs 1000; scheduler.schedule(req);挂上去之后立刻观察三个点第一任务触发时间是否稳定在监控上画一条实际执行时间 - 预计算执行时间的差曲线第二执行耗时是否平稳突然的耗时尖峰往往意味着依赖的下游慢了第三失败重试日志是否正确进入RETRY_WAITING状态而不是直接消失。这三个监控点如果都正常这个调度器就可以开始承接更多任务了。5. 常见问题与排查技巧实录这一章是踩坑实录全部来自实际运行中遇到的故障和排查过程按问题出现概率排序。5.1 时钟回拨引发任务延迟最隐蔽的问题没有之一。现象是某天任务大面积延迟每个任务都慢几十到几百毫秒偶尔慢到秒级。第一反应查线程池没满查时间轮指针看起来在走。后来才发现是宿主机系统时间被 NTP 服务往回拨了一下导致所有基于System.currentTimeMillis()的剩余时间计算全部错乱。排查过程不算难对比宿主机日志里的时间源和任务延迟曲线就能发现。解决方案分两层第一所有业务触发时间计算仍以墙上时钟为准这是语义需要第二时间轮推进节奏改用System.nanoTime()计算 delta它不受系统时钟调整影响保证单调递增。这样即便墙上时钟跳了调度线程本身不会乱走顶多触发时间整体偏差一次下个周期自动纠正。我还顺手加了一个启动检查如果时间轮指针发生回跳立即告警并自动复位到当前时间对应的槽位。5.2 线程池饱和与任务堆积第二个高频问题是执行线程池被占满表现是任务延迟开始累积越积越多直达 10 分钟以上甚至几个小时。排查的时候不要先看 CPUCPU 高反而可能是执行线程在空转等待下游 IO。正确的排查顺序是先看线程池活跃线程数是否长期等于最大线程数再看队列大小是否持续增长最后看任务的执行耗时分布有没有明显右移。ax 在这个问题上加了两个保险。一是执行耗时自动打点超过 P95 阈值就记一条慢任务日志定位是哪个任务在拖累整体二是支持对单个 Job 单独配置并发上限避免某个慢任务无限占用公共线程池。这里我再强调一次CallerRunsPolicy一旦线程池满了调度线程被迫帮忙执行任务会反过来拖慢时间轮推进但这是主动选择的可靠优先策略宁可慢不可丢。5.3 分布式重复执行第三次踩坑是在一个报表任务上。这个任务平常 2 秒跑完锁租约设的 5 秒一直相安无事。某天上游数据量暴增任务执行耗时涨到 8 秒锁在第 5 秒过期另一个调度节点立刻以任务无人执行为由抢锁再跑了一遍于是两份报表同时生成下游对账对不上。这个案例暴露了两个问题第一租约不够长当时只按平时耗时的 2.5 倍设置没预留出足够的波动空间第二没有续租机制执行者持有锁期间不会主动给锁续期一旦超时只能眼睁睁被抢。修复措施是双重保险执行线程内部起一个租约续期守护线程每过租约时长的一半就续一次同时执行器写库时带上executionId唯一索引即使将来锁机制再有闪失数据库也会拒绝第二条重复写。这个案例被团队当作反面教材之后所有任务上线前都要做一次锁租约、任务耗时的推演。5.4 精度与性能的取舍最后聊一个不算 bug 的取舍。一开始把 tick 设成 50ms任务延迟是低了但 CPU 占用比 100ms 时高了不少。原因很简单每个 tick 都要做一次数组定位、链表摘取、并发安全处理即使大部分槽位是空的这些操作的成本省不掉。后来做了个对比压测同一批任务在 50ms tick 和 100ms tick 下的调度延迟差别大约只有几十毫秒但 CPU 开销差了近一倍。于是我默默把 tick 调回 100ms并且只有一个特殊业务对触发时间精度要求 20ms 以内的实时指令下发单独开了一个细粒度轮盘跟主轮盘隔离互不干扰。这个案例说明调度精度不是一个可以无限优化的指标想清楚了业务真实需求再定参数比盲目调小 tick 更值得。最后给一点实际建议。如果你也想写一个类似的调度器不要把精力全花在时间轮算法上那只是入口。真正决定这个项目能不能扛住生产压力的是三件事状态一致性、幂等控制、失败补偿。ax 走到今天我自己估算至少有 70% 的调试时间都花在任务到底算成功没有、到底该不该再执行一次这类问题上剩下的才是时间推算、性能调优和监控告警。把这些底子打牢再谈调度精度才有意义。