☰
自研轻量级调度内核ax:从时间轮到分布式锁的实战拆解
2026/9/25 14:31:36 网站建设 项目流程

ax 这个代号,在我这儿其实是 Action eXecution 的缩写,翻译成大白话就是“动作执行”。从去年开始,我一直维护着这套轻量级调度组件。起因特别朴素:团队从单体脚本转向微服务之后,散落在各个服务里的定时任务变成了一堆没人敢碰的黑盒。凌晨的告警、写死逻辑的重跑脚本、动不动就把任务队列堵死的超时调用,这些问题逼着我把调度这件事从头梳理了一遍。最后沉淀下来的这套东西,就是 ax。现在大家在聊的“ax调度”,大多数场景下指的就是这种面向中小团队、能快速接入业务系统的轻量级调度方案。

这篇文章我不会去推任何商业化产品,只拆解一个自研调度内核在设计、落地和排查问题时踩过的坑。适合谁看?如果你也在为 cron 脚本失控、分布式任务乱跑、重试逻辑一团糟而头疼,那这篇值得你花几分钟读完。没有太高门槛,我尽量把原理和实操揉在一起写,尽量做到每一步你都能照着试。

1. ax从哪来:一次凌晨三点被叫醒之后

1.1 当时的一地鸡毛

最早那阵子,我们服务的定时任务主要靠三样东西:Linux 的 crontab、Java 里的 Quartz,和一堆不知道自己该在哪台机器上跑的 shell 脚本。表面上看大家相安无事,实际上是没人愿意捅这个马蜂窝。

真实情况是这样:订单模块每天凌晨要跑一个结算脚本,脚本里串联了十几个内部接口调用。负责维护的人离职后,这脚本基本处于“黑盒状态”。某天凌晨接口报错,脚本直接中断,但因为是 cron 在跑,没有告警推送,第二天早上十点用户才发现数据不对。然后就是经典的追责、翻日志、手动补数据三件套。

后来我们把脚本改成了给任务中心发消息的模式,本质上还是“定时触发——执行完就忘”。任务重试靠调脚本里的 for 循环,任务状态靠人工盯数据库,任务超时靠猜。这种情况持续到一次凌晨三点的订单异常爆发,我终于决定做一件事:把调度能力集中收敛,做成一个可以内嵌到各服务里的统一调度内核。这就是 ax 的起点。

1.2 方案对比:为什么不是 Quartz / xxl-job

在动手之前,按惯例把市面上的方案过了一遍。Quartz 很成熟,但有几个问题不适合我们:一是对 Quartz 的深度定制需要花费不少时间,二是它默认的持久化和集群模式配置起来略重,调度逻辑和业务代码容易纠缠在一起。xxl-job 这类独立调度平台是另一个极端,功能确实全,但需要部署独立服务端、维护管控台,对已经跑着的微服务来说等于又增加了一个必须保证高可用的中心节点。

当时团队更缺的是“一个能随服务一起启动、把定时任务纳入统一生命周期管理的内嵌组件”,而不是一个重量级调度平台。我们想要的核心能力就五条:

  1. 支持 cron、固定周期、延迟触发这三种基本触发方式。
  2. 任务执行要有超时控制,不能一个慢接口拖死线程池。
  3. 分布式部署下任务不能重复执行,同一个任务同一时刻只能有一个实例在跑。
  4. 失败重试要有退避策略,不能失败后马上用最大频率重新打爆下游。
  5. 接入成本足够低,业务方只需要注册任务函数,其他收尾逻辑尽量内部消化。

这五条定下来,就已经足够说明自己造轮子的理由了。并不是说那些大平台不好,而是我们的场景和团队的维护成本、控制力需求不匹配。

1.3 我理解的“ax调度”到底是什么

搭完一套能用的调度系统之后,再回头看“ax调度”这个词,我的理解会更偏架构一些。所谓调度,本质上是把“时间触发的动作”和“业务系统里要执行的逻辑”解耦,让任务具备可观察、可控制、可恢复这三个属性。

  • 可观察:任务是什么时候触发的、执行状态如何、失败原因是什么,都要有记录。
  • 可控制:任务可以暂停、取消、手动重跑,而不是只能干瞪眼。
  • 可恢复:进程崩溃、断电、网络抖动之后,任务不能凭空丢失,要能从某个标记点恢复。

ax 这个名字后来也变成了我们内部的一个泛指:它既指调度内核本身,也指围绕调度器建立的一套任务治理规范。调度器不是万能的,它替你把时间逻辑管理起来,但业务侧必须配合,做幂等、做超时、做状态标记。两者合在一起才能真正解决线上乱七八糟的执行问题。

2. 调度内核的核心设计拆解

2.1 触发引擎:用时间轮代替裸 cron

实现一个调度器,最常见的起点是照搬cron的定期扫表做法:每分钟检查一次任务表,看哪些任务到了执行时间。这么做简单,但问题很明显,精度粗,秒级任务根本做不了;而且每轮要扫全表,任务量大了以后浪费很严重。

ax 的触发引擎用的是一种类似时间轮的机制。时间轮可以理解成一个循环的、带有槽位的数组,每个槽位代表一个时间刻度。新任务注册时,根据它的下一次执行时间计算该落到哪个槽里。有个指针按周期推进,取落到当前槽的任务,检查是否真的到点,是就丢给执行器。

时间轮方案相比扫表的关键优势在于,任务的检查范围从全表缩小到了当前槽位,复杂度从 O(n) 降到了 O(1) 量级。我们用的刻度是 500ms,对于绝大多数定时任务场景精度已经够用。如果你需要毫秒级甚至微秒级的触法,那要考虑的就不只是调度器了,而是整个事件处理链路的延迟,那又是另一套设计。

容错上还得考虑进程重启。时间轮里的任务都是内存态,重启就丢了。所以在注册任务时我们会同时把任务元数据写到数据库,在启动时做一次“回放”:找出那些到时间但没执行的任务,重新装载进时间轮。这个回放逻辑为了简化,用的是任务表里的next_run_time字段,扫描范围比全表小很多。

2.2 任务注册与执行链路

ax 的任务模型很简单,一个任务就是一个函数加一组配置。在 Python 版本里大致长这样:

@task.register( name="payment.check_timeout", trigger="cron", spec="0 */5 * * * ?", timeout=30, retry=3, retry_delay=5, ) def check_payment_timeout(): """检查支付超时订单""" ...

这个register装饰器背后做的工作并不像表面这么简单,它至少完成四件事:

  1. 解析触发配置,初始化一个Trigger对象,计算并记录下一次执行时间。
  2. 把任务元数据注册进全局任务表,并写入对应的时间槽位。
  3. 把任务状态持久化到数据库或 Redis,供其他实例做一致性判断。
  4. 注册一个统一熔断器,这个任务后续的每次执行都会经过超时和重试策略的过滤。

执行链路则是:时间轮指针推进 -> 取出到点任务 -> 对任务加分布式锁 -> 放进线程池执行 -> 根据执行结果更新状态或安排重试。这里有个细节值得展开:加锁和真正执行之间必须有一段“安全缓冲”,否则如果锁在任务还没跑完时就过期,另一个实例又拿到同一把锁,任务就会重复执行。

2.3 三个关键策略:去重、超时、重试

网上聊任务调度的文章不少,但真到线上扛流量时,决定生死的往往就是这几个策略是否设置得合理。ax 里这三个参数全部是可配置的,而且必须被配置,不允许用默认值蒙混过关。

去重。我们用的是 Redis 分布式锁,锁的 key 就是任务 name,值为本次执行的唯一 id,过期时间默认是任务超时时间的两倍。为什么不是相等,就是为了防止任务因为某些原因没在预估时间内结束,导致锁先于执行释放。锁过期时间设得太短会重复执行,设得太长又会在任务异常退出时导致后续执行被阻塞。我们实际操作中会把任务按执行时长分档,短平快的锁过期时间设为 60s,长任务会单独评估,而不是一刀切。

超时。超时控制靠的是执行线程池里每个 Worker 持有的 Future。执行器提交任务后,通过future.result(timeout=...)强制等待,到点没返回就取消并标记失败。注意cancel只能中断未开始执行的线程,对已经跑起来的线程没有强制杀死的效果,所以超时之后还要把当前线程的标记位设成 interrupted,业务代码配合检查这个标记才能做到真正的中断。

重试。这里有个用血泪换来的结论:重试不能只看次数,还要看退避策略。我们的默认策略是“指数退避 + 抖动”。第一次失败后等 5 秒,第二次等 25 秒,第三次等 125 秒,再叠加 0 到 2 秒的随机抖动,避免大量任务同时在整点失败后的同时涌向重试通道。

def next_retry_delay(attempt: int) -> int: base = min(5 * (2 ** attempt), 300) return base + random.randint(0, 2)

这套策略配合业务方做的幂等设计,基本能做到“不丢任务、不反复轰炸下游”。重点提醒一下:重试的前提是下游接口必须幂等,否则重试就是灾难的加速器。

3. 从零跑通 ax:实操全流程

3.1 目录长什么样

很多人在写调度组件时会犯一个错,把调度逻辑和业务任务全部塞在一个包里,导致后续想单独测试调度器都费劲。ax 的目录结构我尽量保持了边界清晰:

ax/ ├── core/ │ ├── schedule.py # 调度器主类,时间轮推进入口 │ ├── trigger.py # cron / interval / delay 三种触发器 │ ├── task.py # 任务模型与装饰器 │ └── lock.py # Redis 分布式锁封装 ├── registry/ │ └── task_registry.py # 全局任务注册表 ├── worker/ │ └── executor.py # 线程池执行器,超时与取消逻辑 ├── store/ │ └── task_store.py # 任务元数据持久化 └── api/ └── manager.py # 对外管理接口:暂停、取消、手动触发

这个结构拆分是基于一个问题:调度器、注册表、执行器这三者的生命周期是不同步的。调度器要一直转,执行器会根据负载动态调整,注册表则会被业务代码在启动时就填满。如果混在一个模块里,后续加功能非常痛苦,比如你想新加一个触发类型,就得在调度器主文件里到处改。

3.2 核心代码逐段落地

我们先把调度器主类拉出来看。为了看清核心,我略掉锁和存储的细节,只保留主干逻辑:

import time from collections import defaultdict from typing import Callable class TimeWheel: def __init__(self, tick_ms=500, wheel_size=3600): self.tick_ms = tick_ms / 1000.0 self.wheel_size = wheel_size self.slots = defaultdict(list) self.current = 0 def add(self, delay: float, callback: Callable): ticks = int(delay // self.tick_ms) if ticks >= self.wheel_size: raise ValueError("delay too long, need layering") slot = (self.current + ticks) % self.wheel_size self.slots[slot].append(callback) def advance(self): self.current = (self.current + 1) % self.wheel_size for cb in self.slots.pop(self.current, []): cb()

这个实现只支持单圈时间轮,即任务延迟必须小于tick_ms * wheel_size。tick_ms 取 500ms,wheel_size 取 3600,那么一圈就是 30 分钟。超过 30 分钟的延迟任务咋办?简单方案是升级为分层时间轮,但更实际的方案是:对于周期类和 cron 类任务,记录下一次执行时间之后,拆分成“到最近一个整点执行”,在调度器启动时就把这些任务切短;对于长延迟任务,则改用另一个“持久化延迟队列”来兜底,时间轮只处理短延迟任务。

这其实也是一个经验教训:不要试图用一套机制解决所有触发类型。想让时间轮、数据库扫表、分布式延迟队列各干各擅长的活,才是合理的架构。

接下来是任务执行器。这里我用了线程池,因为这是绝大多数业务系统最容易接受的模型:

import concurrent.futures class Executor: def __init__(self, max_workers=10): self.pool = concurrent.futures.ThreadPoolExecutor(max_workers=max_workers) def submit(self, fn, timeout=None): future = self.pool.submit(fn) try: result = future.result(timeout=timeout) return ("success", result) except concurrent.futures.TimeoutError: future.cancel() return ("timeout", None) except Exception as exc: return ("error", exc)

这里有几个细节:

  1. future.result(timeout=...)确实能触发超时异常,但线程本身不一定终止,所以需要在业务函数里检查中断标记,配合实现软超时。
  2. 线程池的max_workers不建议设成固定值,我一般建议按任务的预估耗时做加权:短任务占 0.5,普通任务占 1,长任务占 2。否则一个跑 10 分钟的任务就能把池子占满。
  3. 执行器最好支持队列长度预警,当积压任务超过某个阈值时主动告警,而不是等线程池里的任务全部堆积后再被监控发现。

3.3 接业务系统的两种姿势

ax 接入业务系统基本有两种姿势。第一种是嵌入式:业务服务引入 ax 依赖,服务启动时自动注册任务,调度器随服务进程一起跑。这种方式的优点是部署简单,不需要额外维护调度中心;缺点是任务分散在各个服务实例里,统一管理要靠每个服务的日志和数据库记录配合。

第二种是独立调度节点:单独部署一个 ax 进程,业务服务通过消息队列或 HTTP 回调来领取任务。所有任务元数据集中在调度节点,业务方只暴露一个统一的执行入口。这种方式适合任务归属不明确、跨服务调用的场景,但多了一个部署单元,对可用性要求更高。

我们最后走的是混合路线:普通定时任务用嵌入式,跨服务汇总类的任务由独立调度节点下发。如果你是从零开始,我建议先做嵌入式,它能让你更直观地理解调度器的生命周期,等踩顺了再考虑拆分节点。调度器这类基础设施最怕一上来就搞过度设计,复杂度会吃掉你排查问题的精力。

4. 实战踩坑记录与问题速查

4.1 漏执行为什么总是发生在凌晨

凌晨是定时任务最密集的时刻,也是系统最容易出事的时刻。第一次大规模漏执行发生在某次灰度发布后,我们只重启了一半实例。调度器在启动时会把任务重新装载进内存时间轮,但重启的实例还没抢到锁,而没被重启的旧实例正因为锁被部分释放而误判任务已被执行。结果就是第二天一早数据对不上。

排查到最后发现,问题有两层。第一层是启动顺序:任务装载必须等分布式锁的基础组件 ready 之后再做,启动时的任务回放不能抢在 Redis 连接可用之前。第二层是任务执行标记:重启用例里,任务“已执行”的标记不能只存在内存里,必须同步到数据库。修复方案是任务执行前先更新数据库里的last_heartbeat,执行完再更新last_success_time,重启后通过比较这两个时间戳判断是否需要回放。

这里也提醒大家,启动回放时最容易出的错是把“回放”做成“立刻把所有到点任务并发重跑”。那样轻则下游被瞬时流量打爆,重则产生大量重复数据。正确做法是回放时按最终一致性的原则错峰执行,给回放任务加上小幅度随机启动延迟。

4.2 分布式锁过期引发的重复执行

Redis 分布式锁的经典问题我们在线上遇到过不止一次。有一回,某个任务负责同步会员积分到第三方系统,执行时间平均 3 秒,我图省事把锁过期时间设成了 5 秒。结果那一次下游系统响应特别慢,任务跑了 8 秒,锁在 5 秒时过期了,另一个实例立刻拿到同一把锁,把同一批积分同步了两遍。业务方收到的积分消息直接翻倍,用户那边显示的数字一度对不上。

这个事的教训有三个:

  1. 锁过期时间绝不能按平均耗时设置,要按最大耗时的两倍设置,甚至更长。
  2. 强烈建议开启锁的自动续期机制。在我们自研的 lock 模块里,持锁线程会启动一个后台协程,每隔三分之一过期时间就给锁续期,任务结束再释放。只要任务还在跑,锁就不会过期。
  3. 所有执行逻辑必须做好幂等。分布式里的“绝不重复”是伪命题,永远假设自己会被重复执行,然后让下游忽略重复请求,才是正确姿势。

4.3 线程池耗尽导致的任务排队

另一个高频问题是线程池被长任务占满,导致后面所有短任务集体延迟。我们的一个数据导出功能,因为接口调用比较慢,最坏能跑 20 分钟。而 tasks 表里每两分钟就会生成一个新的导出任务。默认线程池只有 10 个 worker,后来高峰期 9 个 worker 都在跑导出任务,另一个 worker 去跑普通定时任务,结果普通任务的执行时间被拖到十几分钟,业务方反馈“定时任务是不是挂了”。

排查思路是给执行器打点,记录每个 worker 正在跑什么任务。我们发现大量任务长期处于 RUNNING,且线程名统一是ThreadPoolExecutor-0,基本就能判断是线程池内部排队。

解决方式是引入任务分级队列。短任务和长任务使用不同的线程池,短任务线程池的 core 数可以稍大,长任务线程池则单独控制并发上限,同时给长任务加上“等待时间告警”。任务分级做起来很简单,本质就是在任务配置里加一个queue_label字段,注册时就决定进哪条队列。

4.4 问题速查表

这里整理一张问题速查表,方便大家在自己排查时对照:

症状可能原因排查方向快速缓解
任务漏执行调度进程重启后任务未回放检查启动日志中的任务装载数手动触达一次,并补好回放逻辑
任务重复执行分布式锁过期或未加锁查看 Redis 锁租约与执行时长开启锁自动续期,设置更长的过期时间
任务执行超时线程池被占满看线程池存活线程与运行任务任务分级队列,调整超时配置
重试后仍然失败下游接口幂等性问题检查业务日志中同一请求重复到达在下游加去重字段
大量任务积压时间轮槽位设计不合理检查时间轮单圈容量拆分层时间轮或切换持久化队列
任务执行顺序错乱同一触发点多个任务并发检查任务注册时是否指定顺序约束对强依赖任务做执行前置检查

这张表最大的价值不在于答案本身,而在于提醒排查顺序。我们经常犯的错是看到“重复执行”就去改任务函数,看到“漏执行”就去改 cron 配置,而忽略了调度器本身的状态。建议你先把调度器日志、锁状态、线程池状态这三层看全了,再动业务代码。

5. 一些真实的体感

把 ax 这套东西从零搭起来,并在线上稳定跑了大半年之后,我的体会是:调度系统真正的复杂度不在“定时触发”本身,而在于你如何让“时间的约定”和“系统的分布式现实”和平共处。你无法在同一时空里保证任务一定只执行一次,你也无法保证服务器永远不崩溃,你能做的就是让失败可以被发现、可被恢复、且在恢复时不会造成二次伤害。

另一个重要心得是用户侧的体验往往和任务执行链路直接相关,很多看似随机的小问题,追根究底就是某个定时任务没跑对。把任务的状态记录、执行链路、重试策略做成一个透明的整体,比堆各种高深算法更解决问题。我个人目前还在往这个方向扩展:比如给 ax 加更精细的依赖编排,让一个任务跑完再按条件触发下一个任务;也会定期复盘线上误操作的案例,让任务的“手动触发”带更强的约束。现在再遇到凌晨告警,至少第一反应不是手足无措,而是看调度状态去判断是不是又踩了老坑。

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

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

立即咨询