搞了一年多的“ax调度”,总算有底气拿出来给大家说说。有人一听“调度”两个字就发怵,觉得离自己很远,其实说白了就一句话:把该做的事,按正确的时间和顺序,安排好、跑起来、只执行一次。AX 这个名字是我早年顺手起的缩写,A 代表 atomic,X 代表未知数——调度系统最擅长的,就是把那些不确定的、容易出错的事,变成可预期的流程。这篇文章不打算搞大而全的教科书式讲解,就是把我从零设计、开发、上生产、填坑整个链条里的关键决策、代码细节、踩雷记录全部摊开。后台任务编排、定时触发、分布式队列、重试幂等这一类问题,看完基本都能直接落地。
1. 项目整体设计与思路拆解
1.1 我们为什么需要一套正经的调度平台
没接触过调度系统的时候,大多数人会先想到 crontab 或者 Windows 计划任务。早期我也这么干,某个整点数据同步脚本挂在某台服务器上,某个凌晨跑批的统计任务又放在另一台机器。表面上万事大吉,直到你遇到下面这些场景:
- 某台机器宕机了,crontab 跟着没了,任务当天空跑。
- 脚本执行中抛异常,没有报错通知,睡一觉醒来发现数据缺了一天。
- 两个任务之间有先后依赖,前一个没跑完,后一个就拿着脏数据开始算。
- 业务量涨上去之后,一台机器扛不住,任务要拆成多片并行,crontab 完全无能为力。
这些痛点叠在一起,催生了我现在要讲的 ax 调度。ax 调度的定位不是简单替代 crontab,它是一个完整任务生命周期管理平台:任务注册、触发策略、依赖编排、执行器分配、超时重试、结果回传、告警通知。它解决的核心问题是两个:做不做——即该不该触发;怎么做——即谁来做、做完怎么办。
如果你只是个人服务器跑几个脚本,crontab 完全够用。但一旦你的任务数量超过几十个,涉及多部门多个服务的协同,就必须把调度这个环节独立出来,否则排查问题会变成大型事故现场。
1.2 ax 调度整体架构与模块划分
我设计 ax 调度时参考了市面上成熟的分布式调度方案,但没有照搬。整体分成四块:调度中心、执行节点、注册中心、管理端。
调度中心负责接收任务定义,解析触发规则,在里面维护了计时器和队列,决定什么时候把任务推进执行流程。执行节点是真正跑业务代码的 Worker,它们启动后往注册中心注册自己的 IP 和可用状态。调度中心发现一个任务该触发了,会从注册中心选一个合适的执行节点下发指令。注册中心除了负责节点发现,还承担了任务实例状态的存储。管理端就是一个前后端面板,用来提交任务、查看执行历史、手动重跑。
这里有个关键设计原则:调度逻辑和执行逻辑严格分离。调度中心不碰任何业务代码,它只下发任务编号和参数;执行节点不关心什么时候触发,只负责把活干完并上报结果。这样做的最大好处是两边可以独立扩缩容,业务方接入时只要写好处理函数,完全不用理解调度中心内部的时钟和队列逻辑。
模块划分上,我特意把“调度中心”拆成了两个进程:一个管时间触发的 Scheduler,一个管任务实例下发和状态追踪的 Broker。Scheduler 只负责产生任务实例,写入数据库;Broker 负责把实例推给 Worker。拆开的原因很实际:时间密集计算和 IO 密集通信放在一个进程里,压力上来之后互相拖累,尤其在一次推送成百上千个任务的时候。
1.3 功能边界:ax 调度不管什么,反而更重要
分布式调度系统最容易犯的毛病是什么都往里塞。我一开始也想把文件分发、慢 SQL 分析、日志采集全塞进来,后来全部砍掉。ax 调度的定位非常明确:只管任务实例的触发、路由、生命周期和状态一致性。任务真正的业务逻辑、脚本内容、数据来源,都是业务方自己的事。
这听起来像废话,但边界不清导致的问题非常典型。比如有人想把一个大数据同步任务里依赖的外部系统健康检查也放到调度系统来做,结果调度系统只知道外部系统挂没挂,却拿不到业务侧的上下文,最后做出错误的调度决策。再比如,有人希望调度系统能帮他对动态生成的一批数据做重跑,这就不是调度该干的事,该由业务侧自行管理数据版本。把调度边界收窄之后,系统复杂度直线下降,故障定位也快得多——出现问题要么是触达问题,要么是执行问题,没有第三个暧昧地带。
2. 核心细节解析与实操要点
2.1 时间触发机制:从每秒轮询到时间轮算法
任务调度最基础的能力就是定时触发。最初版本我用的是一个很朴素的方案:后台起一个循环,每秒扫一次数据库里的任务表,找出所有满足触发时间的任务。
这个做法在小规模下完全没问题,任务量到几千条之后问题暴露了——每次扫库全表扫描,数据库压力大,而且触发时间只能精确到秒,毫秒级任务想都别想。后来我用内存时间轮替代了扫库方案。
时间轮算法其实不复杂。你可以把它理解成一个钟表盘,表盘上的每个刻度代表一个时间槽,每个槽下面挂一个任务链表。指针每走一格,就把这一格上挂着的任务都取出来执行。我用的是一层秒级时间轮加上一层环形队列的组合:秒级时间轮负责生成“到点”信号,环形队列用来做延迟调度和重试。
下面是简化后的代码:
import time import threading from collections import defaultdict from typing import Callable, Optional class TimeWheel: def __init__(self, tick: float, wheel_size: int): self.tick = tick self.wheel_size = wheel_size self.slots = [defaultdict(list) for _ in range(wheel_size)] self.current = 0 self._lock = threading.Lock() self._thread = threading.Thread(target=self._run, daemon=True) self._thread.start() def add(self, task_id: str, delay: float, action: Callable[["TaskInstance"], None], instance: "TaskInstance"): if delay < 0: raise ValueError(f"task {task_id} delay must be non-negative") ticks = int(delay / self.tick) if ticks >= self.wheel_size: # 超过刻度范围,按整圈数丢弃并提示外部做持久化调度 raise OverflowError(f"task {task_id} delay too large for this wheel") idx = (self.current + ticks) % self.wheel_size with self._lock: self.slots[idx][task_id].append((instance, action)) def _run(self): while True: time.sleep(self.tick) with self._lock: cur_slot = self.slots[self.current] self.slots[self.current] = defaultdict(list) if cur_slot: for task_id, callbacks in cur_slot.items(): for task_instance, callback in callbacks: try: callback(task_instance) except Exception as e: # 生产环境这里要上报告警 print(f"task {task_id} callback failed: {e}") self.current = (self.current + 1) % self.wheel_size实际生产里我不会让时间轮直接跑业务回调,而是让“到点”事件丢进消息队列,由 Broker 去分发。这层解耦很重要:时间轮是纯内存的高吞吐组件,如果下游 Worker 处理不过来,消息队列可以自然削峰。如果你不需要高吞吐,使用简单的轮询扫库也不是不行,但最好扫的是 Redis 里的延迟队列而不是关系库,否则并发一上来数据库就锁死了。
2.2 任务依赖编排:用 DAG 保证不会拿脏数据
定时触达解决了“什么时候干”,但业务上更头疼的是“谁先谁后”。比如统计任务要等数据同步任务完成才能开始,数据同步又要等上游接口取数完成。这些关系用固定顺序去穷举,任务一多就会疯掉。
ax 调度里任务依赖是用 DAG(有向无环图)建模的。每个任务是一个节点,前面的依赖是入边,后面的触发是出边。调度中心每完成一个任务实例,就把它的所有后继节点的入度减一;发现某个节点入度归零,就把它送入待触发队列。
听起来简单,但有两个隐藏的深坑。第一个就是依赖任务失败之后怎么办:默认策略是后继任务直接取消,原因很简单——避免用错误数据跑下游任务。第二个是多个实例之间的依赖关系:上周的任务实例依赖的应该是上周的同步任务,而不是本周的。如果你不加日期维度,直接把任务 ID 确定为依赖对象,每个实例就会互相乱串。我在任务依赖表里强制带了biz_date(业务日期)字段,依赖匹配必须是“同任务类型 + 同业务日期 + 实例状态成功”。
version: v1 job: name: order_stat_daily cron: "0 2 * * *" deps: - name: sync_order_daily biz_date: current check: success - name: sync_pay_daily biz_date: current check: success exec: type: shell command: "python3 run_stat.py --date {{biz_date}}" retry: times: 2 interval: 60 timeout: 120上面是一个任务定义文件的示例。我第一次设计时把依赖信息写死在代码里,后来发现极其痛苦——业务侧想临时加一个依赖,要等版本发布,结果每次调度策略调整都是全量发版。后面才意识到,任务定义必须配置化,并且允许在管理端高权限修改。
2.3 执行器分配与路由策略
一个任务该由哪个 Worker 执行,直接决定了延迟和稳定性。最简单的是随机选择一个存活节点,但这会造成某些机器负载高、某些机器闲着。ax 调度用的是一种带权重的最小负载策略:
每个 Worker 启动时上报自己的 CPU 核心数、当前任务实例数、最近五分钟平均耗时。调度中心给每个 Worker 算一个综合得分:
score = (running_tasks / cores) * 0.7 + (avg_duration_ms / 1000) * 0.3得分越低的节点,被选中的概率越大。这种策略虽不是最优解,但在绝大多数场景下已经能保证任务不会扎堆。
还有一类任务比较特殊:比如某任务的执行需要依赖特定机器上的本地文件或数据库驱动,这种任务必须路由到指定分组。ax 调度支持任务定义tags字段,Worker 启动时声明自己的 tags,调度中心会先按 tags 过滤,再在过滤后的集合里做负载评分。没有 tags 的任务默认在所有 Worker 里选择。
需要留意的是,路由信息不能只看注册中心的最新状态,还要防止推送瞬间 Worker 宕机导致的重复投递。我的经验是在下发热路径里增加确认机制,Worker 收到任务后必须回ACK,超时未 ACK 就自动派发到下一个可用节点。ACK 不是“执行完成”,只是“收到并准备执行”,这能有效区分“任务丢了”和“任务真失败了”。
2.4 超时、重试与幂等性设计
超时和重试是调度系统最容易翻车的环节。超时时间设置过短,慢任务被误杀;设置过长,任务堆积会拖垮 Worker。我一般建议超时时间设为业务预估耗时的 1.5 到 2 倍,但这个值必须由业务方在任务定义里显式声明,而不是调度中心默认。
重试策略也比想象中复杂。第一次碰这个问题的想法是失败了就立即重试,很快发现网络抖动的任务会连续重试三次仍失败,然后直接把消息队列冲爆。正确的做法是给重试加退避时间,第一次失败后等 30 秒,第二次等 60 秒,第三次等 120 秒。重试次数不建议设超过三次,超过之后说明问题不是临时的,而是任务本身有问题。
但真正难的还是幂等。任务重试意味着同一个业务可能被执行两遍,如果在代码层面没有做好幂等,任何重试都是灾难。我给业务方定了三条硬性约束:
- 写操作必须带全局唯一的请求 ID,业务侧要能根据 ID 去重。
- 更新操作尽量用条件更新,例如
UPDATE table SET status='done' WHERE id=xxx AND status='pending',更新行数为 0 就说明已经被执行过。 - 涉及回扣、积分发放等资金敏感操作,必须先写流水,后更新余额,以流水表的唯一索引保证只成功一次。
这些不是调度系统能做进去的功能,但却是调度可靠性的灵魂。在 ax 调度接入文档里,我把幂等检查放在第一页,因为不管调度系统多完善,执行侧的重复处理如果不过关,整体结果依然是错的。
3. 实操过程与核心环节实现
3.1 部署拓扑与基础环境准备
ax 调度的最小部署单元是:一台调度中心(同时跑 Scheduler 和 Broker)、至少两个 Worker、一个 MySQL、一个消息队列。生产环境建议把注册中心单独部署,调度中心也至少做成双活。
我第一次部署时偷懒所有组件塞在同一台 4C8G 的机器上,结果调度中心正常跑,Worker 一旦任务多了,数据库连接先把连接池占满,调度中心连不上库,整个系统雪崩。后面强制做了资源隔离,调度中心和数据存储放在同一内网,Worker 单独放在应用所在的机房,跨机房调度只走消息队列,不直接访问中心数据库。
基础环境列表:
| 组件 | 版本/规格 | 用途 |
|---|---|---|
| MySQL | 8.0,独立实例 | 存储任务定义、实例状态、依赖关系 |
| Redis | 6.x | 时间轮信号缓冲、分布式锁、幂等去重 |
| RabbitMQ | 3.9+ | 任务下发与 ACK 消息传递 |
| 调度中心 | 2C4G × 2 | Scheduler + Broker |
| Worker | 4C8G,按业务量扩容 | 执行业务任务 |
| Nginx | 1.2x | 管理层代理,不做任务转发 |
消息队列不是必须的,小规模直接用 Redis List 也可以。但是一旦涉及慢消费、死信、批量推送,消息队列的成熟语义会省下大量自研代码。我甚至建议过公司的小项目直接用 Redis Stream,后来发现重复消费的问题还是要自己处理,才统一迁到 RabbitMQ。
3.2 任务定义与接入流程
接入 ax 调度前,业务方要做的第一件事是注册执行器。执行器是一个 HTTP 接口,接收调度中心 POST 过来的任务实例 JSON。Worker 启动时会根据执行器清单自动生成这些接口的路由。
{ "taskId": "task_20240101020000", "taskName": "order_stat_daily", "bizDate": "2024-01-01", "execType": "shell", "params": { "source": "order_table", "date": "2024-01-01" }, "traceId": "a1b2c3d4-e5f6-7890-abcd-ef1234567890" }业务方拿到这个 JSON 后,执行自己的逻辑,最后通过 Worker SDK 调用report(traceId, success, message)上报结果。我见过很多人在这里图方便,直接在业务代码里改了数据库状态,却没有上报。这会导致调度中心永远无法知道任务完成情况,后续依赖任务永远不触发。所以接入文档里写得很清楚:任务是否成功,以调度中心最终收到的上报为准,业务内部状态不能替代上报。
在 Worker SDK 内部,上报之前会先检查 traceId 是否已上报过,重复上报会被直接忽略。这一点非常重要,因为偶发的消息队列重复投递会导致同一个任务的 traceId 被 Worker 收到多次。
3.3 调度中心核心循环实现
调度中心最关键的代码是 Scheduler 的主循环。它做的事只有三件:第一,从时间轮拿到期信号;第二,解析任务依赖,找出已满足条件的任务实例;第三,生成待下发消息。
def scheduler_loop(self): while True: task_instances = self.time_wheel.pop_due_tasks() for instance in task_instances: if not self.dependency_satisfied(instance): self.pending_dag_registry.register(instance) continue self.broker.publish( "task.dispatch", { "task_id": instance.task_id, "trace_id": instance.trace_id, "params": instance.params } ) time.sleep(0.1)这段代码看起来简单,但dependency_satisfied这一步背后有好几层玄机。它要查 MySQL 里该任务依赖的上游实例状态,这一步如果每次都实时查询数据库,压力很大。因此我做了一个状态缓存,把最近一小时完成任务实例的状态缓存在 Redis 里,依赖判断优先走缓存,缓存未命中再查库。
这里的坑在于缓存穿透和缓存一致性问题。某个任务实例刚在 MySQL 里更新成成功状态,还没写入 Redis,这时候下游任务来查依赖会得到“未满足”的结论,晚一步又应该触发却没触发。解决方案有两个,一是依赖判断执行后,把结果再次确认;二是把所有任务状态变更做成异步订阅,MySQL 里状态变更是主,Redis 缓存只是加速,出现不一致时以数据库兜底。实际线上我依赖了数据库查库加 Redis 热路径的双层设计,既保证热路径快速判断,也保证了最终一致。
3.4 Worker 端执行与结果上报
Worker 端更像一个通用容器,它自己不做业务,只负责把业务代码跑起来。一个核心设计是任务隔离:每个任务实例运行在独立的 goroutine/线程里,设置独立的上下文超时。任务内禁止自己再开全局线程池,避免两个任务互相影响。
def execute_task(self, task_instance): executor = self.executor_registry.get(task_instance["execType"]) ctx = TaskContext( trace_id=task_instance["traceId"], timeout_seconds=task_instance.get("timeout", 60), params=task_instance.get("params", {}), ) with Timeout(ctx.timeout_seconds): result = executor.run(ctx) self.reporter.report( trace_id=ctx.trace_id, success=True, message=result.message or "ok", )上面的Timeout我用的是信号机制实现的进程级超时。这个方案在单任务并发数上有限制,所以生产里我改成父子进程模型:Worker 主进程收到任务以后 fork 一个子进程来跑业务代码,主进程在超时时间后检查子进程是否结束,未结束就强制 kill 并上报超时失败。这样即使业务代码里出现死循环或者无响应的网络请求,也不会拖死整个 Worker。
结果上报我采用了“先写本地执行日志,再通过网络上报”的方式。本地有失败的,Worker 有一个独立的重试线程,每隔一段时间把未上报的执行日志重新上报。这个设计帮我解决过一次大故障:RabbitMQ 集群短暂不可用,如果上报全部失败,任务明明成功但调度中心全都显示等待,后续任务全部卡死。本地兜底日志让 RabbitMQ 恢复后两分钟内状态全部回补完毕。
3.5 管理端功能要点
管理端不需要做得很花哨,但有一个功能必须有:人工重跑。数据任务注定会有各种不可控因素,上游数据晚了、代码有 bug、外部接口临时改动,都可能导致任务失败然后重跑。
人工重跑有几种模式:
- 按单个任务实例重跑,带上原先的所有参数。
- 按业务日期重跑一整串 DAG,比如把 2024-01-01 当天所有依赖链路上的任务全部重新执行。
- 只重跑失败任务,不触碰已成功的下游。
这里最容易被忽略的是“已成功下游是否应该被重跑”。如果上游某张表数据因为 bug 被修正过,下游统计结果已经写错且是否可覆盖,那么重跑上游时就应该连带下游一起重跑。我在管理端提供了一个开关force_downstream,默认关闭,避免误触下游任务导致重复扣用户量等业务异常。但数据修正场景必须手动打开并确认。
重跑操作本身也要走权限审批。高权限账号才能执行强制重跑,普通运维只能查看日志。这个限制一开始被吐槽麻烦,后来一次偶然有人在管理端误点了全量重跑,跑完发现把月初月报全部覆盖成错误数据,还是靠备份恢复的,打那以后权限审批就成了铁律。
4. 常见问题与排查技巧实录
4.1 任务重复执行:报警和去重如何双管齐下
ax 调度上线后第一个大事故就是任务重复执行。某个订单同步任务在凌晨跑了两次,第一次执行到一半 Worker 进程被 OOM 杀掉,调度中心判断任务失败,开始重试;但第一次执行的进程实际已经把一部分数据写入目标库,第二次执行再次写入,导致目标表里出现了重复记录。
这次事故让我明白,纯靠调度中心记录的“执行状态”来判断是否失败,根本不可靠。真正的防线是业务侧的幂等设计。排查技巧上,我积累了一套顺序:
先查执行历史表,看同一个 traceId 出现了几次分发记录。如果多次分发,看每次分发的 Worker 节点和时间间隔,判断是 Worker 崩溃后的重试还是消息队列重复投递。再看业务日志里的 traceId 去重日志,确认业务代码是否做了重复抑制。
如果确认是调度系统层面的重复,基本原因有两个:一是 Worker 执行完成后上报结果超时,但业务实际已经完成,调度中心按未上报处理并且重试;二是依赖判断时缓存状态不一致导致的下游提前触发。第一种问题的解法是业务侧在真正执行前先调claim(traceId)接口获取执行权,谁拿到分布式锁谁执行,谁执行完谁上报,未拿到锁的任务直接跳过。
4.2 任务卡死不结束:罪魁祸首是超时设置不合理
很多初用 ax 调度的人会奇怪:任务明明卡死了,调度中心却不把它踢掉,后面一大堆依赖任务全都排队等它。答案通常是超时时间没设置。任务定义里如果没有显式给 timeout,我默认只给 60 秒。但对一些跑全量数据回溯的任务来说,60 秒远远不够,于是会被误杀;反过来,有人把超时设成两个小时,结果某个任务内存泄漏卡住,下游整整等了两个小时才发现。
排查卡死任务有个技巧,别只看 Worker 状态,先看任务实例的执行日志。如果日志一直停留在“开始执行”,但进程还存在,大概率是任务在等一个永远不会来的结果,例如网络请求连接池耗尽、数据库锁等待、死循环。这时候从调度中心强制终止任务只解决了表面问题,真正要做的是进入 Worker 看线程 dump。在 Worker 端我会事前开启 JMX 或者 Python 的faulthandler,卡死的时候直接发送SIGABRT拿到线程堆栈,十次里有八次能立刻定位到阻塞点。
4.3 时间不准导致触发混乱:单机时钟漂移问题
ax 调度早期只在单机部署时没发现问题,后面加了多个调度中心节点做双活,诡异的事来了:同一个任务有两个触发实例,时间相差了十几秒。
排查后才意识到是服务器系统时间漂移了。标准做法是给所有调度中心节点配置 NTP 同步,但这还不够,因为应用层无法感知时钟是否跳变。我在时间轮上加了一个节流保护:触发任务时记录当前系统时间,任务实例生成时必须带上expected_time,Worker 收到后校验当前系统时间和 expected_time 的差值,超过五秒就拒绝执行并上报时钟异常。
这个设计帮我们在一次多云环境割接时避免了大范围误触发。那次迁移后部分新节点 NTP 配置没生效,系统时间慢了十分钟,如果 Worker 不校验,所有任务都会按错误时间提前或延后执行,数据统计会全乱。
4.4 任务积压:队列监听的延迟与扩容
某个大促期间,消息队列里的任务量一下子暴涨,消费者线程处理不过来,任务下发到执行之间的延迟从正常几百毫秒变成几分钟。业务方反馈订单统计迟迟不出来,但调度中心看任务状态全是“运行中”。
这个问题的排查路径比较清晰:看队列堆积数、看 Worker 消费速率、看单任务平均耗时。但我们的定位方法更提前一步,在 RabbitMQ 每个队列上挂了延迟监控指标,当队列里的消息数量超过消费者数量的 100 倍时,自动触发告警。扩容时不是盲目加 Worker,而是要跟任务类型匹配,比如 CPU 密集任务需要加核心数,IO 密集任务加线程数就够了。
另外任务积压不一定是性能不够,也可能是死循环把 Worker 全占住。遇到大批量任务积压,我建议先把队列暂停消费五分钟,观察 Worker 进程的 CPU 是否降到低位。如果降低,说明业务侧有死循环或非法占用;如果依然跑满,才是真正的扩容需求。
4.5 任务实例状态不更新:数据库连接池与事务边界
有些任务执行成功了,但调度中心的任务实例状态一直停留在“执行中”或“待下发”。这个现象很常见,原因往往不是上报失败,而是状态更新的事务边界不对。
典型的错误是业务代码在自己的数据库事务里调用了上报接口。假设业务事务还没提交,上报接口已经把成功状态发给了调度中心,调度中心随即触发下游任务;下游任务去查业务数据时,这个事务还没提交,查到的还是旧数据,整个链条就错了。而且如果业务事务回滚,状态已经被上报为成功,那问题就更严重了。
排查思路是检查任务实例状态变化时间与业务写入时间是否一致。如果状态成功时间早于业务数据落库时间,大概率就是事务顺序问题。ax 调度接入规范里明确要求:上报动作必须放在业务事务提交之后。如果需要在事务中先发状态,必须改成“事务提交后的异步上报”,或者把上报接口设计成可以被补偿和撤销。
4.6 快速排查工具与常用命令
这里整理一份我平时排查 ax 调度问题的命令清单,都是最直接有效的动作:
| 症状 | 第一步动作 |
|---|---|
| 任务没触发 | 查调度中心日志里时间轮任务是否到期,再查pending_dag_registry里依赖是否阻塞 |
| 任务触发但 Worker 没执行 | 查队列消费者连接数,确认 Worker 是否已注册,rabbitmqctl list_consumers |
| Worker 执行但不上报 | 查业务日志中 traceId,确认是否进入上报代码;再查本地执行日志缓存 |
| 状态上报成功但下游没触发 | 查依赖缓存 Redis 键,确认状态是否成功写入 |
| 时间轮指针不动 | 查系统时钟和 NTP 状态,chronyc tracking |
| 任务分发到多个节点 | 查 Worker ACK 超时配置和网络延迟,重点排查跨机房链路 |
这些命令未必能一步定位,但基本能迅速把问题范围缩小到一两个模块,省去到处抓日志的迷茫感。
5. 生产环境高可用与容量规划思路
5.1 调度中心多活方案与脑裂规避
ax 调度要做多活,最先要考虑的是调度中心多节点之间如何避免同时触发同一任务。因为它们都挂着时间轮,如果都把同一个定时任务发出去了,任务就重复执行了。
我用的是 Redis 分布式锁加数据库唯一约束的双保险。Redis 锁是快速路径,拿到锁的节点才允许从时间轮取出某个任务;同时数据库中任务实例表的task_name + expected_time建了唯一索引,就算两个节点同时拿到锁,后写入的也会因为唯一索引冲突而失败。
这套方案的关键点在于锁释放时机和唯一索引冲突后的处理。一个节点拿到锁之后如果还没处理完就宕机,锁会因过期时间自动释放,另一个节点接管后会发现数据库里已经存在相同实例,把它标记为“孤儿实例”并告警。我不会自动清理孤儿实例,因为无法确定原节点到底执行了多少,交给人工确认更安全。
5.2 容量评估的经验公式
做了这么久调度系统,我总结了一套相对粗犷但实用的容量估算公式。假设你有 N 个任务,平均执行时长是 T 秒,每小时任务执行总量是 R,那么 Worker 并发线程数至少是 R 乘以单任务平均耗时的积。
举个例子:每小时有 1000 个任务实例,平均每个跑 5 秒,那么这 1000 个任务在同一小时内的总耗时为 5000 秒。一个 Worker 有 8 个线程,一小时能提供 28800 秒的执行能力,显然一个 Worker 就够;但如果平均耗时变成 30 秒,总耗时 30000 秒,8 线程的 Worker 还是会超时,需要两个 Worker 或者上 16 线程。
调度中心侧的容量相对好算,因为它的职责轻,核心关注点是按时产生的任务实例数量。时间轮的每个刻度处理任务的时间只要在刻度间隔内完成即可。假如刻度是 1 秒,单刻度内最多 100 个任务,每任务处理耗时 2ms,那么 1 秒的处理时间完全够。你需要预留三倍以上余量,防止突发任务量和网络抖动阻塞主循环。
5.3 插件化扩展:如何接入新任务类型
ax 调度默认支持 shell、http、python、sql 四类执行器。但业务系统五花八门,很快有人提出要接入 Spark 任务、Flink 任务、甚至自定义的 Java 类。一开始我打算把这些都内置进去,后来发现这会让 Worker 变得无限重。
插件化是更优雅的方案。我定义了执行器接口,只要实现四个方法就能注册新的任务类型:init、prepare、run、describe。第三方业务方只要打成 JAR 或 Python 包放到 Worker 指定目录,重启后自动被加载,完全不需要改动调度中心代码。
接口大概长这样:
class Executor(ABC): @abstractmethod def init(self, config: dict): ... @abstractmethod def prepare(self, ctx: TaskContext): ... @abstractmethod def run(self, ctx: TaskContext) -> ExecResult: ... @abstractmethod def describe(self) -> str: ...插件化带来一个附加好处:不同团队可以维护自己的插件包,职责清晰,版本独立。坏处是插件之间可能依赖冲突,比如两个插件都引用了同一个库的不同版本。对这个问题的妥协方案是为每个插件定义独立的虚拟环境,Worker 启动时按插件名切换到对应环境。代价是启动时间变长了一些,但隔离效果值得。
5.4 高可用演练与降级策略
技术方案做得再漂亮,不演练到生产环境里还是会翻车。我给自己定了一个原则:每个月做一次故障演练,专门破坏基础设施。
演练的典型场景包括:停掉一个调度中心节点、停掉一个 Worker、停掉整个 RabbitMQ、让 MySQL 一段时间不可用。每次演练都有对应的可接受影响时间,例如 worker 单节点宕机要求两分钟内自动恢复任务执行,RabbitMQ 不可用要求十分钟内积压消息不丢失且恢复后自动回放。
演练过程中最容易暴露的问题就是依赖外部的隐蔽调用。比如调度中心判断依赖的时候,顺手调了注册中心的健康检查接口,检查接口又依赖数据库,数据库挂了就一路雪崩。这类问题在静态代码里很难发现,只有断掉某个组件才能暴露。排查总结的降级顺序是:优先保证调度中心自身存活,次要保证消息不丢,最后才考虑业务任务及时执行。调度中心如果自身都存活不了,后面的一切都是空谈。
5.5 监控体系建设
ax 调度配套的监控指标,核心是四类:
第一类是延迟指标,包括任务下发到 Worker 收到的时间间隔、任务执行完成耗时。第二类是错误指标,包括执行失败数、重试次数、超时次数。第三类是堆积指标,包括队列积压量、时间轮待触发任务数、依赖未满足任务数。第四类是资源指标,包括各节点的 CPU、内存、磁盘。
单纯列出指标没用,关键是阈值和联动。队列积压大了,不只是告警,要能触发自动扩容,或者自动暂停某些低优先级任务。我实现的优先级策略里,高优任务有独立的队列,低优任务不可占用高优任务的消费线程。这样一个大报表任务卡住时,不会影响订单实时状态同步任务。
监控面板上的字段不用多,五个足够:今日任务总数、失败总数、平均触发延迟、最大执行耗时、失败 Top5 任务。这个面板是我日常看的最勤的一个,因为它能快速回答“系统今天到底健不健康”这个根本问题。
6. 常见误区与后续演进
6.1 误区一:调度系统万能,什么任务都往里塞
最典型的一句话是“我们那个接口慢,用调度的重试机制做补偿吧”。调度系统的重试机制解决的是执行失败后的临时故障,如果接口本身设计有问题,无论重试多少次都会失败,只是把问题延后。把调度当成万能工具箱,是项目后期维护成本飙升的根源。
6.2 误区二:任务状态越多越精细
初始设计里我定义了十几个状态:初始化、排队中、分发中、收到 ACK、执行中、暂停中、重试中、执行失败、执行成功、已取消、已跳过、超时未确认、超时已取消。真实生产里,运维和业务方根本记不住这么多状态,最后判断任务是否正常只看一个问题:它到底成没成功。
后面我把状态收敛成五个核心状态:运行中、成功、失败、等待重试、已取消。其他所有细分状态都是这五个的补充标签。运维看状态一眼就知道该做什么—失败了看失败原因,等待重试就等,成功就放心。过度设计的状态机只会让系统更难用。
6.3 后续可扩展的方向
ax 调度目前能支撑公司每天上百万次任务触发,但我知道它仍然有大量问题需要继续演进。比如多集群任务的血缘关系追踪,目前只做了实例级依赖,还没有做数据血缘。再比如更精细的任务策略:根据业务优先级动态缩容或扩容执行线程,而不是简单地按队列优先级调度。还有任务的结果回传目前只支持成功失败标志和消息文本,后续会支持用户自定义结果结构,方便下游业务直接解析。
不过在所有演进规划里,我认为最该做的还是把接入体验做得更顺滑。代码生成器、任务模板、一键本地调试,这些开发体验层面的东西往往比新奇的技术方案更能提高效率。
6.4 最后的一点私人经验
如果你打算自己动手做一个调度系统,我的最大建议是先把失败路径设计完整再写触发逻辑。初学者往往把精力放在如何准时触发任务上,真正到了生产环境才发现,失败恢复、重复处理、状态一致性才是每天打交道最多的部分。先想清楚 Worker 宕机怎么办、消息丢了怎么办、任务跑一半进程崩溃怎么办,再考虑怎么快速触发一万个任务,你的调度系统才真正立得住。
踩过这么多次坑之后我的体会是,调度系统其实没有多高深,它的本质就是对抗不确定性。时间不确定、节点不确定、依赖不确定、结果不确定,而调度要做的是在这么多不确定中提供一个相对确定的执行框架。把这句话想透了,你设计的每个细节都会不一样。