我平时不太喜欢追热点,但“ax调度”这个词在圈子里连续几天被刷到之后,我还是没忍住去翻了翻上下文。结果发现大家讨论的并不是什么神秘的新框架,而是一个很典型的场景:业务起来了、任务变多、定时器越来越乱,然后在某次大促或活动流量冲击下,线上调度直接崩了。我自己的项目里也遇到过一模一样的事故,后来花了一个周末重构出一套轻量级调度内核,代号就叫 AX(取 Agile eXecution 的意思)。这篇文章就把我从需求拆解、模型设计到落地实现的完整过程写出来,代码都是可以直接拿回去改的版本,适合正在被 Cron 和零散任务折磨的团队参考。
先交代一下背景。AX 不是一个大而全的分布式工作流平台,它解决的痛点非常聚焦:定时任务、延迟任务、失败重试,以及最重要的“谁先跑谁后跑”。大多数业务系统早期都用 Linux Crontab 或者应用内置的 Schedule 组件,任务一多就会出现三个问题——任务互相挤占、重复执行、失败后没人管。AX 的定位就是把这堆烂摊子收拢起来,用一个足够简单的调度模型去承接,同时把性能控制在一个合理的范围内。它的核心代码量不大,但是设计上能扛住数万级任务、秒级调度延迟,这也是我反复测试之后比较满意的地方。
1. 从一次线上事故说起:AX 调度器要解决什么问题
1.1 根因不是一个 Bug,是调度模型缺失
那段时间我们内部有个数据对账系统,每天凌晨要跑上百个任务:拉取流水、清洗数据、更新报表、推送通知。最初是用一个线程池加上@Scheduled注解硬撑,任务多了以后开始加队列,然后有人为了赶时间在代码里写了Thread.sleep,有人把重试逻辑堆在 catch 块里,最后整个系统变得像一团乱麻。有一天上游系统延迟,导致凌晨的任务全部积压,正常上班时间用户还在收到昨晚的对账通知,那次事故让我彻底下决心做自己的调度器。
排查现场的时候我发现,问题根本不是某个任务写错了,而是整个系统缺少“调度语义”。什么叫调度语义?就是一个任务应该什么时候触发、触发之后由谁执行、执行失败要不要重试、多个任务同时就绪时该优先跑哪个、任务执行超时怎么办。Cron 表达式只回答了“什么时候”这一个问题,剩余的四个问题完全没有覆盖。所以 AX 的第一版需求列表直接把这些空白全部填上,不追求花哨的 DAG 编排,先把单任务的触发、排队、优先级和重试做好。
1.2 为什么不自研还要自研:三个硬约束
当时团队里也有同事提议直接引入现成的开源调度框架。说实话现在能用的方案很多,但如果仔细对照业务需求,会发现它们都有一个尴尬的错位。首先是部署成本,很多重型调度平台依赖独立的存储、独立的 Web 控制台、特定的运维流程,我们团队只有三四个后端,没有专职运维,引入一套体系等于给自己造一个新系统要维护。其次是定制成本,我们需要的“优先级抢占”和“任务依赖”这两个能力在通用平台里往往都是高级功能,配置起来不直观,出了问题反而不容易排查。最后是调度延迟,有的平台最短调度周期只能到分钟级,可我们有些延迟任务需要秒级触发。
自研的风险我当然也评估过。调度器是典型的“看起来简单、做起来容易翻车”的组件,分布式锁、时间轮、故障转移、幂等控制,任何一个细节没写对都会埋坑。但换个角度想,我需要的不是一个完整的调度平台,而是一个内核:把调度决策和执行分离,存储层尽可能薄,逻辑层全部自己掌控。这样出了问题我能直接看懂日志,能改代码,不用猜黑盒。AX 的命名也是这个意思——我们不追求管理系统,我们只做那个“敏捷的执行内核”。
2. AX 的整体架构与调度模型设计
2.1 三层职责拆分
AX 从逻辑上拆成三层:触发层、调度层、执行层。触发层负责接收外部信号,包括 Cron 时间到点、HTTP 调用、内部事件通知,以及延迟任务的时间到达。调度层是核心,它从触发层拿到“任务可以开始”的信号,然后根据任务的优先级和当前系统的负载情况,决定把任务放进哪个执行队列。执行层是真正跑业务逻辑的地方,它不在调度进程内部跑任务,而是通过一个 Worker 去拉取任务,这样调度器就算重启,Worker 上的任务也不会丢失。
这三层的边界一开始就必须划清楚,否则很容易退化成一个分布式线程池。触发层不关心任务怎么跑,调度层不关心业务逻辑,执行层不关心任务什么时候触发。调错的时候只需要判断:是没触发?触发了没进队列?进队列了没被执行?还是执行了但结果没回写?每一层有单独的日志前缀,排查效率能提升好几倍。
2.2 核心设计:时间轮 + 优先级队列 + 乐观锁
AX 内部最关键的三个组件,我分别用不同的策略实现。时间轮用来处理定时任务和延迟任务,它本质上就是一个环形数组,每个槽位放一批到期的任务指针,时钟每走一格就去触发对应的槽位。相比直接用数据库轮询扫描,时间轮的触发延迟更稳定、也几乎没有空转成本。优先级队列用来处理同时就绪的任务,底层用堆实现,每个任务有一个动态优先级数值,调度循环每次从堆顶取任务。乐观锁用来保证集群环境下多个调度节点不会把同一个任务重复发给 Worker。
这三个组件串起来的逻辑是这样的:任务注册时先放进时间轮,时间到达之后被推入优先级队列;调度循环从优先级队列弹出堆顶任务,通过 Redis 的SETNX做一次加锁标记;拿到锁的任务写入执行队列,等待 Worker 消费。这里最关键的设计是“加锁只是标记、不是占用”,任务不是被调度器主动推给 Worker 的,而是 Worker 主动来拉取,拉取时再次校验任务状态。这样可以避免调度器崩溃导致任务永远卡在内存队列里。
2.3 几个关键技术选型的取舍
我最早其实想过用纯 Redis 的 ZSET 来做延迟队列,Score 存触发时间戳,再用一个轮询线程去查。这在任务量不大的时候完全够用,但 ZSET 的 range 操作在高并发下会有性能瓶颈,而且 Redis 一旦重启,没有落盘的延迟任务就直接丢了。后来折中了一下:定时任务的元数据始终放在数据库里,Redis 只做“唤醒通知”和“分布式锁”,就算 Redis 重启,数据库里的任务还在,下一次调度周期依然会触发。这个取舍牺牲了一点点实时性,换来了很高的可靠性。
还有一个取舍是执行结果的回写方式。AX 用的是数据库状态流转:任务被调度时置为RUNNING,Worker 执行结束之后回调接口把状态改成SUCCESS或FAILED,失败的任务根据重试策略重新计算下次触发时间。有人觉得这样太慢,应该用消息队列,但实际跑下来数据库状态流转完全够用——因为任务执行的耗时通常以秒为单位,状态回写的频率远低于任务执行本身的频率。把数据库当状态机用,天然就有了审计日志和恢复依据。
3. 核心实现细节与可直接抄作业的代码
3.1 任务表结构设计
AX 的任务表是核心的数据模型,我一开始就把所有调度需要的字段都塞进去了,但后续迭代中真正用得上的就那么几个。表结构如下:
CREATE TABLE `ax_task` ( `id` BIGINT PRIMARY KEY AUTO_INCREMENT, `task_key` VARCHAR(128) NOT NULL COMMENT '业务唯一标识,用于幂等', `type` TINYINT NOT NULL COMMENT '任务类型:1-定时任务 2-延迟任务 3-事件任务', `status` TINYINT NOT NULL DEFAULT 0 COMMENT '0-待调度 1-运行中 2-成功 3-失败 4-取消', `priority` INT NOT NULL DEFAULT 100 COMMENT '数值越小优先级越高', `trigger_at` BIGINT NOT NULL COMMENT '计划触发时间戳(ms)', `next_retry_at` BIGINT DEFAULT 0 COMMENT '下次重试时间戳(ms)', `retry_count` INT NOT NULL DEFAULT 0 COMMENT '已重试次数', `max_retry` INT NOT NULL DEFAULT 3 COMMENT '最大重试次数', `handler` VARCHAR(255) NOT NULL COMMENT '执行器标识', `payload` JSON DEFAULT NULL COMMENT '业务参数', `version` INT NOT NULL DEFAULT 0 COMMENT '乐观锁版本号', `created_at` DATETIME DEFAULT CURRENT_TIMESTAMP, `updated_at` DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, UNIQUE KEY `uk_task_key` (`task_key`), KEY `idx_trigger` (`status`, `trigger_at`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;关键就在uk_task_key这个唯一索引上。它为整个系统的幂等提供了最后一道防线。无论调度器发了多少次调度信号,同一个task_key在数据库里只能有一条待调度记录,重复插入直接报错。version字段我用来做乐观锁,Worker 回写状态时带上读取到的 version,更新时校验 version,防止旧数据覆盖新数据。
3.2 调度主循环实现
调度循环是整个 AX 的心脏,我用 Python 演示核心逻辑,换成 Go 或者其他语言思路完全一致。这个循环每隔 200ms 跑一次,每次做三件事:把时间轮里到期任务推入优先级队列、从优先级队列弹出任务并尝试加锁、把加锁成功的任务状态更新到数据库。整体结构非常简洁:
import time import redis from queue import PriorityQueue from datetime import datetime class AXScheduler: def __init__(self, db, redis_client): self.db = db self.redis = redis_client self.ready_queue = PriorityQueue(maxsize=10000) self.schedule_lock = 'ax:schedule:lock' def schedule_once(self): # 1. 将时间轮到期的任务拉入优先级队列 due_tasks = self.pop_due_tasks_from_timing_wheel() for task in due_tasks: self.ready_queue.put((task['priority'], task['id'], task)) # 2. 一批批处理队列中的任务 while not self.ready_queue.empty(): _, task_id, task = self.ready_queue.get() # 3. 用 Redis 做分布式加锁,避免多节点重复调度 lock_ok = self.redis.set( f"ax:lock:task:{task_id}", "1", nx=True, ex=max(30, task.get('timeout', 60)) ) if not lock_ok: continue # 4. 乐观锁更新数据库状态 rows = self.db.execute( "UPDATE ax_task SET status=1, version=version+1 " "WHERE id=%s AND version=%s AND status=0", (task_id, task['version']) ) if rows == 0: continue # 5. 写入 Worker 可拉取的执行队列 self.redis.lpush(f"ax:queue:worker:{task['handler']}", task_id) def run_forever(self): while True: try: self.schedule_once() except Exception as exc: # 落盘错误日志,别吞异常 self.log_error(exc) time.sleep(0.2)每一步都有对应的失败分支:时间轮到期但队列已满就暂时跳过,等下一轮;Redis 加锁失败说明别的节点已经处理了;数据库乐观锁更新失败说明任务状态已经被改变。这三个失败分支基本覆盖了集群环境下 99% 的重复调度问题。有一个细节是加锁的过期时间要带上任务超时时间,且至少 30 秒,避免任务还没跑完锁就过期,导致另一个调度器又把任务捞起来跑一遍。
3.3 Worker 执行与幂等控制
Worker 侧的逻辑比很多人想象中要简单。它不需要实现任何调度算法,只需要做三件事:从队列拉任务、执行业务代码、回写执行结果。但回写结果之前的幂等校验绝对不能省,因为网络超时可能导致 Worker 已经执行成功,但回写接口的请求没到达调度器,调度器就会判定失败并触发重试。这种情况下同一笔业务会被执行两次。
我在 Worker 里加了一个去重表,专门记录最近执行过的任务指纹:
def execute_task(task_id, task_key, handler, payload): if not can_execute(task_id, task_key): return mark_start(task_id, task_key) try: result = handle(handler, payload) report_success(task_id, result) except Exception: report_failure(task_id, exception_to_string()) finally: mark_finish(task_id, task_key) def can_execute(task_id, task_key): # 1. Redis 里有执行标记,说明刚执行过,直接忽略 # 2. 数据库里有 task_key 的最近成功记录,也直接忽略 pass这套实现的代价是每个任务多一次 Redis 查询和一次数据库查询,换来的是整个系统的幂等保证。对于很多电商场景来说,重复发一次通知可以忍,重复扣一次库存是不能忍的。宁可多一次查询,也不要让重复消费的风险敞口开着。
4. 线上踩坑实录:排查与调优经验
4.1 重复调度:从 Redis 锁到唯一索引
AX 上线后的第一周,我就在监控里发现有个别任务在凌晨被触发了两次。第一次排查时我以为是时间轮的指针重复拨动导致的,后来加上日志才发现根本原因是 Worker 执行超时,调度器判定任务失败后重试,而第一次执行其实在超时前的最后一刻已经完成了。业务方看到的结果是“同一个任务跑了两遍,产生了重复数据”。
解决办法是双管齐下。首先在任务表上加uk_task_key唯一索引,让业务层面的同名任务在数据库层面就只能存在一个;然后在 Worker 里加了一个任务去重表,以“执行器标识 + 当天日期 + task_key”作为去重维度,既能防止跨天重试带来的重复,又不会让去重表无限膨胀。加完这两个东西之后,重复调度的告警就彻底消失了。
4.2 队头阻塞与任务饿死
上线之后又遇到一个新问题:高优任务一直能跑,低优任务偶尔被饿死。原因是优先级队列用的是静态优先级数值,某个业务方把自己的任务优先级都调到了 1,大批普通任务虽然也到了触发时间,但堆顶永远是那几个 1,普通任务只能排队。尤其当高优任务执行时间比较长的时候,低优任务可能等上十几分钟甚至更久。
这让我意识到光有“优先级”还不够,还得有“老化机制”。我在任务对象里加了一个进入队列的时间戳,每次计算优先级时做一次惩罚性降级:等待超过 30 秒的任务,优先级数值每多等 30 秒就减 1。这样即使高优任务源源不断进来,低优任务也会随着等待时间上升而逐渐获得调度机会。这个机制实现成本很低,但直接解决了任务饿死的问题。
4.3 故障转移与漏调度
集群部署之后还有一个很隐蔽的坑:调度器有多个节点,如果某个节点在加锁之后、更新数据库之前宕机了,Redis 锁会持有 60 秒,这段时间内其他节点都会认为任务正在处理中,不会接手。等锁过期之后任务状态还是 0,才会被再次调度。这里的延迟最多达到锁的过期时间,对于秒级任务来说是不可接受的。
我的处理方式是给调度节点加了一个“自愈检查”协程:每隔 5 秒扫描一次“正在运行中但超过心跳时间没有更新”的任务,直接把状态重置为待调度。每个调度节点写心跳,Worker 执行前也写心跳,超过两倍心跳周期没有更新的任务就被视为“僵尸任务”,允许被重新调度。这套心跳检查机制和 Redis 锁配合起来,既防止了重复调度,又防止了漏调度。
4.4 监控与报警阈值怎么设置
调度器本身的监控指标不需要多,我最终只保留了五个核心指标:调度延迟(从触发时间到进入执行队列的时间差)、执行耗时、失败率、队列积压量、心跳超时数。报警阈值我调了一个多月才稳定下来:调度延迟的报警线是 5 秒,超过就说明时间轮或者优先级队列可能卡了;失败率是相对值,用近 10 分钟内的失败任务数除以总执行数,超过 5% 触发;队列积压量超过 2000 就报警,因为积压通常意味着 Worker 数量不够或者执行耗时暴涨。
有一个最容易忽略的指标是时间轮的时钟拨动。如果部署环境有人手动改系统时间,哪怕是改回 5 分钟,也可能会导致定时任务瞬间全部到期、同时涌入队列。AX 在调度循环里记录每次执行时的系统时间偏移,如果发现偏移量超过 2 秒就输出告警日志,虽然没办法完全避免环境问题,但至少能第一时间定位到“所有任务突然集中到点”的原因。
5. ax调度背后的思考与实战建议
5.1 调度的本质与边界
把 AX 完整跑起来之后,我对调度这件事本身的理解深了一层。所谓调度,本质上就是“在资源有限、时间有限、需求多变的情况下,做出执行顺序的决策”。这和我们每天在生活里做的排队和优先级判断一模一样:食堂窗口再多也总要排队,紧急的订单先做,卡了很久的单子要人工介入。调度器做的就是把这一套规则固化成代码,用统一的方式去处理数万个任务之间的竞争关系。
但调度器也有明确的边界。它不负责业务代码的正确性,不负责执行器的资源隔离,也不负责最终的数据一致性。过度依赖调度器去兜底业务问题,一定会把调度器本身拖垮。AX 能保持简单的最大原因就是我不停地拒绝新需求:任务编排不做、流程审批不做、执行日志的全文检索不做,这些功能周边系统更擅长。守住边界,核心调度逻辑才能一直稳定。
5.2 给你的三点实战建议
根据我自己从设计、编码到运维 AX 的完整过程,有三条建议想送给准备自研调度器的人。第一条,别一上来就做分布式,先把单机版本跑通、把状态流转理清楚,分布式是在单机无法承受之后才需要做的事。第二条,把“幂等”当成系统的第一特性来设计,而不是后期补救措施,前期多一个唯一索引,后期少无数个不眠夜。第三条,监控告警要跟随调度器同步上线,不要等出故障了再补,没有监控的调度器等同于盲飞。
5.3 一件值得尝试的小事
最后分享一个我在 AX 里坚持的小设计:每个调度决策都会打印一行日志,包含任务 ID、决策原因、本次优先级、等待时长。初期会有很多人觉得日志太多,但线上定位问题时,这行日志就是唯一的线索。现在我自己排查问题时会先按时间线把所有任务的决策日志拉出来,通过对齐时间戳就能还原出当时的竞争全貌。调度器这种系统,逻辑不复杂,难的是出事之后能不能快速定位问题,而能把问题定位清楚,就已经赢了一半。