资讯动态

自研轻量级调度内核ax:从时间轮到分布式锁的实战拆解

发布时间:2026/9/25 14:31:41 来源:尧图企业网站定制
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 这类独立调度平台是另一个极端功能确实全但需要部署独立服务端、维护管控台对已经跑着的微服务来说等于又增加了一个必须保证高可用的中心节点。当时团队更缺的是“一个能随服务一起启动、把定时任务纳入统一生命周期管理的内嵌组件”而不是一个重量级调度平台。我们想要的核心能力就五条支持 cron、固定周期、延迟触发这三种基本触发方式。任务执行要有超时控制不能一个慢接口拖死线程池。分布式部署下任务不能重复执行同一个任务同一时刻只能有一个实例在跑。失败重试要有退避策略不能失败后马上用最大频率重新打爆下游。接入成本足够低业务方只需要注册任务函数其他收尾逻辑尽量内部消化。这五条定下来就已经足够说明自己造轮子的理由了。并不是说那些大平台不好而是我们的场景和团队的维护成本、控制力需求不匹配。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( namepayment.check_timeout, triggercron, spec0 */5 * * * ?, timeout30, retry3, retry_delay5, ) def check_payment_timeout(): 检查支付超时订单 ...这个register装饰器背后做的工作并不像表面这么简单它至少完成四件事解析触发配置初始化一个Trigger对象计算并记录下一次执行时间。把任务元数据注册进全局任务表并写入对应的时间槽位。把任务状态持久化到数据库或 Redis供其他实例做一致性判断。注册一个统一熔断器这个任务后续的每次执行都会经过超时和重试策略的过滤。执行链路则是时间轮指针推进 - 取出到点任务 - 对任务加分布式锁 - 放进线程池执行 - 根据执行结果更新状态或安排重试。这里有个细节值得展开加锁和真正执行之间必须有一段“安全缓冲”否则如果锁在任务还没跑完时就过期另一个实例又拿到同一把锁任务就会重复执行。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_ms500, wheel_size3600): 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 取 500mswheel_size 取 3600那么一圈就是 30 分钟。超过 30 分钟的延迟任务咋办简单方案是升级为分层时间轮但更实际的方案是对于周期类和 cron 类任务记录下一次执行时间之后拆分成“到最近一个整点执行”在调度器启动时就把这些任务切短对于长延迟任务则改用另一个“持久化延迟队列”来兜底时间轮只处理短延迟任务。这其实也是一个经验教训不要试图用一套机制解决所有触发类型。想让时间轮、数据库扫表、分布式延迟队列各干各擅长的活才是合理的架构。接下来是任务执行器。这里我用了线程池因为这是绝大多数业务系统最容易接受的模型import concurrent.futures class Executor: def __init__(self, max_workers10): self.pool concurrent.futures.ThreadPoolExecutor(max_workersmax_workers) def submit(self, fn, timeoutNone): future self.pool.submit(fn) try: result future.result(timeouttimeout) return (success, result) except concurrent.futures.TimeoutError: future.cancel() return (timeout, None) except Exception as exc: return (error, exc)这里有几个细节future.result(timeout...)确实能触发超时异常但线程本身不一定终止所以需要在业务函数里检查中断标记配合实现软超时。线程池的max_workers不建议设成固定值我一般建议按任务的预估耗时做加权短任务占 0.5普通任务占 1长任务占 2。否则一个跑 10 分钟的任务就能把池子占满。执行器最好支持队列长度预警当积压任务超过某个阈值时主动告警而不是等线程池里的任务全部堆积后再被监控发现。3.3 接业务系统的两种姿势ax 接入业务系统基本有两种姿势。第一种是嵌入式业务服务引入 ax 依赖服务启动时自动注册任务调度器随服务进程一起跑。这种方式的优点是部署简单不需要额外维护调度中心缺点是任务分散在各个服务实例里统一管理要靠每个服务的日志和数据库记录配合。第二种是独立调度节点单独部署一个 ax 进程业务服务通过消息队列或 HTTP 回调来领取任务。所有任务元数据集中在调度节点业务方只暴露一个统一的执行入口。这种方式适合任务归属不明确、跨服务调用的场景但多了一个部署单元对可用性要求更高。我们最后走的是混合路线普通定时任务用嵌入式跨服务汇总类的任务由独立调度节点下发。如果你是从零开始我建议先做嵌入式它能让你更直观地理解调度器的生命周期等踩顺了再考虑拆分节点。调度器这类基础设施最怕一上来就搞过度设计复杂度会吃掉你排查问题的精力。4. 实战踩坑记录与问题速查4.1 漏执行为什么总是发生在凌晨凌晨是定时任务最密集的时刻也是系统最容易出事的时刻。第一次大规模漏执行发生在某次灰度发布后我们只重启了一半实例。调度器在启动时会把任务重新装载进内存时间轮但重启的实例还没抢到锁而没被重启的旧实例正因为锁被部分释放而误判任务已被执行。结果就是第二天一早数据对不上。排查到最后发现问题有两层。第一层是启动顺序任务装载必须等分布式锁的基础组件 ready 之后再做启动时的任务回放不能抢在 Redis 连接可用之前。第二层是任务执行标记重启用例里任务“已执行”的标记不能只存在内存里必须同步到数据库。修复方案是任务执行前先更新数据库里的last_heartbeat执行完再更新last_success_time重启后通过比较这两个时间戳判断是否需要回放。这里也提醒大家启动回放时最容易出的错是把“回放”做成“立刻把所有到点任务并发重跑”。那样轻则下游被瞬时流量打爆重则产生大量重复数据。正确做法是回放时按最终一致性的原则错峰执行给回放任务加上小幅度随机启动延迟。4.2 分布式锁过期引发的重复执行Redis 分布式锁的经典问题我们在线上遇到过不止一次。有一回某个任务负责同步会员积分到第三方系统执行时间平均 3 秒我图省事把锁过期时间设成了 5 秒。结果那一次下游系统响应特别慢任务跑了 8 秒锁在 5 秒时过期了另一个实例立刻拿到同一把锁把同一批积分同步了两遍。业务方收到的积分消息直接翻倍用户那边显示的数字一度对不上。这个事的教训有三个锁过期时间绝不能按平均耗时设置要按最大耗时的两倍设置甚至更长。强烈建议开启锁的自动续期机制。在我们自研的 lock 模块里持锁线程会启动一个后台协程每隔三分之一过期时间就给锁续期任务结束再释放。只要任务还在跑锁就不会过期。所有执行逻辑必须做好幂等。分布式里的“绝不重复”是伪命题永远假设自己会被重复执行然后让下游忽略重复请求才是正确姿势。4.3 线程池耗尽导致的任务排队另一个高频问题是线程池被长任务占满导致后面所有短任务集体延迟。我们的一个数据导出功能因为接口调用比较慢最坏能跑 20 分钟。而 tasks 表里每两分钟就会生成一个新的导出任务。默认线程池只有 10 个 worker后来高峰期 9 个 worker 都在跑导出任务另一个 worker 去跑普通定时任务结果普通任务的执行时间被拖到十几分钟业务方反馈“定时任务是不是挂了”。排查思路是给执行器打点记录每个 worker 正在跑什么任务。我们发现大量任务长期处于 RUNNING且线程名统一是ThreadPoolExecutor-0基本就能判断是线程池内部排队。解决方式是引入任务分级队列。短任务和长任务使用不同的线程池短任务线程池的 core 数可以稍大长任务线程池则单独控制并发上限同时给长任务加上“等待时间告警”。任务分级做起来很简单本质就是在任务配置里加一个queue_label字段注册时就决定进哪条队列。4.4 问题速查表这里整理一张问题速查表方便大家在自己排查时对照症状可能原因排查方向快速缓解任务漏执行调度进程重启后任务未回放检查启动日志中的任务装载数手动触达一次并补好回放逻辑任务重复执行分布式锁过期或未加锁查看 Redis 锁租约与执行时长开启锁自动续期设置更长的过期时间任务执行超时线程池被占满看线程池存活线程与运行任务任务分级队列调整超时配置重试后仍然失败下游接口幂等性问题检查业务日志中同一请求重复到达在下游加去重字段大量任务积压时间轮槽位设计不合理检查时间轮单圈容量拆分层时间轮或切换持久化队列任务执行顺序错乱同一触发点多个任务并发检查任务注册时是否指定顺序约束对强依赖任务做执行前置检查这张表最大的价值不在于答案本身而在于提醒排查顺序。我们经常犯的错是看到“重复执行”就去改任务函数看到“漏执行”就去改 cron 配置而忽略了调度器本身的状态。建议你先把调度器日志、锁状态、线程池状态这三层看全了再动业务代码。5. 一些真实的体感把 ax 这套东西从零搭起来并在线上稳定跑了大半年之后我的体会是调度系统真正的复杂度不在“定时触发”本身而在于你如何让“时间的约定”和“系统的分布式现实”和平共处。你无法在同一时空里保证任务一定只执行一次你也无法保证服务器永远不崩溃你能做的就是让失败可以被发现、可被恢复、且在恢复时不会造成二次伤害。另一个重要心得是用户侧的体验往往和任务执行链路直接相关很多看似随机的小问题追根究底就是某个定时任务没跑对。把任务的状态记录、执行链路、重试策略做成一个透明的整体比堆各种高深算法更解决问题。我个人目前还在往这个方向扩展比如给 ax 加更精细的依赖编排让一个任务跑完再按条件触发下一个任务也会定期复盘线上误操作的案例让任务的“手动触发”带更强的约束。现在再遇到凌晨告警至少第一反应不是手足无措而是看调度状态去判断是不是又踩了老坑。

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑