后端干久了,你会发现“定时任务”四个字能撑起半个分布式系统的复杂度。我主要负责交易链路那块,报警群里被问得最多的就是“怎么又有任务没跑”“谁把同一张表跑重了”。前前后后折腾过好几套方案之后,我整理出一套适合自己团队的调度设计,内部代号就是“ax”,我们一般叫它 ax调度。它不是什么黑科技,核心就解决三件事:任务什么时候触发、交给哪台机器执行、失败之后怎么办。这篇文章把 ax调度的设计思路、实现要点和落地过程完整拆一遍,适合正在用 Quartz 救火、或者准备自研调度系统的后端同学参考。
1. 为什么我不直接用现成的调度框架,而是自己搞了套 ax调度
每次聊任务调度,第一反应肯定是“为什么不用现成的”。市面上 Quartz、ElasticJob、XXL-JOB 都成熟得不行,直接拿来用不香吗?我一开始也是这样想的,但真正在业务里跑一段时间后,发现“能用”和“好用”之间有很长一段路。
1.1 Quartz 的分布式困境
Quartz 在单机时代是绝对王者,但一旦上集群,痛点就非常明显。它的分布式方案是靠数据库行锁实现的:调度线程定时去 triggers 表里抢记录,抢到锁的节点才有资格执行触发逻辑。这套机制在小规模场景没问题,可当触发器表到了几万条、数据库连接一旦紧张,锁竞争就会让调度延迟从秒级恶化到分钟级。我在压测环境里实测过,仅仅 50 个并发调度线程、单表 5 万条 trigger 记录,就有接近 3% 的触发动作延迟超过 30 秒。对于财务对账、库存同步这类对时效敏感的任务,30 秒意味着业务已经出了可感知的误差。
更头疼的是双机房部署时的网络抖动。节点 A 刚拿到锁准备触发,节点 B 因为网络分区误判自己失联,转而竞争同一把锁;等网络恢复,两个节点同时认为自己是 Leader,结果就是同一个任务被触发两次。Quartz 本身没有很好的幂等保护,业务侧只能自己加分布式锁,这个成本往往被很多人低估。
1.2 大而全调度平台的改造成本
后来我也认真调研过 XXL-JOB 这类功能齐全的调度平台。它们确实很强,有可视化控制台、有权限体系、有动态任务管理,拿来即用。但我们的团队规模并不大,核心服务不到 20 个,真要引入一套完整的平台,学习成本和维护成本会压在本来就紧张的研发资源上。控制台、权限、多租户、归档报表这些功能,对我来说属于“偶尔用到、平时添乱”的范畴。
而且这类平台的扩展点虽然多,但真要改核心调度逻辑,需要熟读大量源码,对团队里每个接手的人来说都是负担。我需要的不是大而全,而是小而可靠:调度器专注时间计算和任务分发,执行器专注跑任务和回传状态,中间通信尽可能薄。这就促成了 ax调度的雏形。
1.3 ax调度的设计取舍
ax调度的架构就三个角色:调度中心(ax-server)、执行器(ax-worker)、元数据库。调度中心只干两件事:维护任务元数据、计算触发时间并下发执行指令。执行器只干三件事:注册自身地址、接收指令执行任务、上报日志和结果。通信走 HTTP 端点,请求携带全局唯一的 batchId(执行批次号),从触发到执行结束,全链路可以通过 batchId 把日志串起来。
这套拆分的直接收益是执行器可以做到语言无关。调度中心下发的是通用指令,任何语言只要实现一个 HTTP 回调接口就能接入。我在团队里就有两个 Python 写的任务进程,注册方式跟 Java 服务一样,只是回调接口略有区别。这个特性在微服务多语言环境下特别实用。
2. 核心细节解析与实操要点
很多调度框架的原理文档写得抽象,落地时全靠自己猜。下面把 ax调度里我认为最核心的几个模块拆开讲,偏实现向,尽量说清楚“为什么这么设计”。
2.1 调度引擎的“多层时间轮 + 延迟队列”混合模型
时间轮是解决大量定时任务扫描问题的经典方案,核心思想是把时间分成一个个 slot,用数组存待触发的任务,指针每 tick 跳一个 slot,命中哪个槽就触发哪个槽上的任务链表。相比每秒钟全表扫一遍数据库,时间轮在插入和触发上都能做到 O(1) 级复杂度。
ax调度里用的是两层时间轮:第一层 512 个 slot,每个 slot 对应 200ms,所以能覆盖 102.4 秒;第二层 60 个 slot,每个 slot 对应 1 分钟,用于存放更远的触发任务。第一层指针转完一圈,会把第二层中最近一分钟的任务降级到第一层里。这种设计避免了单层时间轮内存占用过大的问题。
但纯时间轮有个缺陷:任务量稀疏的时候,指针空转浪费 CPU。所以在 ax调度里增加了一个容量约 1000 的 DelayQueue,用来缓存最近 5 分钟内需要触发的任务。调度线程从 DelayQueue 里阻塞取最近到期的任务,到期后再根据任务精确时间放入时间轮。这样即使任务数量很少,线程也能休眠等待而不是空转。
这里有一个很容易踩的坑:调度中心的时间精度不能完全依赖系统 tick。调度进程需要对 NTP 时间同步敏感,一旦宿主机时钟漂移超过 500ms,调度触发的准确性就无从谈起。ax调度专门加了一个“时钟偏差检测”模块,每隔 30 秒跟元数据库时间做对比,偏差过大直接在当前任务日志里打 WARN。
2.2 任务分片、路由与故障转移策略
任务触发之后,调度中心要决定把这批任务发给谁。ax调度支持三种基本路由策略:轮询、一致性哈希、分片广播。
轮询适合任务短时间内多、但单个执行成本低的场景,比如批量发送通知,自动均匀分散到各个 worker。一致性哈希适合需要把同一业务维度积累到同一节点的场景,比如按用户 ID 做缓存预热,保证同一个用户总是落在同一台执行器上,提高本地缓存命中率。分片广播则是我用得最多的策略:把一个大任务切分成 M 个分片,每个执行器拿到分片号(shardingId)和总分片数(shardingTotal),自行处理属于自己的一部分数据。
分片任务最大的坑是数据边界不清晰。如果分片逻辑是“订单号取模后等于当前分片号”,那新增 worker 会导致分片总数变化,老任务可能漏数据。所以我在 ax调度里规定:执行器启动后分片总数是「当前在线 worker 数 × 每个 worker 允许并发数」,动态扩容时并不立即改变已经下发任务的 shardingTotal,而是等下一轮触发再生效。虽然短期可能存在资源不均,但换来了数据一致性。
故障转移的处理我设计成“先标记后剔除”。执行器每 10 秒发一次心跳,调度中心连续 3 次没收到心跳,就把它标记为“可能失联”,此时不立即摘除,而是再等 15 秒尝试一次主动探测。如果探测失败,才从在线列表摘除,并把尚未完成执行的批次重新路由到其他节点。这种两段式方案能有效避免因为网络瞬断导致的大规模任务迁移。
2.3 触发规则、线程池和重试参数怎么定
调度系统的参数配置如果不解释清楚,就是黑魔法。先说触发规则,ax调度支持 cron 表达式和固定间隔两种。cron 表达式解析后先转成下一次触发时间点,再通过时间轮排程。值得注意的是,cron 表达式解析要显式指定时区,默认用 Asia/Shanghai,不要用服务器系统时区,否则服务器时区被改一下,所有任务秒变“时差任务”。
执行器收到指令后,会丢进自己的业务线程池。线程池参数我一般这样建议:核心线程数 = CPU 核心数 × 2 + 1,最大线程数 = 核心线程数 × 2,有界队列容量 1024,拒绝策略用 CallerRunsPolicy 的改进版:不阻塞本次任务执行,而是立即把失败原因写进 execution_log 表,同时触发告警。
这里要特别说明,调度场景最怕的是“队列积压但任务还在排队”的假象。如果你把队列设成无界,遇到下游接口变慢,任务会一直在队列里堆着,执行时间远晚于调度时间,排查时根本看不出是哪一环拖的。有界队列配合拒绝策略,至少能让问题立刻暴露出来,而不是让延迟“温水煮青蛙”。
重试策略我经历过从“无脑重试”到“谨慎重试”的转变。默认重试 3 次,重试间隔指数退避:1 秒、2 秒、4 秒。超过 3 次就不重试,转人工/告警。为什么不用更大的重试次数?因为任务失败的重试是有成本的,如果目标系统已经过载,重试只会加剧雪崩;如果任务是不可幂等的(比如发短信),重试还会造成重复发送。所有重试逻辑必须配合幂等令牌使用,这个令牌就是 batchId。
2.4 两套核心表结构设计
调度中心能稳定运行,元数据表结构非常关键。ax调度只保留两张核心表:任务表 job_info 和 执行日志表 execution_log。
job_info 主要字段包括:任务主键、应用名 app_name、处理器标识 handler、cron 表达式、路由策略、超时时间、重试次数、是否开启并发执行、最近触发时间、下次触发时间、状态。索引重点落在 app_name 和 next_trigger_time 上,调度扫描时就靠这两个字段定位任务。
execution_log 主要字段包括:日志主键、任务主键、batch_id、触发时间、实际执行时间、worker 地址、执行状态、耗时、错误信息、日志快照。索引建议落在 job_id + batch_id 和 trigger_time 上。batch_id 一定要有唯一索引,这是防止调度中心重复下发的最后一道防线:即使调度侧逻辑出错,插入同一 batch_id 也会失败,从而中止重复执行。
表设计要说一个血泪经验:不要把任务详情和执行日志混在一张表。最开始我图省事,把 handler 参数直接拼在任务表的字段里,结果日志查询稍微变多,任务表就出现锁竞争,调度扫描也被拖慢。拆成两张表后,调度链路的查询路径非常干净,日志表的膨胀也不会影响调度性能。
3. 实操过程与核心环节实现
讲完原理,进入正文操作。以一个订单服务为例,从零搭建 ax调度环境、注册执行器、跑通第一个调度任务。
3.1 十分钟搭起调度中心环境
调度中心我打包成了一个独立可执行 Jar,依赖一个 MySQL 库。启动前先初始化 SQL 脚本,创建 ax_dispatch 数据库,脚本里建好 job_info 和 execution_log 两张表以及基础索引。
mysql -u root -p -e "CREATE DATABASE ax_dispatch DEFAULT CHARACTER SET utf8mb4;" mysql -u root -p ax_dispatch < sql/init.sql然后启动调度中心:
wget https://mirrors.example.com/ax-server-1.0.0.tar.gz tar xzf ax-server-1.0.0.tar.gz cd ax-server sh bin/start.sh --server.port=8081 --spring.datasource.url="jdbc:mysql://127.0.0.1:3306/ax_dispatch"启动后访问 http://127.0.0.1:8081/actuator/health ,返回 UP 说明调度中心起来了。调度中心自身不存储执行逻辑,只负责计算、下发行令和维护 worker 列表,所以它的重启对已注册执行器是无感的,正在执行的任务不会被中断。
3.2 注册执行器并跑通第一个任务
执行器端接入很简单。我是 Maven 项目,加依赖:
<dependency> <groupId>com.ax</groupId> <artifactId>ax-executor-spring-boot-starter</artifactId> <version>1.0.0</version> </dependency>配置文件 application.yml 里加:
ax: executor: app-name: order-service address: 127.0.0.1:9099 registry-url: http://127.0.0.1:8081/ax/registry注意 address 是执行器对外暴露的回调地址,在容器环境里要配置成宿主 IP 和映射端口,不能配成 pod 内部 IP,否则调度中心无法回调。这个我在后面容器化章节还会再提。
写一个最简单的任务:
@Component public class HealthCheckJob implements SimpleJob { @Override public JobResult execute(JobContext ctx) { log.info("health check start, batchId={}", ctx.getBatchId()); return JobResult.success(); } }执行器启动后,日志中会出现“register to scheduler success”的字样。接着在调度中心控制台新建任务:应用名 order-service、处理器标识 healthCheckJob、cron 表达式每 5 分钟一次,保存后即可看到下一次触发时间。到点后,执行日志里会出现 batchId、起止耗时和状态码。到这里,第一个 ax调度任务就跑通了。
3.3 分片任务:批量订单数据处理的完整示例
光跑通 Hello World 没意思,看一个真实场景:每天晚上要同步前一天的订单数据到数仓,订单量有几百万,单台机器跑要几个小时,必须分片。
任务实现采用分片模式:
@Component public class OrderSyncShardingJob implements ShardingJob { @Override public JobResult execute(JobContext ctx) { int shardId = ctx.getShardingId(); int shardTotal = ctx.getShardingTotal(); int batchSize = 5000; long minOrderId = ctx.getJobParam().getLong("minOrderId"); long maxOrderId = ctx.getJobParam().getLong("maxOrderId"); // 每个分片只处理属于自己范围内的订单 long rangeLen = (maxOrderId - minOrderId) / shardTotal; long startId = minOrderId + rangeLen * shardId; long endId = (shardId == shardTotal - 1) ? maxOrderId : startId + rangeLen; List<OrderDO> orders = orderMapper.selectRange(startId, endId, batchSize); for (OrderDO order : orders) { syncToWarehouse(order, ctx.getBatchId()); } return JobResult.success(); } }控制台配置路由策略为“分片广播”,假设当前 4 台 worker,任务触发后每台机器拿到的 shardingId 分别是 0、1、2、3,shardingTotal = 4。由于上面代码里每个分片只处理订单 ID 区间的一部分,合计起来正好覆盖全量,不会重复也不会漏单。
实际跑起来我遇到过一个边界问题:如果订单 ID 不是连续分布(中间有空洞),按 ID 区间分片会导致分片 1 可能比分片 2 多处理几十万条。后来我改成先统计 min/max,再在任务参数里携带全量 ID 集合的布隆过滤器,每个分片查询时本地过滤一遍空洞。虽然多了一点内存消耗,但负载均衡效果好了很多。
3.4 配置监控告警,提前发现问题
分布式调度最怕“没跑”和“跑重”,而这两类问题靠日志事后查是非常痛苦的。ax调度通过暴露 Prometheus 指标让你提前感知风险:
# HELP ax_dispatch_task_trigger_total 触发任务总数 # TYPE ax_dispatch_task_trigger_total counter ax_dispatch_task_trigger_total{app="order-service"} 1280 # HELP ax_dispatch_task_failure_total 失败任务总数 # TYPE ax_dispatch_task_failure_total counter ax_dispatch_task_failure_total{app="order-service"} 3在 Prometheus 里配置抓取调度中心的 /actuator/prometheus 端点,然后配两条核心告警规则:一是调度延迟:任务实际执行时间减去触发时间超过 15 秒触发 Warning;二是失败率:5 分钟内失败率超过 10% 触发 Critical。我自己的习惯是,给每个重要任务额外加一条“SLA 未触发告警”——如果一个任务在预定窗口内完全没有触发记录,说明调度链路本身可能挂了,这比失败告警能更早暴露问题。
4. 常见问题与排查技巧实录
不管设计多完整,落地过程总会有各种幺蛾子。这里把我在 ax调度使用中踩过的坑集中列一下,附上排查思路,方便遇到类似问题的人直接定位。
4.1 任务总是延迟几十秒才触发
现象:cron 设的是整点触发,但 execution_log 里记录的实际执行时间比触发时间晚了 30 秒甚至 1 分钟。排查第一步看调度中心的线程池是否打满。ax调度的下发动作会经过一个独立的 IO 线程池,如果业务任务里有人直接写了阻塞式 HTTP 调用且超时设成 60 秒,这个线程池会被占满,后续任务全部排队。
解决办法有两个层面:任务代码里禁止在调度执行线程中调用不可控外部接口,改成异步提交;把同步下发改成“批量聚合下发”,同一秒触发的任务合并成一条请求发给执行器,降低 IO 次数。我在 4.1 节中实测发现,聚合下发能减少约 40% 的下发延迟波动。
另一个隐蔽原因是线程池配置把核心线程数设得太小,默认 2 个线程扛不住瞬时 20 个任务同时触发的尖峰。调度线程池的队列要有界,但线程数不能太保守,建议按“单机任务触发峰值 × 1.5”来设计。
4.2 同一个任务被重复执行
重复执行是调度系统最严重的故障之一。我遇到过一次典型的重复场景:调度中心在任务执行超时后判定失败,立即重试路由,但第一次执行的线程并没有被真正终止(数据库查询卡住了而已),导致同一个批次两个 worker 同时在跑。这是“超时误判 + 不可中断操作”共同作用的结果。
对应的防御手段我总结成三个层级:第一层,数据库约束。execution_log 表为 job_id + batch_id 建立唯一索引,重复下发会直接插入失败;第二层,执行器内部并发控制。每个任务在执行前先对 job_id 加分布式锁,锁的 key 是 job_id + batch_id,用 Redis 的 SETNX 实现,过期时间设为任务超时时间的两倍;第三层,严格禁止并发。在任务配置里强制开启“并发运行阻止”,同一任务即使被触发多次,也只会有一个执行中的实例。
三层都做了,也不能说 100% 安全,但至少能挡住绝大多数误操作。
4.3 长任务导致的任务堆积
有些任务是跑 3 小时才算完的批处理任务,如果 cron 是 2 小时一次,第二次触发时第一次还没结束。任务会越积越多,每个任务还都在抢数据库连接,最终拖垮整个服务。
ax调度里我建议对这类任务单独设置“最大执行时间”,超时后触发中断信号。注意中断不是简单的 Thread.stop,而是标记线程中断状态并把执行结果写成“超时失败”,让业务代码在数据库操作处感知到异常退出。如果中断后任务还是没停(比如卡在了不可中断的锁等待),则配合线程池的 discardPolicy 强制丢弃该批次的新触发请求。
更治本的手段是分片 + 增量。把全量同步改成上一天增量,用 offset 记录进度,每次任务只处理新增数据,这样单次执行时间能压到 10 分钟以内,彻底摆脱“任务跑不完”。
4.4 调度中心与执行器时钟偏差
分布式的世界里“时间”是个伪命题。我遇到过一个问题:执行日志里显示任务实际执行时间比触发时间早 2 秒,导致链路追踪的时间线倒挂,排查问题的时候非常误导。原因是执行器的宿主机时钟快了 2 秒,它记录的是本地时间,而调度中心记录的是自身时钟时间。
解决方案是统一以调度中心的时间作为事实来源。触发时间、计划时间都用调度中心生成的时间戳,执行器只记录 CPU 耗时、内存等相对指标,不再记录本地时间用于链路排障。需要真实时间做业务判断时,从调度中心下发的报文头里取,而不是本地 System.currentTimeMillis()。同时所有机器必须强制配置 NTP 同步,偏差超过 500ms 时监控告警。
4.5 调度问题排查速查表
我把上面这些经验浓缩成一张排查表,贴在团队内部文档里,遇到问题直接对照:
| 现象 | 可能原因 | 排查方向 |
|---|---|---|
| 任务触发延迟大 | 调度中心线程池阻塞 | 查调度中心活跃线程数、队列堆积数 |
| 同一任务执行多次 | 超时误判、路由重复下发 | 查 execution_log 中 job_id + batch_id 是否唯一 |
| 任务长时间排队 | 执行器线程池满或任务不可中断 | 查执行器活跃线程数、等待队列长度 |
| 执行时间倒挂 | 执行器本地时钟漂移 | 统一走调度中心时间戳,校准 NTP |
| 日志数据缺失 | 日志上报链路被阻塞 | 查执行器到调度中心日志端点的网络状态 |
| 任务分片数据不均 | 分片区间数据分布不均 | 改用一致性哈希或实现消费进度记录 |
这张表最大的价值是帮团队节省了“猜”的时间。遇到调度问题,先看表中对应的最可能原因,再通过命令查指标,基本能在 5 分钟内定位到问题层。
5. 容器化部署与容量规划
现在部署基本都往 K8s 走了,调度系统也要适配容器环境。这一章把我自己在容器化落地过程中的要点写清楚。
5.1 K8s 下怎么部署 ax调度
调度中心无状态,可以做成 Deployment 跑两个副本,通过 Service 对外暴露端口。但底层依赖的 MySQL 不建议也容器化,至少生产环境用云数据库会更稳。调度中心多副本部署时,副本之间会通过数据库锁做任务触发互斥,所以不用担心两个副本同时触发同一个任务;但如果用了数据库锁,连接池大小要给足,否则高并发下锁等待会拖慢调度。
执行器部署稍微复杂点:每个 Pod 就是一个执行器实例,注册地址必须配成宿主机 IP 加 NodePort 映射出来的端口。如果直接用 Pod IP,调度中心从集群外访问不到。我用的是 StatefulSet 部署执行器,每个 Pod 有个稳定网络标识,配合 headless service 做回调。Pod 生命周期结束时,preStop 钩子里调用执行器提供的 unregister 接口,把这个节点从调度中心在线列表里摘除,避免已经下线的执行器继续被路由。
lifecycle: preStop: exec: command: ["sh", "-c", "curl -X POST http://127.0.0.1:9099/ax/executor/unregister; sleep 5"]这个 5 秒的 sleep 是为了等调度中心完成摘除操作,再真正终止容器,否则正在执行的任务会被强制杀掉。
5.2 多租户与权限隔离
团队大了之后,不同业务线都要用调度平台,直接混在一个命名空间里很不安全。ax调度在应用层面支持简单的多租户隔离:每个任务绑定一个 namespace,执行器注册时声明自己属于哪个 namespace。调度中心只把任务路由到相同 namespace 下的执行器。
控制台权限上,按角色分为管理员、开发者、只读三种。管理员能管理所有 namespace;开发者只能在自己的 namespace 里创建和修改任务;只读只能查看日志和监控。这种隔离粒度对中小团队足够用了,没必要一开始就上完整的 RBAC 权限系统,不然又是维护负担。等需要更细粒度的权限时,再基于 OpenID Connect 做一层对接也不迟。
5.3 生产环境需要多大规格的机器
容量规划不能靠感觉。我自己跑过一轮压测,作为参考:调度中心使用 8C16G 的虚机、MySQL 使用 4C8G 的云数据库,单调度中心能支撑 5000 个任务,峰值触发频率约 120 次/秒,端到端调度延迟 P99 在 80ms 以内。这个数据是在任务逻辑都比较简单的情况下测的,如果你的执行器处理很重,瓶颈一般不在调度中心而在执行器线程池。
建议的规划比例是:一个执行器实例最多注册 200 个任务,如果超过 200 个,就拆成多个执行器应用。一个调度中心最多支撑 20 个执行器应用在线,超过之后建议做多调度中心分片,按业务线拆成多个独立集群。这样出问题时影响面能控制在一个集群内,排查也方便。
我个人在实际操作中的体会是,调度系统最怕的不是负载高,而是你把所有鸡蛋放在一个篮子里。宁可初期拆成两个小集群,也不要等到一个集群里几百个任务互相影响时再迁移。另外,最后分享一个小技巧:给每个任务设置一个“SLA 窗口”,如果任务在这个窗口内没有被触发就立即告警,这比观察失败率更能提前发现调度链路的问题。把调度当第一公民来看待,而不是最后一环,系统稳定性的提升是肉眼可见的。