ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

从零自研分布式定时任务调度器:ax的设计与踩坑复盘

从零自研分布式定时任务调度器:ax的设计与踩坑复盘 ax 是随手敲的两个字母也是我这个调度器项目的代号。最早它只是代码仓库名后来同事在需求单里写了“接入 ax 调度”这个名字就这么叫开了。调度这个领域听起来不就是“定时跑任务”嘛可真要自己在生产环境里做一套才会发现坑全在细节里任务状态不可见、重复执行、节点一挂就没人管、重试把下游打爆。这篇就把 ax 从零到落地过程中那些关键设计和踩过的坑一次性倒出来给想自研或者正在选型调度系统的朋友做个参照。1. ax 出现的背景现成调度框架治不了的三个毛病1.1 第一个毛病Quartz 在微服务环境下越用越别扭说实话我们最开始不是没试过现成方案。项目里最早用的是 Quartz单机跑没什么问题但服务一拆成多实例就麻烦了。Quartz 集群模式靠数据库锁来保证同一个 Job 只有一个节点执行随着任务量涨上来锁竞争和数据库压力越来越明显。而且 Quartz 的任务信息存在 RAMStore 里的默认配置会导致重启丢任务换成 JDBCStore 又要多维护十几张表。最难受的是它的执行历史很弱任务跑失败了、跑了多久、重试了几次这些信息基本要靠自己另写日志去凑。我并不是否定 Quartz它在传统单体应用里非常成熟可靠。只是到了服务拆细、实例随意扩缩容的阶段我们需要的是一个“调度中心”而不是“嵌入在应用里的调度库”。1.2 第二个毛病XXL-Job 和 DolphinScheduler“功能过剩”我也认真评估过 XXL-Job 和 DolphinScheduler。XXL-Job 功能确实全有可视化后台、有告警、有分片广播社区也活跃。但引入它意味着要额外部署一个调度中心还要维护它的数据库、权限体系、执行器 SDK。对我们当时只有三四个后端、任务总量几千个的团队来说属于“为了喝奶养头牛”。DolphinScheduler 更偏工作流编排DAG 血缘、定时调度、依赖关系这些能力很强但我们 90% 的需求只是“到点跑一个接口或者一条命令”根本不需要那么重的工作流引擎。还有一个很现实的问题公司内网环境对第三方系统依赖管控得很严引入一个新中间件要走审批、要申请机器、要定升级机制。这种情况下自己写一个满足核心需求的轻量调度器反而比推一个重框架更顺。1.3 第三个毛病业务方要的是“可观测、可干预”业务侧真正不满意的点是任务跑没跑、为什么没跑、失败了怎么补救这些问题没人能快速回答。我们原来的定时任务散落在各个服务里有的用 crontab有的用代码里的 Timer有的写在 CI 流水线里。出问题的时候排查路径完全不同运维同学想重跑一个任务都不知道去哪点。所以 ax 立项时我给自己定了三条原则调度与执行彻底分离调度器只负责“到点了告诉你该跑了”执行由各业务方自己的执行器负责。任务是一条可追溯的记录每一次触发都有唯一 ID状态、耗时、日志、重试次数全部落表。默认允许手动干预随时可以手动触发、取消、标记成功或失败。这三条后来被证明是 ax 能活下来的根本。调度器最大的价值不是“定时”本身而是让“定时”这件事变得透明。2. 先把时间算明白ax 的调度核心一个调度器最先要解决的就是“下一次执行时间”怎么算。这部分看起来是纯函数问题实际做起来全是边界。2.1 触发规则我为什么最终留了两种表达方式ax 的触发规则支持两种写法。第一种是标准的 cron 表达式适合“每天凌晨两点”“每周一上午十点”这种固定节奏。第二种是自定义的简单 DSL比如every 5m表示每五分钟一次after 2025-06-01T00:00:00Z表示某时间点之后开始执行主要在测试和临时任务时用。cron 解析的难点在于计算下一次触发时间。网上现成的解析库很多但直接拿来用的前提是你要搞清楚它返回的是本地时间还是 UTC是否处理了夏令时秒位是否支持。ax 里我规定所有时间统一按 UTC 存储展示层再转本地时区。这个决定让我少踩了很多坑——后文会专门说。核心伪代码大概是这样的def next_time(expr, base): # base 是当前时间返回下一个合法的执行时间 parts expr.split() # 分别解析 分、时、日、月、周 for candidate in generate_candidates(base, parts): if candidate base: return candidate这个函数本身不复杂但要注意日和周的关系。cron 里日和周同时有值时不同实现处理方式不一样有的取并集有的取交集。我在 ax 里明确按“两者都满足才算合法”来实现并且在文档里写清楚免得使用者困惑。2.2 用分级时间轮保存“到点任务”cron 表达式算出的是一个个绝对时间点调度器要做的就是在这些时间点到达时触发任务。最简单的办法是维护一个优先队列每次取堆顶元素看时间到没到。任务量小的时候完全够用但任务量涨到几万、几十万时堆的插入和删除都是 O(logN)性能会变得难看。ax 选择了时间轮方案。基本原理是维护一个环形数组每个槽位放一个任务链表一个指针每秒向前移动一格移到哪个槽就把那个槽里的任务全部取出来执行。插入任务时根据延迟时间算出应该放哪个槽时间复杂度是 O(1)。class SimpleTimeWheel: def __init__(self, slots60): self.slots [[] for _ in range(slots)] self.current 0 def add(self, task, delay_seconds): idx (self.current delay_seconds) % len(self.slots) self.slots[idx].append((task, time.time() delay_seconds))但单层时间轮有个问题精度越高槽位越多。如果精度要 1 秒最多只能覆盖 60 秒的任务超过 60 秒的任务就没地方放了。ax 的实现是三层时间轮秒级轮覆盖 1 到 60 秒分钟级轮覆盖 1 到 60 分钟小时级轮覆盖更长的延迟。指针每转完一圈就把下一级轮的任务降级到当前级。实际实现里每个槽位存的不是任务对象本身而是任务的到期绝对时间戳。这样即使调度进程因为 GC 停顿迟了几毫秒也能把到期时间落在当前 tick 范围内的任务捞回来执行而不是傻乎乎地只处理“当前指针指向的槽”。2.3 时钟漂移与补偿执行时间轮方案有个天然问题如果指针只处理当前槽位一旦 tick 发生漂移比如进程卡了 5 秒那这 5 秒内到期的任务就会被漏掉。ax 的解决办法是每一轮 tick 都记录当前真实时间now处理槽位时把所有scheduled_time now的任务都拿出来而不是只取恰好等于当前 tick 的任务。对于已经过期超过一定阈值的任务比如 5 分钟默认不执行直接标记为 missed并触发告警。这样避免了“追任务”追出一堆并发。选择 5 分钟这个值是因为我们业务里大多数任务延迟几分钟跑也没关系如果是秒级精确的场景这个阈值就要调小。这里我想特别提醒时间轮的 tick 不能依赖sleep(1.0)。sleep 并不保证精确应该每次循环记录真实时间算出下一次 tick 应该等待多少毫秒。ax 第一个版本就是 sleep 出来的结果线上任务偶尔莫名其妙延迟几秒排查了半天才发现是系统时钟本身就是跳跃的。3. 从单机到集群ax 的分布式三件套单机版 ax 跑通之后我很快就遇到了新的问题一个调度节点挂了怎么办两个调度节点同时跑会不会重复触发业务执行器多了之后任务怎么均匀分发这些都是分布式调度绕不开的三座山。3.1 选主先用 MySQL 锁后来才换 Redis第一版 ax 只允许一个调度节点通过 MySQL 的SELECT ... FOR UPDATE抢一把“主节点锁”。这个方案简单粗暴但有几个问题锁的续期机制要自己写节点优雅退出时锁释放不及时数据库抖动会导致误判主节点挂掉。后来我换成了 Redis 的SET key value NX EX来实现选主。每个节点启动时尝试写入同一个 key写成功的成为主节点并启动一个后台 goroutine 每隔几秒续期。如果主节点挂了锁自然过期其他节点再次竞争。方案优点缺点MySQL FOR UPDATE实现简单事务内可靠数据库压力大故障恢复慢Redis NX EX快续期简单需要保证 Redis 可用性etcd lease语义清晰分布式协调正统引入额外组件运维成本高ax 现在的结论是任务量在几万级别以下Redis 选主完全够用如果公司已经有 etcd 集群用 etcd 更省心千万别为了选主专门引入一套新系统。3.2 幂等宁可重复执行不可丢失执行分布式调度最怕的不是“任务没执行”而是“任务执行了两次”。前者可以通过告警发现后者往往造成数据错误且极难追踪。ax 默认采用 at-least-once 语义调度器保证任务至少触发一次执行器侧自己去重。具体做法是调度器每次触发时生成一个execution_id它由task_id trigger_time attempt组成幂等表的主键就是这个execution_id。执行器接到触发请求后先把execution_id插入本地执行记录表如果插入冲突说明这条执行已经处理过直接忽略。INSERT INTO execution_log(exec_id, task_id, trigger_time, attempt, status) VALUES (:exec_id, :task_id, :trigger_time, :attempt, running) ON CONFLICT(exec_id) DO NOTHING如果DO NOTHING那就说明是重复请求直接返回成功不执行业务逻辑。这里有个容易犯的错不要把幂等键设计成task_id trigger_time因为同一次触发可能因为重试被调度器再次发出。必须把attempt加进去保证“同一次触发的不同尝试”也能被正确区分而不是误判成重复执行。3.3 任务分片让所有调度节点都有活干选主解决了“谁调度”的问题但只有主节点干活的话从节点就浪费了。ax 的做法是任务分片。每个任务根据task_id做哈希映射到不同的分片每个调度节点负责一部分分片。这样所有节点都能参与调度单个节点挂了它负责的分片会被重新分配。分片算法的细节值得注意。直接取模hash(task_id) % shard_count在分片数变化时会引发大面积迁移。ax 采用的是带虚拟槽的哈希一致性把整个哈希空间分成 1024 个槽每个节点负责一段连续区间。任务哈希后先定位到槽再找到负责该槽的节点。节点增减时只需要迁移部分槽而不是全量重算。3.4 一个节点挂掉之后的完整时间线把上面三件事串起来一个节点挂掉的完整流程大概是这样的t0 秒主节点最后一次上报心跳t5 秒其他节点发现心跳超时发起选主竞争t10 秒新主节点选出t15 秒新主节点重新分配原节点负责的分片t20 秒新主节点扫描执行记录把“已触发但未完成”且允许重试的任务重新加入调度队列。这个流程里最关键的是第 20 秒。如果不对未完成任务做补偿那故障期间本该执行的任务就彻底丢了。我见过不少调度系统选主做得很好但忘了“补偿执行”这一步结果每次主节点切换都会吞掉一批任务业务方过了半天才反馈说数据没更新。4. 放大镜看执行细节超时、重试、限流调度器把任务“送出去”只是第一步真正难的是把执行过程中的异常处理干净。这一章说三个高频问题超时怎么切、重试怎么退、流量怎么挡。4.1 超时必须在执行侧切不能指望调度器很多调度系统把超时做成“到了时间就标记失败”但这是自欺欺人。如果执行器里那段代码还在跑你只是把状态改成失败线程照样泄漏、资源照样占用。ax 的做法是调度器把超时时间放在触发请求里执行器必须用 context 或者等价机制来真正中止任务。在 Go 里就是context.WithTimeout在 Python 里建议用子进程而不是线程来跑耗时的任务因为线程超时后没法真正杀掉只有进程才能被 kill。import subprocess proc subprocess.Popen(cmd, shellFalse) try: out, err proc.communicate(timeout60) except subprocess.TimeoutExpired: proc.kill() out, err proc.communicate() raise TimeoutError(task killed after 60s)调度器的超时判断只是保险丝真正的熔断开关必须装在执行器这一端。4.2 重试不是简单地重跑一遍重试策略 ax 用的是指数退避加抖动next_delay base_delay * (2 ** attempt) next_delay min(next_delay, max_delay) next_delay random.uniform(0, jitter_seconds)指数退避的道理大家都懂但两个参数容易被忽略。一是max_delay必须设上限否则任务失败二十次之后重试间隔可能变成几天业务上根本接受不了。二是抖动必须有不加抖动的指数退避会造成“重试风暴”一大批任务同时失败同时按同样的退避曲线重试到点之后又同时打垮下游。我们的线上配置一般是base_delay 2smax_delay 300sjitter 1s最多重试 3 次。重试次数记录在 execution 表里方便后面排查。手动补偿另走接口不占用自动重试的次数。4.3 背压执行器忙不过来时必须告诉调度器有一类问题是“调度器疯狂发任务执行器接不住”。任务量突然上涨、下游接口变慢都会导致执行器线程池被打满。ax 的执行器内部有一个有界队列触发请求进来先放到队列里队列满了直接返回busy。调度器收到busy响应后会把后续分发给该执行器的任务临时降速比如原本每秒发 10 个降到每秒 2 个并且把busy次数记录到监控里。这样做的目的是让执行器有时间消化积压而不是一边积压一边继续灌新的。这个设计虽然简单但效果极好。线上出过几次事故都是因为某条链路变慢执行器线程池打满任务全部超时触发新一轮重试把下游彻底压垮。有了背压信号之后顶多是任务延迟执行不会出现雪崩。4.4 调度器必须盯着的六个指标ax 上线后我整理了一套必看的监控指标这里直接列出来指标含义报警建议schedule_latency_seconds实际触发时间与计划触发时间之差超过 30 秒就要关注execute_duration_seconds任务执行耗时分布按 P95 观察不设固定阈值execution_success_rate执行成功率低于 90% 报警queue_backlog执行器待处理队列深度持续上涨报警time_wheel_gap时间轮单次 tick 处理耗时超过 1 秒说明任务积压严重retry_total重试次数突增时尽快排查如果这些指标都是齐的调度器的大部分问题都可以在业务方感知之前先暴露出来。4.5 ax 的配置示例和核心 API一个任务的配置长这样{ name: order-sync, trigger: { cron: 0 */10 * * * ? }, executor: http://executor-a:8080/run, timeout_seconds: 60, retry: { max_attempts: 3, base_delay_seconds: 2, jitter_seconds: 1 } }核心 API 只有三个几乎全公司后端都能一小时内上手POST /api/v1/tasks创建或更新任务POST /api/v1/tasks/{id}/trigger手动触发一次GET /api/v1/executions?task_idxx查询执行历史。我不做复杂的可视化后台因为把几十个任务和几千条执行记录用表格列清楚比做一个漂亮的大屏有用得多。5. ax 踩坑日志与复盘如果再写一遍我会动哪里5.1 三个印象最深的事故第一个和时间存储有关。早期版本把触发时间用本地时间存进数据库测试环境没问题但部署到跨时区机房后所有定时任务都偏了几个小时。后来统一改成 UTC 存储只在展示层做转换。这个改动本身不大但定位过程折腾了一天半。第二个和重试风暴有关。某天下游依赖的 Kafka 集群故障大量任务同时失败自动重试策略又都集中在一个时间窗口里恢复。结果 Kafka 刚恢复就被我们的重试流量再次打垮服务雪崩持续了近 40 分钟。后来加了全局熔断单任务并发超过阈值直接跳过重试并发了告警。另外把重试的抖动范围加大保证不同任务的重试时间不会扎堆。第三个是幂等键设计问题。第一版用task_id trigger_time做幂等键结果同一触发时间的重试请求被当成重复执行过滤掉丢了真正的重试机会。后来改成task_id trigger_time attempt才把问题解决。这个错误属于经验问题纸上谈兵很难发现只有在线上的执行记录里看到“明明重试了但状态没变化”才会意识到。5.2 现在我对 ax 适用边界的判断经常有人问我 ax 能不能用在某某场景我的回答很直接适合的场景不适合的场景定时任务几千到几万级大规模工作流 DAG 编排任务执行耗时在秒级到分钟级毫秒级精确调度需要快速定位“任务到底跑没跑”复杂任务依赖和条件分支中小团队不想维护重框架需要可视化编排、血缘追溯已有基础监控体系愿意接指标需要开箱即用的完善后台界面如果你要做的是数据平台级的调度还是老老实实用 DolphinScheduler 或 Airflow。ax 适合的是“我就想把一批定时任务管清楚”的场景。5.3 给想自研调度器的人一个最低配建议如果你也动了自研的念头我建议不要一上来就搞分布式。最简方案只需要三样东西一张任务表、一个轮询 cron 的循环、一张执行记录表。先用这个最小闭环跑几个月把业务方的需求摸清楚再加时间轮、加选主、加分片。ax 就是这么一点点长出来的一开始的版本丑得很但不影响它解决了真实问题。最后再分享一个小技巧给调度器起接口名字的时候不要让所有人都看不懂。我们内部所有接口都叫tasks和executions新同事上手基本不需要文档。ax 这个名字虽然随意但“ax 调度”在今天已经是群里大家默认会用的词了。对我个人来说这个项目最大的收获不是技术有多深而是所有踩过的坑都转化成了可复盘的规则——这些规则比任何架构设计图都值钱。
返回列表