ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

自研分布式调度系统ax:从cron痛点到大促流量控制的完整复盘

自研分布式调度系统ax:从cron痛点到大促流量控制的完整复盘 去年的一个晚上线上结算任务的告警把我从饭桌上拽回工位。排查了一圈发现是一台凌晨扩容的机器没有同步crontab配置整整四十分钟的账务数据没跑出来。也是从那天起我开始动手搞一个内部代号叫“ax”的任务调度系统。当时根本没想什么框架选型目标只有一个让所有定时任务、延迟任务、批量分片任务都有统一入口、统一监控、统一恢复手段。后来这套东西在组里跑了一年多接了几万个任务我把它踩过的坑和设计取舍完整复盘一遍希望能给同样被cron折磨的团队一点参考。这里说的“ax调度”不是某个现成开源软件的改名而是一个带流量控制和资源感知的分布式调度组件。它解决的核心问题有三个任务不丢、定时可依赖、高峰期不会把后端冲垮。适合那些已经觉得“cron 脚本 人肉运维”撑不住、但又不想直接上重型工作流引擎的团队看。1. 为什么会有 ax 这个项目1.1 一点就着的现实痛点最开始我们和大多数团队一样用的是服务器crontab外加一个简单的时间轮工具。业务量小的时候完全够用但一旦任务量涨到几千个就会暴露出几个非常恶心的场景一是任务散落在几十台机器上谁在跑、跑没跑、要不要重跑全靠人肉记忆。二是节点重启、扩容、缩容之后cron配置要么丢失要么重复没有任何自愈机制。三是业务方开始提各种“扭七扭八”的需求比如某个数据文件想在每天的9点30分之后延迟两分钟再拉取比如大促期间希望特定任务自动降级、跳过非核心批次。这些东西用原生cron根本表达不出来。说白了我们需要的不是一个“定时触发器”而是一套能把任务送到正确机器、在正确时间执行、并且对结果负责的调度基础设施。ax就是在这个背景下立项的。1.2 选型对比为什么不直接上现成框架网上开源的调度系统一抓一大把我也认真对比过Quartz、XXL-JOB、Airflow这几个主流方案但最后都没有直接用原因挺现实Quartz是单机起家的虽然支持JDBC集群模式但它的锁机制在高并发调度下会成为瓶颈而且 Quartz 对“分片”“动态路由”“失败重试策略”这些概念支持得很弱。XXL-JOB功能确实全调度中心加执行器的模型也很清晰但整个体系偏重对我们这种很多任务只是简单HTTP回调的团队来说部署运维成本和二次开发成本都不低。Airflow则更适合数据管道场景分钟级起步对于秒级定时和业务性任务调度并没有做得更好。还有一个更实际的原因我们很多任务并不是传统意义的“定时任务”而是“事件到了之后立刻找一个空闲worker执行”。这本质上更像一个消息队列加路由器的结合体。拿现成调度框架去适配这种场景就像拿轿车底盘去跑越野能走但总不得劲。1.3 定义清边界ax 到底调度什么立项时我们给ax划了三条边界这决定了后面所有设计第一层是时间调度包括定时任务、延迟任务和周期任务。第二层是资源调度一个任务拆成N片分发到不同worker并行执行同时感知每个worker的负载避免把任务全扔给同一台机器。第三层是流量调度在业务高峰期通过令牌桶控制任务下发速率相当于给任务执行装了一个水龙头防止下游数据库、外部接口被瞬时打爆。这个边界定义非常重要。如果只是做一个定时任务框架那核心逻辑就是“到点触发”复杂度和价值都有限。把资源调度和流量调度一起塞进去之后ax就从“闹钟”进化成了“交通指挥”这也是后来团队内部说“ax调度”时真正想表达的东西。2. ax 整体架构与核心模型2.1 从调度中心到执行代理的分层设计ax的整体结构分为三层接入层、调度层、执行层。接入层提供统一的API和SDK业务方可以通过HTTP接口或者Java/Python的SDK创建任务、查询状态、手动触发。调度层是核心由一个或多个无状态的scheduler组成它们共同消费同一个任务元数据存储通过分布式锁保证同一时刻只有一个scheduler对某个任务实例负责。执行层就是部署在业务机器上的ax agent它负责拉取分配给自己的任务分片并真正执行。这个分层参考了Master-Worker模型但没有做强Master。scheduler之间是对等的任何一个调度节点宕机其他节点能通过leader竞选接管任务不会出现“调度中心挂了全盘停摆”的恐怖局面。agent和scheduler之间通过长轮询做任务拉取也就是说调度器不会主动把任务推给worker而是worker主动来认领。做成长轮询而不是长连接推送是我个人比较坚持的一个点。长连接推送对网络波动太敏感一旦连接断开服务端很难判断是网络问题还是worker死亡容易误判重试。而长轮询简单粗暴worker每两秒来问一次“有没有我的活”有就把任务领走没有就挂起请求等调度器这边hold住等有任务了再响应。这样即使worker短时间断网重连之后也能靠轮询自然恢复不需要额外的对账机制。2.2 任务模型里那些必须有的字段任务数据模型是ax最核心的地基。第一版测试的时候我只设计了id、name、cron、url四个字段结果上线一周就翻车了因为根本没法回答“这个任务失败后该怎么办”“这个任务能不能并发跑多个实例”这些问题。后来我参考了多个成熟系统的做法把任务字段重新梳理成几组基础信息包括id、name、group、ownerowner尤其重要任务出问题第一时间能够找到责任人。触发信息包括trigger_type取值有cron、interval、delay和manual四种cron表达式统一存cron_str还要配上timezone字段解决跨时区调度的老大难问题。执行信息包括executor_type和executor_info比如HTTP任务就把url、method、header放进去脚本任务就把命令行放进去。控制信息包括timeout_ms、retry_times、retry_interval_ms、priority这四项决定了任务异常时的表现。分片信息包括shard_count和shard_params例如任务要处理全量用户shard_count10每个分片会自动获得一个0到9的分片序号业务代码按序号取模就行。整套模型核心思路是把“任务是什么”和“任务怎么调”分离开。调度器只关心trigger、timeout、retry这些调度属性完全不理解任务内部的业务逻辑。执行器只看executor_info拿到就执行。这样即使以后业务方接了一个完全不明所以的任务类型调度体系本身也不用改。2.3 分片和优先级两个最容易被忽略的设计分片功能看起来很简单但细节决定成败。我们早期有一个批处理任务每天要跑全量用户的数据刷新单机跑要两个小时。后来给这个任务配置了shard_count10按user_id的hash值分片十个worker同时干整个执行时间降到15分钟以内。但分片不能光有数量还得有“分片感知”。worker在启动一个分片之前必须通过调度器拿到分片参数比如分片序号、总分片数、以及该分片对应的数据范围。我们实际设计了一个shard_param的生成器允许任务注册时自定义分片规则有的任务按用户尾号有的任务按订单时间区间有的任务直接让调度器分配连续区间。这样不同业务场景都能复用同一套拆分机制。优先级则更像一个砖家级别的隐藏设定。业务方通常会认为自己的任务都是最高优先级但资源是有限的。ax采用了两级优先队列第一级先按priority值从高到低排队第二级在同一优先级内按创建时间从早到晚排队。调度器每次从队列头部拿出一个任务实例进行下发保证高优任务不会被低优任务的长队列饿死。3. 调度引擎核心细节拆解3.1 时间轮让百万定时任务不靠扫表触发最早实现ax触发逻辑时我用的最笨办法每秒扫描一次任务表把当前时间范围内所有到期的任务捞出来。这个方案在几百个任务时完全没问题等到任务量到几万、几十万级别每秒一次全表扫描让数据库压力直接拉满而且扫描周期决定了任务触发的精度最多只有一秒很多秒级任务根本排不上。最后换成了时间轮算法并参考Netty的HashedWheelTimer实现了分级时间轮。默认基础轮是256个槽每个刻度0.5秒这样一轮覆盖128秒不够用就在上一级再加一个更大的轮共同构成一个多层级时间轮。每次新增一个延迟任务就计算它和当前时间的差值根据差值决定放入哪一级时间轮的哪个槽位。调度器每秒或者每0.5秒走动一个刻度只处理当前刻度槽位上挂载的任务即可。这样处理复杂度从“所有任务量”变成“每刻度需要触发的任务数”几十万定时任务轻松扛住。时间轮还有一个额外的好处天然支持延迟任务。比如业务方希望通过ax实现“下单30分钟未支付自动关闭订单”只需要创建一条trigger_typedelay、delay_ms1800000的任务它会安静地待在时间轮里等待完全不需要额外的Redis过期监听或者数据库轮询。3.2 任务分配一致性哈希和负载感知的结合调度器确认一个任务实例触发之后接下来要解决的是“把这个任务给哪个worker执行”。最简单的做法是把所有worker拍扁到列表里轮流分配或者随机分配但很快发现两个问题第一某些任务在一个worker上有本地缓存频繁换节点会导致缓存命中率骤降第二worker配置不一样有的机器16核有的机器4核平均分配会累死小机器。ax的任务路由策略是两阶段的第一阶段用一致性哈希把任务id映射到哈希环保证同一个任务在正常情况下总是路由到同一批节点给任务执行的局部性留出空间。第二阶段是负载纠正worker会周期性上报自己的CPU、内存和排队任务数agent在认领任务之后调度器会比对worker最近3分钟的负载均值如果发现某个节点明显超载就临时把后续任务偏移到其他备选节点。后来我们干脆为两种场景提供了不同route_mode。任务本身偏向IO密集型的大文件处理就设成LOAD_BALANCED任务内部重度依赖本地缓存就设成HASH_BY_TASK。这样既照顾了通用性又给业务方留下灵活性不用为了性能来求我们改调度策略。3.3 超时、重试和幂等保证任务不丢不漏任务调度世界里最尴尬的问题是任务到底执行成功没有网络抖动可能导致结果回传失败worker在任务执行到一半时宕机消息语义就变成未知。ax一开始就明确采用at-least-once语义也就是说任务可以重复执行但不能不执行。基于这个语义我们给每个任务实例生成了一个全局唯一的execution_id。无论是调度器本身还是worker所有日志和状态流转都带上这个ID。任务执行完成后worker会往调度器回传结果调度器更新状态为SUCCESS。如果worker在超时时间内没有回传调度器就会判定超时把任务重新放回时间轮等待重试。超时时间的设置很关键。如果设得太短稍微慢一点的业务请求会被误判超时造成大量重复执行。如果设得太长遇到真正卡死的任务下游会被拖死。我们的经验值是取任务近期P95执行耗时的三倍最低5秒最高不设上限并在首次超时后自动按指数退避递增。exponential backoff的底数设成2初始间隔5秒最多重试3次总共等待时间不会超过100秒业务方不用面对“任务失败之后突然又活过来”的灵异事件。幂等保护则分两层调度器通过Redis的SETNX记录execution_id保证同一个任务实例不会被重复调度worker端在执行业务逻辑之前先检查本地数据库的幂等表同样以execution_id作为唯一键已经处理过的直接跳过。双保险之后重复执行的概率被压到极低。3.4 流量调度令牌桶给底层系统加一道保险时间轮把任务按时触发出来一致性哈希把任务送到合适的worker但这还没完。大促场景下整点齐刷刷几百个任务同时触发每个任务可能都要查一次数据库这种瞬时尖峰足以把数据库连接数打满。ax的流量调度模块其实很简单一个全局令牌桶单位时间只能释放rate个令牌。调度器在进入分派流程之前必须先去令牌桶申请令牌拿到令牌的任务才能继续往下走拿不到的一律排队等下一轮。rate不是静态配置而是根据任务分级动态调整比如P0任务的权重高P2任务碰到高峰期会被降级令牌只给P0和P1用。这个能力上线之后效果很直观。以前每天0点数据库CPU飙到90%加了流量调度以后峰值降到70%左右虽然任务总耗时变长了但每个任务都在预期时间内完成没有出现连接池打爆导致的雪崩。这其实就是把“所有任务一拥而上”变成了“均匀洒水”底层系统的压力曲线平滑了很多。4. 实操从零接入一个 AX 定时任务4.1 环境准备和最小化部署要复现ax调度这套流程不需要三台机器一台4核8G的服务器加一台执行任务的worker就够跑通全流程。环境上准备MySQL 8.0和Redis 6.x调度器和agent都是打包好的JAR包。数据库初始化只需要三张核心表task_def保存任务定义task_instance保存每次触发后的执行实例task_log保存分片和worker的执行日志。建表语句不展开说几个容易踩的坑task_def的cron字段要区分cron_expression和timezone两个column否则后面跨时区调度一定踩坑task_instance的execution_id要建唯一索引这是幂等实现的基础task_log的worker_ip和shard_id要建组合索引排查问题全靠它。agent部署很简单配置好server地址和机器唯一标识就可以启动。agent启动后会自动向调度器注册同时上报自己的CPU、内存、磁盘等信息。为了调试方便第一版可以关闭自动注册改用手动配置worker_id等确认运转正常再放开。4.2 定义一个任务从API到SDK两种方式最快速接入ax的方式是通过HTTP API创建任务。比如业务方希望每天早上8点整拉取一次商户账单数据只需要一个POST请求curl -X POST http://ax-server:8080/api/task/register \ -H Content-Type: application/json \ -d { name: merchant_bill_daily_pull, group: bill, owner: zhangsan, trigger_type: cron, cron_str: 0 0 8 * * ?, timezone: Asia/Shanghai, executor_type: http, executor_info: {\url\:\http://billing.internal:8082/pull\,\method\:\POST\}, timeout_ms: 60000, retry_times: 3, retry_interval_ms: 5000 }接口返回值里会带着task_id。调度器收到请求后会校验cron表达式然后写入task_def表并启动任务定义级别的触发注册。整个过程无需重启任何服务新任务在下一个时间轮刻度就会生效。不过HTTP API适合简单场景涉及复杂分片和自定义逻辑我用得更多的是SDK方式。Java SDK里任务就是一个加了注解的方法AxJob(name user_bucket_rebuild, cron 0 0 3 * * ?, shardCount 10, timeoutMs 300000, retryTimes 2) public void rebuildUserBucket(AxShardParam int shard, AxShardTotal int shardTotal) { // 只处理属于当前分片的那部分用户 ListLong userIds userService.queryUsersByShard(shard, shardTotal); bucketService.rebuild(userIds); }SDK在执行前会把shard和shardTotal注入进去业务方法对“自己在跑哪个分片”完全透明。这里有个易错点shardCount调整之后积压的分片任务要观察一段时间再下线否则会出现“老任务还在用10个分片跑新任务却已经按12个分片跑”的数据错乱。4.3 任务下发和执行日志的核对方法任务定义注册好不代表调度成功。我建议接入后第一件事是盯紧task_log表确认agent真的认领到了任务。ax的日志链路是execution_id贯穿始终调度器生成execution_id - agent收到任务时记录RECEIVED - 业务方法执行结束记录SUCCESS/FAILED - 调度器更新task_instance状态。如果发现task_instance状态一直停留在SCHEDULED说明调度器没有真正把任务下发出去。这时候优先检查agent的长轮询请求是否正常通常是因为agent机器无法反向访问scheduler的端口多见于安全组配置问题。如果日志里看到FAILED再看是TIMEOUT还是EXECUTION_ERROR前者说明任务本身耗时太猛或者线程卡死后者说明业务代码抛异常两个排查方向截然不同。为了降低排查成本我们后来给task_log表加了trace_id字段并把worker的stdout和stderr都定向到本地文件文件名带上execution_id。这样拿到一个失败任务直接按execution_id搜日志就能把链路串起来。4.4 监控指标别等业务方来找你调度系统最忌讳“无声失败”。ax从第一版就埋了Prometheus指标重点看四组数据调度器层面的scheduled_total和delayed_total反映时间轮是否堆积下发层面的dispatch_total和dispatch_failed_total反映worker和调度器之间网络是否健康执行层面的execute_success_total、execute_timeout_total、execute_fail_total反映业务任务自身质量等待层面的queue_depth反映流量调度模块是否过度拦截。告警规则是我的血泪教训换来的。一开始只设置了失败数大于0就告警结果高峰期每个小时告警几百条人直接麻了。后来改成按任务分组统计单个任务10分钟内连续失败3次且失败率超过50%才告警同时带上owner字段自动往对应企业微信群里发消息。这样既不漏报也不轰炸。5. 常见问题排查与避坑心得5.1 任务不触发八成是时区和cron表达式的锅“任务怎么没跑”是最常见的工单。排查了很久发现大多数情况不是调度系统坏了而是业务方不熟悉cron表达式细节。alarm灵的案例比较典型业务方写了个“0 0 10 * * ?”以为每天10点执行结果服务器时间UTC比北京时间晚了8小时任务每天凌晨2点就跑了数据当然不对。规范流程是所有cron表达式强制要求带timezone字段默认Asia/Shanghai但允许业务方显式指定调度器会把表达式转化到UTC后再写入时间轮。另外cron表达式里秒和年这两个字段经常被漏掉Ax的校验逻辑会把缺少秒字段的表达式直接拒绝虽然最开始被吐槽不友好但确实避免了大量线上乌龙。5.2 任务重复执行execution_id是唯一的解药有一段时间一个支付回调任务反复出现重复执行数据库里出现了两条一模一样的通知记录。定位发现问题出在agent处理完业务后回传结果时网络超时调度器判定任务失败并触发重试但第一个线程其实还在做收尾工作最终产生了双写。这个问题的根源是at-least-once语义下网络超时和真正的执行失败无法区分。解决方式就是我前面说的两层幂等调度器侧用Redis的SETNX锁住execution_idworker侧在业务入口查询幂等表。这里需要特别提醒幂等表插入动作必须和业务更新放在同一个本地事务里否则查完幂等表之后的事务还没提交重试线程又来了照样会重复执行。5.3 调度延迟抖动优先怀疑GC和连接池有段时间用户反馈一个秒级任务经常延迟三五秒才出数。查调度器日志发现时间轮指针走得正常但任务下发和agent回执之间的耗时出现明显尖刺这基本是JVM Full GC的信号。给scheduler堆内存调大同时把CMS换成G1并设置目标停顿时间尖刺立刻下降。另一个隐蔽点是MySQL连接池。调度器每次状态更新都要写库如果连接池最大连接数设得不够高峰期所有线程都在等连接直接拖垮调度延迟。建议连接池配置的初始大小覆盖正常峰值线程数最大大小放宽到两倍并且在监控里加上active_connection指标防患于未然。5.4 灰度发布时的任务大面积失败优雅下线救了一命Agent升级的时候如果直接kill进程正在执行的任务全部中断而且未回传结果的任务会被调度器判定失败自动重试。旧进程抢不到任务新进程又没有完全启动容易出现“任务的补偿风暴”。后来我们在agent里实现了优雅下线协议收到SIGTERM后停止拉取新任务但等待当前运行任务完成后再退出超时上限90秒。同时调度器侧增加watchdog连续三次长轮询没有响应的agent会被标记为离线任务自动漂移到其他节点。经过这次改造之后Agent发布对线上的影响基本可以忽略不计。5.5 经验速查表现象优先排查项常用解决办法任务完全不触发时区、cron表达式、任务是否设置disable统一timezone严格校验cron格式任务偶发延迟调度器GC、数据库连接池、Redis抖动优化GC参数扩大连接池增加监控任务重复执行网络超时后重试execution_id幂等本地事务保护幂等表任务堆积worker负载过高、优先级设置不当检查queue_depth拆分分片或调整优先级发布时失败剧增agent未优雅下线实现SIGTERM优雅退出90秒缓冲分片后数据错乱分片数变更后新旧任务混跑修改分片数时先观察积压再下线旧任务6. 从 ax 调度到统一调度平台整套系统跑了一年多之后组里对它的定位悄悄变了。以前大家只把它当定时任务工具后来接入了事件驱动型任务、CI/CD流水线的触发节点、甚至部分消息转发的路由规则。ax调度的核心价值已经不只是“到点执行”而是成为业务和底层资源之间的一道逻辑闸门。我在实际维护中有个很深的体会调度系统的复杂度不会因为业务简单而消失只会延迟爆发。与其等cron崩了再补窟窿不如一开始就把任务的模型定义清楚把超时、重试、幂等、分片、优先级这些机制做成公共能力。当然这套方案不一定适合所有团队如果你们只有几十个任务一个Quartz加完善的告警完全够用。但如果你也面临任务规模增长、执行环境动态变化的阶段希望这份复盘能帮你少走些弯路。对了最后再分享一个小技巧不要急着在项目初期实现可视化的任务编排界面先把任务的API和日志查询接好让业务方用起来等需求攒到一定量再考虑要不要上一个拖拽画布。调度系统最怕的一是过度设计二是没有反馈这俩你都躲开系统就不会难用。
返回列表