ARTICLE DETAIL

资讯详情

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

自研分布式调度器ax:从Cron到时间轮与高可用集群实践

自研分布式调度器ax:从Cron到时间轮与高可用集群实践 1. 从Cron到ax一次凌晨两点的漏调度事故逼出来的自研调度器先说背景。我在团队里负责基础设施这块主要工作是维护各类业务跑批任务、数据同步任务和定时报表。最早用的是Linux自带的Cron后面业务量上来了换成Quartz集群再后来又试过开源的分布式调度平台。这套链路上踩的坑基本能写一本《定时任务灾难题集》。真正让我下定决心自研的是那次凌晨两点的漏调度事故。那是一个数据仓库的日级同步任务每天凌晨两点准时启动把前一天的业务库增量数据同步到分析库。某天早上运营反馈报表数据没更新我上去一看任务居然没跑。查Cron日志发现配置还在但进程因为OOM被系统杀掉之后Cron完全没感知。任务就这么静悄悄漏掉了。那天的数据补跑花了大半天而原因仅仅是“进程死了没人通知”。那次之后我把目光转向了市面上成熟的调度框架但调研一圈下来发现每个方案都有那么几个让我不舒服的地方Quartz集群依赖数据库锁锁竞争一上去调度延迟就飘忽不定。xxl-job之类的平台功能全但偏重要部署调度中心还要写一堆业务对接代码。用Celery做定时任务它在分布式场景下对“同一时刻只跑一次”的全局锁做得不够干脆。最麻烦的是没有一套方案能同时满足几个我们特有的需求秒级任务、任务依赖编排、精确一次的语义、以及调度器自身状态的可观测性。于是就有了ax这个项目。ax这个名字没什么深意就是当时在文档里随手敲的两个字母后来说“ax调度”说顺口了就一直沿用。很多人都问我为什么不用现成的轮子我的回答是不是轮子不好是你的车床形状太怪。我们当时有大量短周期任务30秒到5分钟一次有跨系统依赖的DAG需求又要求调度器本身不能有单点。市面上的通用调度器单点问题能解决但依赖编排和秒级精度很难同时满足。这一篇我就把ax调度的完整实践写出来从架构设计、核心模块、集群模式、任务编排到生产环境踩过的坑和最终的压测结果全程干货。如果你也在纠结“自研调度器到底值不值得”这篇应该能帮你想清楚。2. 架构设计的取舍逻辑为什么用Go写调度核心而不是Java或Pythonax调度的整体架构一句话概括是调度中心负责算“什么时候该跑哪些任务”执行器负责“具体把任务跑起来”两者通过gRPC通信元数据放MySQL分布式锁和实时状态用Redis。这套架构看着不稀奇但里面每个选型背后都有真实的踩坑经历。2.1 调度中心与执行器分离是清理历史债务的第一步最早用Cron的时候调度逻辑和业务代码在同一个进程里。任务多了之后一个任务OOM会拖垮整个进程所有定时任务一起完蛋。后来用Quartz虽然支持集群但任务代码还是要和调度器耦合在一个应用里升级一个任务要重启整个调度服务。所以ax从设计第一天就把调度和执行彻底分开了。调度中心是纯状态机不跑任何业务代码只负责触发执行器是一个独立的Agent进程部署在业务机器上收到调度指令后拉起具体的Worker子进程去执行业务逻辑。这样一来业务任务再怎么写都不会污染调度核心的稳定性。换任务代码只需要重启Agent调度中心完全不受影响。这么说吧调度中心就像快递分拨中心它只管把包裹按地址分到对应线路执行器是各个站点的快递员具体送上门的路线、敲门方式都是快递员自己的事。分拨中心不需要知道快递员怎么骑电动车。2.2 Go在调度场景下的优势并发、编译期控制和内存占用语言选型我权衡了很久。Java的生态最成熟Quartz、XXL-Job都是现成的但如果自研还要选Java等于背着包袱跑。Python写起来爽但GIL在并发调度上容易卡脖子而且部署分发不方便。最终选了Go理由如下调度器本质是高并发的“计时器分发器”Go的goroutine天然适合这个模型。一个时间轮上有几万个待触发任务每个触发动作起一个goroutine去投递内存开销很小。Go编译出来是单二进制部署到服务器上看不到一堆JAR依赖。静态类型在维护调度规则这种核心领域里比动态语言安全得多。改一个字段有编译错误兜底不会上线后才炸。我见过太多团队在自研调度器时用Java硬写结果光是把Spring那一套配置搬进调度节点就耗了大半个月。用Go写从零到跑通核心调度链路我当时的实际耗时是一个周末。但Go也有一个不算缺点的缺点生态里没有特别成熟的分布式调度基础设施几乎每一块都得自己拼。这正好也是自研的意义所在。2.3 存储层选型MySQL存元数据Redis存实时状态调度器需要持久化的数据有两类任务定义Cron表达式、执行器地址、超时时间、重试次数和调度历史每次触发的记录。这类数据用MySQL最方便事务、索引、备份体系都是现成的。我见过有人非要把这类元数据放MongoDB或ES里查询和事务都变得很别扭。任务定义表的核心字段大概是这样的CREATE TABLE task_def ( id bigint(20) NOT NULL AUTO_INCREMENT, task_id varchar(64) NOT NULL COMMENT 全局唯一任务ID, name varchar(128) NOT NULL, cron_expr varchar(64) DEFAULT NULL COMMENT Cron表达式定时任务用, delay_seconds int(11) DEFAULT NULL COMMENT 延迟秒数延迟任务用, executor_group varchar(64) NOT NULL COMMENT 执行器分组, payload text COMMENT 业务参数透传给Worker, timeout_ms int(11) NOT NULL DEFAULT 60000, max_retry int(11) NOT NULL DEFAULT 0, priority int(11) NOT NULL DEFAULT 5 COMMENT 1-10越大越优先, status tinyint(4) NOT NULL DEFAULT 1 COMMENT 1启用 0暂停, create_time datetime NOT NULL, update_time datetime NOT NULL, PRIMARY KEY (id), UNIQUE KEY uk_task_id (task_id) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;实时状态哪些任务正在跑、哪些节点活着、锁信息放Redis因为这类数据变更频率高MySQL扛不住高频UPDATE而且Redis天然支持分布式锁和过期时间省了自己写超时清理的逻辑。我见过有人在Redis里存调度记录做持久化这在任务量大的时候会非常浪费内存。Redis只做实时状态历史记录走MySQL异步落库这是最省心的组合。3. 时间轮与调度循环一个调度器最核心的“心跳”调度器的灵魂是它如何高效地知道“现在有哪些任务该被触发”。最简单的做法是每个任务起一个定时器goroutine任务多了之后几万个goroutine同时挂在那里等触发内存和调度开销都很夸张。更严重的是Go的定时器在高并发场景下会有精度毛刺集群节点一多整体触发时间的抖动根本控制不住。ax的调度核心改成层级时间轮之后这个问题才算彻底解决。时间轮的基本思路和钟表很像一层60个槽位代表秒一层60个槽位代表分一层24个槽位代表小时任务根据到期时间落在对应层级的槽位上指针每跳一下只处理当前槽位上挂着的任务。3.1 层级时间轮的实现要点我参考了Kafka的层级时间轮设计但针对调度场景改了不少东西。核心数据结构长这样type TimingWheel struct { tickMs int64 // 每一格代表的时间秒级任务设为1000ms wheelSize int64 // 格子数量默认60 currentTime int64 // 当前指针位置 buckets []*TaskBucket // 每个格子挂一个任务桶 }每次调度循环的伪逻辑指针走到当前时间所在的格子。取出该格子上的所有任务。判断任务是否已经到期到期就触发未到期就降级到更细的层级重新入轮。触发方式不是直接调用业务代码而是往执行器的消息队列投递一条“执行指令”异步解耦。层级时间轮的好处是一万个任务挂上去不管是最近一秒要跑的还是明天凌晨三点要跑的占用的内存都差不多不会因为任务数量线性膨胀。这块我实测过5万条任务挂在时间轮上调度器本身的常驻内存增量只有40MB左右。3.2 调度循环里的“精准触发”细节时间轮解决了“谁该触发”的问题但“什么时候触发”还有个隐藏的坑节点上的时钟可能不准。云服务器默认启用NTP但NTP同步本身有延迟而且宿主机时钟跳变会直接导致调度提前或延后。ax的做法是不直接用系统时钟做触发判断而是用系统时钟加一个校正偏移量。调度中心每个节点启动时会向一组时间基准节点发起时间同步算出本地时钟与基准时钟的偏差每次触发判断都减掉这个偏差。虽然做不到绝对精准但实测能把跨节点的调度时间差控制在50ms以内。对于绝大多数业务场景这个精度完全够用。再补充一个细节触发时间到了之后调度指令不是立刻发出去而是先进一个内部投递队列。投递队列的消费者负责执行真正的RPC调用。为什么中间要隔一层因为如果任务集中到期比如整点批量任务直接并发发RPC会把执行器的连接池打满。投递队列配合信号量控制并发数等于给执行器加了一层限流保护。3.3 秒级任务必须单独处理的一个坑时间轮能处理秒级任务但Cron表达式本身在秒级场景下有精度限制。Linux Cron只能精确到分钟我们早期业务里有个每30秒跑一次的任务当时用Cron根本没法表达只能写两个重叠的分钟级表达式再配合去重逻辑非常不优雅。ax支持两种任务类型Cron表达式分钟级及以上和周期秒级表达式EVERY 30s。秒级任务在时间轮上的实现逻辑是每次触发后计算下一次触发时间点当前触发时间间隔然后把任务重新挂回时间轮。这里有个细节必须注意计算下一次触发时间点时不能用“当前时间间隔”因为任务执行可能耗时长如果执行了5秒才结束下一次触发就会顺延5秒。正确做法是记录上次触发的时间点下一次触发时间上次触发时间间隔。这样哪怕任务执行了50秒间隔30秒它也会把错过的触发次数补上。这块我写成了一段测试代码专门验证func TestNextTriggerTime(t *testing.T) { lastTriggerTime : int64(1000) // 上次在1000ms触发 interval : int64(30000) // 间隔30秒 now : int64(36000) // 当前时间已经过了36秒 next : lastTriggerTime interval // 下一次应该是31000ms if next now { next now // 如果已经严重过期不追赶直接从当前时间开始重新计算 } t.Logf(next trigger: %d, next) }这里还有个业务侧的选择问题任务积压严重时是补跑错过的次数还是直接跳过从头算我们的答案是基于任务的业务属性配置。数据同步类任务通常要补跑告警监控类任务通常直接跳过因为告警讲求实时性补一个5分钟前的告警没有意义。4. 集群模式与任务不重不漏全局锁、心跳和故障转移调度器是高可用要求的核心。如果调度器挂了所有定时任务全部瘫痪这是绝对无法接受的。ax的集群模式围绕三个核心问题展开谁来调度、谁在干活、挂了怎么办。4.1 Leader选举调度动作收口避免多节点重复触发时间轮每个节点都必须跑但如果每个节点都触发同一个任务就会造成重复执行。ax采用了典型的Leader-Follower模式同一时刻只有一个节点Leader负责从时间轮触发任务Follower节点只同步任务定义不执行触发动作。Leader选举用了Redis分布式锁。每个节点启动时尝试获取一个名为ax:leader的锁拿到锁的节点成为Leader。锁的过期时间设了30秒Leader每隔10秒续期一次。如果Leader宕机锁在30秒后自动过期其他节点竞争获取完成Leader切换。这里有个坑必须得说Redis锁获取到了但Leader节点上的任务定义和内存状态可能还是不完整的。所以新Leader上位后不是立刻开始调度而是先从MySQL全量加载一遍任务定义重建时间轮再启动调度循环。整个过程大约需要2到5秒业务侧会感知到一次短暂的调度空窗。对于绝大多数分钟级任务来说这个空窗可以接受。如果需要秒级任务也不中断那就要上真正的热备方案了我们当时没有这个刚需。4.2 任务不重不漏MySQL唯一索引 Redis原子操作双重保险Leader只负责触发但触发后任务可能失败、执行器可能没收到、网络可能抖动。这种情况下调度器既要避免“一个任务被多个节点同时执行”又要避免“任务失败了却没人处理”。ax的任务状态流转是这样的待触发 - 已触发 - 执行中 - 成功/失败/超时用户任务侧执行器收到调度指令后会在Redis里执行一段Lua脚本尝试“抢任务”只有抢成功的才真正拉起Worker子进程执行。抢任务的脚本核心逻辑是这样lockKey KEYS[1] -- ax:task_lock:{taskId} statusKey KEYS[2] -- ax:task_status:{taskId} -- 如果任务正在执行中说明前一次还没跑完不重复触发 if redis.call(GET, statusKey) running then return 0 end -- 原子写入锁和状态 redis.call(SET, lockKey, ARGV[1], PX, ARGV[2]) redis.call(SET, statusKey, running, PX, ARGV[2]) return 1这样设计之后同一时刻一个任务在全集群范围内只会被一个Worker执行。“不重”解决了但“不漏”还需要另外一套机制兜底。调度器在触发任务之前会先在MySQL的任务执行记录表里插入一条状态为TRIGGERED的记录带有唯一键(task_id, trigger_time)。如果两个节点同时尝试插入唯一键会拦住一条。等执行器确认开始执行后这条记录才被更新为RUNNING状态。谁会去更新这条记录执行器在Worker启动成功后会回调调度中心调度中心再更新状态。这套双保险的代价是每次触发都多了一次MySQL写入和一次Redis调用。实测单次触发链路增加约3ms延迟换来的是明确的“不重不漏”语义值。4.3 执行器失联了怎么办超时、重试和最终一致执行器节点可能会突然宕机、网络分区、或者机器被回收。如果Worker正在跑一个耗时任务而执行器失联了调度器怎么处理ax的做法是每个任务定义里都有timeout_ms参数调度中心会记录任务触发的时间点如果超过超时时间还没有收到完成回调判断任务执行失败。失败后根据max_retry决定是否重试重试会重新投递给执行器分组里的其他节点。整个过程的最终一致性靠的是一个每分钟跑一次的“扫描补偿任务”。它会扫描所有处于执行中状态且超时未完成的任务记录将它们重置为待触发或标记失败。这个补偿任务的实现有一点需要注意扫描不能太频繁否则会在任务量大的时候给MySQL带来额外压力。每分钟扫描一次配合超时判断已经能让最坏情况下的漏调度时间控制在1到2分钟。5. 任务依赖编排DAG调度和动态参数传递做调度器最容易忽略、但业务方最痛的一个需求是任务B必须在任务A成功后才能跑而且B要用A的输出作为入参。Cron解决不了这个Quartz处理起来也很别扭。ax支持轻量级的DAG编排这部分是业务接入后反馈最好的一块。5.1 依赖触发模型不预定义整张图而是每个任务声明上游很多调度平台的做法是画一张DAG图节点连线表示依赖。这种大而全的方案学习成本高而且任务的增删改都牵一发动全身。ax选择了更轻的模型每个任务定义里加一个upstream_task_ids字段声明自己依赖哪些上游任务。{ task_id: sync_order_to_dws, name: 订单数据同步到数仓, upstream_task_ids: [sync_order_to_ods, sync_payment_to_ods], executor_group: data-sync, cron_expr: 0 2 * * * }调度器扫描发现sync_order_to_ods和sync_payment_to_ods都执行成功后自动触发sync_order_to_dws。依赖状态存在Redis里每个DAG节点执行结束后会更新所有下游的依赖完成度完成度达到100%且上游状态均为成功下游任务进入待触发队列。这种模型的缺点是没法表达复杂分支但胜在理解成本低。业务方不用学习任何DAG概念只需要在配置里写“我是谁的下游”就够了。5.2 上游输出传给下游上下文透传的实际做法任务B需要任务A的执行结果这是最常见的真实需求。ax的设计是每个任务执行完成后Worker可以把一段JSON写入到执行记录里调度器会把它存进task_result字段。下游任务触发时调度器会把所有已完成上游的task_result合并成一个JSON对象放进下游任务的Payload里一起下发。这个机制在实际跑批场景里帮了大忙。我们有条链路是抽取订单数据 - 清洗转换 - 汇总统计。每一步的输出都是下一步的输入以前靠人工拼接参数现在全部自动传递。有一点要注意上游结果不能无限大我们限制单个上游的task_result不能超过512KB超出部分走对象存储通道只把存储路径传给下游。这个上限在配置文档里写清楚业务方就很少踩坑了。5.3 失败了整条链路怎么办失败策略选项DAG中一个节点失败链路怎么处理ax给了三个选项失败后下游不触发默认最安全、失败后跳过该节点继续触发下游适合非关键依赖、失败后等待手动重试成功再触发适合核心链路。这块没有银弹只能把选项开放给业务方。我见过不少团队在这一步卡住因为他们把所有任务的失败策略都配成“失败后跳过”结果上游漏了数据下游还在照常跑报表出来以后全是对不上的数字。最终我们定了一个原则数据链路全部用默认策略只有监控链路允许跳过。6. 生产环境的硬故障复盘从状态不一致到重复调度自研调度器最大的风险不是功能不够而是生产环境里出现的各种你在设计时压根没想到的状态。下面几个故障每一个都直接打到了我们的设计盲区也是ax迭代中最有价值的部分。6.1 重复调度Redis锁过期时间设置不当引发的“双跑”第一次线上事故现象是同一个任务跑了两次而且两次都成功。排查链路如下打开执行器日志发现两次执行指令来自不同的调度节点。确认两个节点都持有过Leader锁。查Redis锁记录发现锁在30秒后过期释放但原Leader节点的任务处理能力弱触发消息还在内部队列里排队。新Leader上位后重新触发同一任务两条消息先后到达执行器。执行器抢锁时第一条消息把状态写成了running并设置了锁过期时间但第二条消息到达时锁已经过期因为拖得太久于是也成功抢锁。根因是锁的过期时间必须大于“从触发到执行器完成抢锁”的最长耗时。而我们的内部队列在消息积压超过30秒时就打破了这个约束。修复方案加了一个“触发时间防重窗口”。任务定义上配置了dedup_window_seconds执行器收到触发指令后先用Redis检查这个任务在最近N秒内是否已经触发过如果是就直接丢弃。这个窗口现在默认设置成60秒比所有任务的锁过期时间都长彻底堵住了双跑的可能。6.2 任务消失Leader切换时丢了内存中的时间轮第二次严重事故现象是某些隔天任务在Leader切换后完全消失值班日志里看不到任何触发记录。原因很直接Leader的内存时间轮是纯内存结构切换前没有持久化。旧Leader挂了以后它内存里的待触发任务随之丢失新Leader从MySQL重建时只重建了任务定义没有重建“已经进入时间轮等待触发的实例”比如手动触发的临时任务、延迟任务。修复把所有进入时间轮的任务实例都额外同步写一份到Rediskey是ax:wheel:task:{triggerTime}:{taskId}新Leader启动时扫描Redis重建时间轮。用Redis做时间轮的持久化层有一个额外好处Follower节点也能看到待触发任务列表做运维复盘的时候方便查看。6.3 时钟漂移某台机器的调度时间比别人快了8秒第三次事故比较隐蔽。某个数据对账任务每天凌晨0点触发有几天凌晨0点0分8秒就有结果了但调度日志显示触发时间是23点59分52秒。排查发现跑调度中心的这台云服务器时钟比时间基准快了约8秒。调度器判断“当前时间已到触发点”时用的是本地时钟本地时间比真实时间快导致提前触发。更麻烦的是NTP客户端默认在校时时会逐步调整而调校期间这台机器的本地时间一直是偏快的。ax现在的处理方式是调度节点禁用NTP自动校时改为调度器应用层主动向基准节点同步时间偏移量这个偏移量在触发判断时动态扣减。这里有个细节偏移量不是一次性校正的而是每隔30秒重新校准一次因为云服务器和基准之间的网络延迟本身也会浮动。最终效果是任务触发时间误差从最坏8秒降到了平均30ms以内。6.4 状态不一致Redis数据没了MySQL里还是成功还有一种低频但让人头疼的问题Redis因为内存策略清理了任务状态或者执行器调用回调时网络抖动任务实际上成功了但调度中心标记为超时失败触发了重试。重试导致业务侧幂等性问题。解决方案是在业务侧约定一个原则所有任务处理函数必须幂等即同一份数据执行两次结果必须一致。调度器能保证的是“大多数情况下不会重”但任何分布式系统都没法做到100%不重所以幂等是业务方的责任。这个原则写进了开发规范里比任何技术方案都管用。因为我们见过太多团队把精力耗在“如何保证永不重试”上而实际上只要业务处理函数本身幂等重复执行根本不影响结果。7. 性能压测与容量规划一万个任务跑起来到底要吃多少资源自研调度器最终要回答的问题是它能撑住我们业务量的多少倍。我们做了一轮相对完整的压测这里把数据放出来给打算参考ax设计的人一个量级概念。压测环境是三台4核8G的云服务器跑调度中心五台4核8G的执行器MySQL和Redis都是标准云产品不算高性能规格。压测里我们建了1万个定时任务其中9000个是分钟级任务每分钟触发一批1000个是30秒级任务。每个任务触发后执行一个模拟业务处理处理耗时控制在50ms以内。结果数据指标数值触发QPS峰值约330次/秒调度中心CPU峰值约60%三台节点均分调度中心内存峰值每节点约800MB时间轮任务容量5万加挂后内存增量约40MB触发到执行器收到指令的P95延迟18ms触发到执行器收到指令的P99延迟42msRedis操作QPS峰值约1600次/秒看数据能得出几个结论第一调度器本身的瓶颈不在CPU而在Redis操作和MySQL写入。每次触发至少有一次MySQL插入和一次Redis状态查询这是“不重不漏”语义的代价。如果对高并发有更高要求可以考虑把记录写入改成批量异步但换来的是一致性弱化。第二执行器的压力取决于任务本身的逻辑调度中心再怎么压都很难成为瓶颈。实际上调度中心跑在1核1G的容器上都能勉强撑住万级任务但为了Leader切换的稳定性我们还是建议至少4核。第三如果触发QPS超过500次/秒就不要再勉强用单组Redis了要把锁和状态拆到不同的Redis实例里避免互相干扰。容量规划的建议是按业务未来一年任务量的2倍做基准调度中心节点数不少于3个每节点至少4核8G。MySQL的task定义表10万条以内非常轻松执行记录表建议定期归档我们保留30天的在线记录超过30天的转冷存储。8. 从自研到上线最被低估的其实是“可观测性”最后聊一个技术之外的话题自研调度器最大的隐性成本不是写代码不是搞架构而是可观测性。为什么这么说因为调度器是“幕后角色”它出问题的时候业务方只会发现自己的任务没跑而不会第一时间联想到调度器挂了。ax在这块做了三件事第一调度事件全链路日志。从任务触发、指令下发、执行器确认、Worker启动、执行完成、结果回传每个环节都有一条结构化日志带上trace_id。业务方报问题的时候只要把任务ID和时间发过来通过trace_id就能串出完整调用链。第二核心指标暴露给监控系统。包括每个节点的触发QPS、时间轮容量、投递队列长度、Redis调用耗时、MySQL写入耗时、Leader状态。投递队列长度这个指标特别重要它一旦持续上涨就说明调度中心在积压任务迟早会出问题。我们配置的告警阈值是队列长度超过1000持续5分钟。第三调度大盘。每次触发、成功、失败、重试的实时统计按执行器分组、按任务类型拆分。大盘不是为了好看是为了在问题发生时能第一时间判断影响范围——是某个任务的问题还是某个执行器分组的问题还是整个调度器的问题。这个能力在故障响应的前5分钟特别值钱。如果你也在搞自研调度器我强烈建议把可观测性放到和核心功能同等重要的位置从上线的第一天就开始做不要等出事了再补。这比任何花哨的调度算法都更能帮你保住周末的睡眠时间。
返回列表