资讯详情

自研调度内核ax:时间轮、状态机与分布式一致性解析

📅 2026/9/28 21:49:51 | 华诺云谱 👁 阅读
自研调度内核ax:时间轮、状态机与分布式一致性解析
最近在基础架构圈子里大家开始频繁提起“ax调度”这四个字。如果你还没接触过我简单交代一下背景ax是我大半年一直在维护的一个轻量级调度内核的代号取自Adaptive eXecution的缩写。市面上调度框架并不少但真把业务放进去跑一遍就会明白现成的轮子看着省事落到具体场景里往往不是过重就是不够灵活。这篇文章会把ax从设计动机、内核模型、分布式一致性到落地案例完整拆开来讲适合正在做任务调度选型、或者想自研一个调度器的后端同学参考。我自己是从业务研发转过来做基础组件的所以写的东西一般都比较“土”不绕弯子。你跟着读下来至少能知道一个调度系统真正难在哪里以及遇到问题时该从哪个切入点排查。1. 为什么放着现成的调度框架不用非要自己写 ax先说结论并不是现成框架不好而是大部分团队的调度需求并没有复杂到需要一套完整平台的程度。我们最初评估过 Quartz、XXL-Job、Temporal 这几类方案每个都有各自的隐性成本最后才决定自己写一个内嵌式的调度内核。1.1 四类现成方案的隐性成本Quartz 最大的问题是集群模式靠数据库锁竞争来协调节点调度量一上来数据库压力和分布式锁等待就开始变成瓶颈。你在单机上跑 Quartz 确实很舒服但一旦要做高可用、多节点抢跑它那一套 RAMJobStore 和 JDBCJobStore 的切换逻辑会把你折腾得够呛。XXL-Job 的定位是完整的分布式调度平台调度中心、执行器、控制台一应俱全。功能确实全但部署和运维成本也高调度中心要单独维护执行器要接入它的 SDK 并注册心跳整个链路被固定成“中心派发到执行器”的模型。对我们几十个微服务来说为了几个延迟任务和周期任务就要多维护一套中心化服务的生命周期性价比太低。Temporal 就更重了。它是完整的工作流引擎支持补偿、人审、子流程编排学习曲线非常陡。如果团队只有基础 CRUD 的背景投入产出比非常不划算。我们当时需要的不是“工作流”只是“到点了把事情跑起来”而已。1.2 我们真正想要的一颗可嵌入的调度内核帮团队做技术选型时我会先列一个“不要什么”清单不要独立的调度中心、不要强制改业务工程的启动方式、不要复杂的 SDK 依赖。于是ax的核心定位就很清晰了——它是一颗能直接塞进业务进程里的调度内核。你可以在任何 Go 服务里ax.New()创建一个调度器实例然后注册 handler、投递延迟任务所有状态持久化在你自己指定的数据库里。这带来两个直接好处一是部署模型非常简单没有多余的组件要维护二是扩展路径灵活单个实例跑不下了可以再加节点做选主派发而不用推翻重来。2. ax 的内核时间轮与一张持久化任务表的配合调度器最核心的能力就一句话任务到点了能被触发。听起来简单真实现起来会发现“高效地知道任务到点”这件事并不容易。2.1 为什么选时间轮而不是轮询 Cron 表达式最朴素的实现方式是定期扫表把所有状态为待执行且next_run_at now的任务捞出来。任务量小的时候这完全没问题但任务量上到十万级之后每次全表扫描的代价就很大。即便加索引高频扫描对数据库的压力还是让人心疼。所以内核里我用了分层时间轮。你可以把它理解成钟表的时分针一个 60 格的秒轮加一个 60 格的分轮任务往对应的格子放指针每走一秒就摘掉当前格子里的到期任务复杂度是 O(1)。拿“每分钟执行一次”来举例任务会被放入分轮的某个槽位等到分针走到那个位置时触发不需要每一秒都去扫描整个任务表。时间轮只负责高效地发现“哪些任务到点了”它不自己存业务任务详情。真正的任务主记录在数据库里时间轮里放的是任务 ID 和内存中的到期时间。这是一个非常关键的权衡——纯内存方案丢了状态纯数据库方案扛不住高频两者配合才是工程上比较合理的解法。2.2 内存时间轮 数据库任务表的混合模型具体启动流程是这样调度器启动后先恢复数据库里所有处于待执行状态且还没有跑完的任务把未来一段窗口内的任务灌入时间轮。平时的时间轮靠 leader 节点的调度循环从数据库批量拉取“未来 5 分钟内到点”的任务补充进去。这样每次数据库查询都是带索引的局部查询而不是无差别全表扫描。数据库这张ax_tasks表是状态的权威来源时间轮只是加速发现的本地缓存。即使进程重启、时间轮全部丢失重启后也能从数据库恢复。这里要特别强调一下时间精度设置。时间轮 tick 我设置的是 1 秒这足够覆盖大部分业务场景。如果你的场景需要毫秒级延迟触发可以把 tick 改成 100ms但要注意内存开销和数据库拉取频率会同步上涨。2.3 核心调度循环的简化示意调度循环用 Go 写大概就是这样一个结构func (a *Ax) scheduleLoop(ctx context.Context) { timer : time.NewTicker(a.cfg.Tick) defer timer.Stop() for { select { case -ctx.Done(): return case -timer.C: tasks : a.wheel.Expire(time.Now()) for _, t : range tasks { a.Dispatch(ctx, t) } a.refill(ctx) // 从DB补充未来窗口内的任务到时间轮 } } }Expire取出当前刻度上的任务 ID 列表Dispatch负责走状态机推进和派发。补轮的动作放在同一个循环里是为了避免多线程对时间轮内部状态产生并发竞争。实际开发中refill需要控制每次拉取的量避免一次性灌入太多任务导致内存抖动。3. 任务状态机与派发链路从“到点”到“回调成功”中间发生了什么调度器最怕的事情是“模棱两可”——任务到底派发出去没有执行器到底跑了没有跑的结果是什么如果这些问题没有明确答案后面所有重试和补偿都是无根之木。所以ax从一开始就设计了一套严格的状态机。3.1 状态机设计与迁移规则一张表把这些状态说清楚状态含义谁写入可迁移到SCHEDULED已持久化等待到时投递方 / 调度器DISPATCHINGDISPATCHING已到点尝试派发中调度器RUNNING/SCHEDULEDRUNNING执行器确认执行中执行器SUCCEEDED/FAILEDSUCCEEDED成功结束执行器无FAILED业务执行失败或重试耗尽调度器 / 执行器SCHEDULED重试 /FAILED放弃状态迁移最关键的一条铁律状态的每一次变化都必须通过条件更新 SQL 实现不能先读再写。比如从SCHEDULED迁移到DISPATCHINGSQL 条件是WHERE id ? AND status SCHEDULED更新影响行数为 1 才说明当前节点抢到了这个任务的派发权。这种乐观锁式的做法天然避免了两个节点同时派发同一个任务。3.2 一次完整派发的链路细节一个周期任务“到点”后完整的链路是这样的时间轮到期摘出任务 ID调度器先把它从SCHEDULEDCAS 成DISPATCHING然后根据任务注册时选择的执行器地址发起调用。执行器收到请求后先落一条执行记录表再把状态更新为RUNNING业务逻辑跑完后再回调调度器把结果写回SUCCEEDED或FAILED。这里有个很容易被忽略的点执行器必须先落执行记录再执行业务逻辑。如果把“执行记录”放在业务逻辑之后一旦业务逻辑崩溃或进程被杀这条任务会进入无限重试循环而且你根本查不到上一次执行到哪一步了。执行记录表里的execution_id是每次派发独立生成的 UUID后面幂等全靠它。3.3 失败重试这样设计不把自己坑死失败重试不能傻傻地立即重试否则一个下游故障能把执行器打挂。重试策略沿用了指数退避第一次失败等 5 秒第二次等 25 秒第三次等 125 秒最多重试 5 次。每次重试都只是把任务状态 CAS 回SCHEDULED并更新next_run_at让它重新进入时间轮。这里有个非常实用的经验重试次数一定要和服务商给出的“最大重试”语义区分开。max_retry指业务重试次数不包含触发阶段失败的重试。触发阶段失败比如执行器地址不通应该单独用dispatch_retry_count控制两个计数字段分开维护不然你会发现任务没过业务重试就已经被系统放弃了。4. 多实例不会重复调度吗租约、幂等与孤儿任务回收ax支持多节点部署节点之间通过选主决定谁来跑调度循环。但选主只是第一步真正的复杂度在“选主失败”“锁过期”“执行器宕机”这些边缘场景里。4.1 选主锁与租约续期我们用了 Redis 做分布式选主核心就是SET ax:leader nodeId NX PX 10000这把锁。拿到锁的节点成为 leader负责从数据库拉任务、填时间轮、派发调度其他节点进入待命状态每 5 秒尝试抢一次锁。锁不能设了就不管。Redis 锁最常见的坑是“过期时间到了但任务还没执行完”另一个节点抢到锁两个 leader 同时存在任务被重复派发。所以我专门写了一个续约协程每隔TTL/3的时间刷新一次锁的过期时间。简单说就是把租约续住而不是寄希望于业务逻辑能在锁过期前跑完。4.2 fencing token 才是防脑裂的终极手段光靠续约仍然不够踏实。极端情况下leader 节点发生长 GC 或网络分区锁过期释放了新 leader 上位但老 leader 又缓过来继续派发任务两个节点同时操作任务表。这时候唯一可靠的防线就是 fencing token。每次抢锁成功锁值里的 token 会自增。leader 在派发请求时把这个 token 放进 header 里执行器收到请求后检查 token如果发现小于自己本地记录过的最大 token直接拒绝执行。这个机制保证哪怕老 leader 还在“自嗨”它发出的一切派发请求都会被执行器挡掉双主造成重复执行的窗口被压到最小。4.3 幂等和孤儿任务回收一个都不能少即便做到上面的防护网络超时、进程崩溃这些情况还是会带来重复投递。所以每个任务派发时都会带execution_id执行器端用它对结果表做唯一索引。重复请求进来时查询到execution_id已经存在就直接返回旧结果不会重复执行业务逻辑。执行器宕机的场景更隐蔽。调用方已经发出请求执行器进程却没了任务停留在RUNNING状态永远无法终结。为此账户任务表里额外加了last_heartbeat字段执行器执行期间每 10 秒更新一次。调度器巡检线程会扫描RUNNING超过 5 分钟且心跳停止的任务把它们强制标记为FAILED并触发重试。代价是心跳的额外写入但换来的是“死任务”可以被及时发现和处理。5. 两个落地案例订单超时关闭和定时数据对账讲了这么多原理落到具体业务里才更有感觉。我挑两个已经在线上稳定跑的场景展开。5.1 订单超时未支付自动关闭这个场景是典型的延迟任务。用户下单创建的订单30 分钟内未支付就要自动关闭。传统做法是定时批量扫订单表然后挨个判断既浪费数据库资源又没办法做到精确的 30 分钟级触发。ax的延迟投递只需要一行调用ax.Delay(order.close.timeout, orderID, 30*time.Minute, orderCloseHandler)订单创建时把任务投递出去任务记录通过orderID作为业务唯一键。orderCloseHandler在处理前先查一次订单状态如果已经支付就直接返回nil不再做任何关闭动作。这个“先查再改”的习惯非常重要因为延迟任务一定存在迟到触发的情况业务方必须自己处理“状态已经变化”的场景。时间轮的精度能保证触发时间非常接近 30 分钟整点实测 p99 误差在 200ms 以内远好于过去那种每分钟扫一次表的批量方案。更重要的是数据库负载明显降下来了订单表不再需要频繁被全表扫。5.2 分钟级数据对账任务另一个场景是跨系统数据对账。每个两分钟跑一次把订单库和账务系统的数据拉出来比对不一致的记录要写告警表。对账任务执行时间可能超过 2 分钟如果任务还在跑下一次又触发了就要做并发控制。这里我用的是ax的互斥执行能力同一个task_key在同一时间点最多只允许一个实例执行。实现上也是在状态机里做文章调度前检查是否有其他执行中的同task_key任务有就直接跳过本次触发。这种“跳过”策略对于周期任务来说比“排队”更合理因为对账要看的是最新快照旧的数据等对完了再跑一次就行了。5.3 手动补偿接口不能少调度系统哪怕再稳总有需要人工介入的时刻。ax提供了一个手动触发接口允许运维人员绕过时间约束直接把某个任务推到DISPATCHING状态。它的核心用途不是日常操作而是故障恢复后的补偿——比如数据修复脚本、漏掉的报表重跑。这个接口的权限控制必须严格不能让普通业务人员随便调。否则一个误操作就可能让一个任务瞬间执行成千上万次这种事故在调度系统领域太常见了。6. 上线三个月我们踩过的三个坑自研调度器的好处是可控坏处是坑全靠自己踩。这三个问题每一个都让我们半夜爬起来过分享出来帮你避雷。6.1 任务“凭空消失”时钟漂移在背后捣鬼第一阶段测试时批跑任务一直正常结果上线后出现一个诡异现象某些任务比设定的执行时间提前了十几秒触发。排查了很久最后发现是数据库服务器的时间比应用服务器快了十几秒。任务恢复流程读取的是数据库时间而时间轮触发用的是应用本地时间两边一旦有偏差就会出现“未来时间任务被当成到期任务执行”的问题。修复方案是把时间标准统一所有调度决策一律以 leader 节点的单调时钟为准数据库时间只作为持久化恢复的参考。恢复流程里加了一个偏差修正逻辑节点启动时先从任务表中读取最新一条任务的updated_at和本地时间做差值修正后续所有到期判断都基于修正后的时间。从此这个坑再没踩过。6.2 回调接口被长任务堵死健康检查先挂了执行器接口既要处理任务回调又要响应注册中心的健康检查。任务量一大长耗时任务占满了 HTTP 连接池健康检查请求也在排队执行器节点被注册中心误判为离线流量被摘除。更麻烦的是摘除后的流量又转移到其他节点引发连锁效应。解决思路是给执行器接口分两类线程池一类专门处理健康检查永远不排队另一类处理任务派发允许一定程度的积压。同时给任务调用设置客户端超时超过指定时间直接放弃不要无限期占用连接。这是我在排查告警风暴时总结出来的调度系统的健康检查和业务处理必须物理隔离不能共享资源。6.3 批量重试触发数据库行锁死锁大促那段时间任务量激增数据库开始频繁报死锁。看日志发现全部集中在工作批次 CAS 更新任务表状态时。原因是批量更新 SQL 里没有对任务 ID 排序两个并发事务分别持有对方需要的行锁形成环形等待。这个问题的修复很简单批量更新前先按task_id升序排列所有节点按同样的顺序获取行锁。同时把单次批量更新的行数限制在 200 行以内降低单个事务的持锁时间。从那以后死锁日志基本绝迹。这个经验不是调度系统独有的任何大规模状态更新的场景都适用。7. 最后再分享两点实在的如果你打算自研或者选型调度器我最后的建议分两层。一是技术上的把时间当作一等公民去建模所有跟时间相关的字段必须区分语义是到期时间、实际执行时间、心跳时间还是补偿时间混在一起迟早出事。二是工程上的调度故障往往不是单点问题而是锁、状态机、心跳、幂等这些能力缺一环导致的所以要重视可观测性核心指标至少包括调度延迟、派发失败率、重试分布和任务执行耗时。有个小技巧也一并说了吧。我们给每个任务都加了一个owner_service字段记录投递方归属。这个字段平时没什么用但每次任务异常需要拉群找人时它直接告诉你该找哪个服务的人少了很多沟通成本。ax这套东西没有用什么高深算法贵在把工程细节抠到位。希望你读完后对调度系统的设计取舍有自己的判断也能避开那些反复出现的暗坑。
📝

华诺云谱内容团队

资深建站顾问 · 行业研究员

10年+企业数字化服务经验,专注智能建站、SEO优化与品牌营销,持续输出建站技巧、行业洞察与营销干货,已帮助5000+企业实现数字化增长。

你可能需要的服务

订阅华诺云谱资讯周报

每周一封,精选建站技巧、SEO与营销干货,直达邮箱。已有 8,000+ 企业主订阅,助你少走弯路。

↑