分布式任务调度系统实践:从架构设计到踩坑实录
1. 项目背景从定时任务散落一地说起手头这个项目代号叫 AX最初的需求听起来特别简单把团队里所有定时任务统一管起来。当时线上起码有三套跑任务的姿势——有一套挂在 Jenkins 上靠构建触发有一套直接写在每台服务器的 crontab 里还有一套是业务代码里用线程池加 timer 自己 sleep 出来的。每套各有各的问题Jenkins 那套一到发版高峰期就排队crontab 那套根本看不到失败率和执行日志sleep 那套更离谱进程只要一重启任务就全丢了还没人记得补。真正让我下决心做 AX 的是一次凌晨事故一个数据同步任务因为 crontab 配置失误加上网关重试在 8 台机器上几乎同时跑了一遍重复写入把业务表直接写花了。当时从日志倒推折腾了三个多小时才定位到根因。那次之后我意识到团队缺的不是定时任务工具而是一个统一的调度大脑——能看清楚任务什么时候触发、在哪台机器上执行、跑没跑成功、失败了该找谁。所以先明确一下本文里 AX 到底指什么它是一个面向批处理场景的分布式任务调度系统AX 就是这个项目的代号。所谓ax调度在我这个项目语境里指的是围绕任务注册、触发、分发、执行、重试、监控这一整套链路的设计与实现。如果你现在面对的也是任务多了管不过来、执行结果不可控、出了问题全靠猜这类状况这篇内容的参考价值会很大。下文我会把架构设计、核心实现和踩过的坑完整记录下来不会藏着掖着。2. 架构设计与关键取舍2.1 为什么没有直接选开源方案动手之前我其实认真评估过 Quartz、Elastic-Job 和 XXL-JOB 这几个主流方案最后还是没有直接用原因倒不是它们不好而是我们的场景有几个硬约束。第一任务来源不只是定时触发还有上游事件触发和上一任务完成后的链式触发。Quartz 这类框架的核心是 cron 调度对事件驱动和任务编排的支持需要二次开发改起来成本不低。第二我们的任务数量虽然不大但单个任务执行时间跨度很长最长的报表任务能跑 40 分钟这要求调度器不能依赖单机内存里维护一个 tick这种简单模型节点一挂状态就没了。第三也是最重要的一点团队希望把执行链路的数据全部统一收口——日志、指标、告警都要在一套体系里开源方案往往要把监控部分外包给 Prometheus 自研代码等于还是得写不少东西。所以最后决定自研但架构上明确了一个原则调度器本身要尽量轻复杂逻辑下沉到执行器和存储层。调度器只负责三件事——算时间、发指令、收状态。业务逻辑和任务依赖一律放执行端。这样后期不管是接新的任务类型还是把执行节点横向扩容都不需要动调度核心。2.2 核心模块拆分AX 整体上拆成了六个模块各管一摊模块职责关键点任务注册中心维护任务元数据、版本、启停状态用数据库行级锁保证并发注册安全触发引擎计算下次触发时间生成调度事件支持 cron、延迟、链式、事件四类触发分发器把任务实例分配给某个执行节点采用一致性哈希 租约机制执行器真正跑业务代码上报进度和结果嵌入业务服务通过 SDK 接入存储层保存任务定义、实例、日志、告警规则MySQL Redis 配合使用监控告警采集指标、检测异常、推送通知独立进程不依赖执行器存活每个模块都设计成无状态服务可以独立部署多副本。调度器集群内部通过选主的方式保证同时只有一个节点在写触发事件避免两个节点同时把同一个任务推出去。这个选主我一开始想用 ZooKeeper后来评估了一下发现用 MySQL 的GET_LOCK配合心跳续期就够了省掉一个中间件运维负担也小。2.3 分布式一致性的处理方式调度系统最难的地方不在定时触发这个动作本身而在分布式环境下的状态一致性。比如分发器把任务发给执行器之后执行器挂了怎么办执行器执行完了但回调结果因为网络抖动没送达怎么办这些都是必须想清楚的问题。AX 的方案是事件驱动 状态机校验。任务实例的状态只在 DB 里流转CREATED → DISPATCHED → RUNNING → SUCCEEDED/FAILED/RETRYING。所有节点都通过读取数据库状态来判断下一步动作而不是依赖内存里的印象。分发器收到触发事件后不是直接把任务推给某个执行器而是先在 DB 里插入一条CREATED状态的任务实例然后由执行器轮询待领取的任务。执行器领到任务后把状态改成RUNNING同时记录自己的worker_id。其他节点如果发现某个RUNNING状态的任务对应的工作节点心跳超时了就自动把它标记为RETRYING并重新进入分发池。说白了就是状态以数据库为准任何节点的本地状态都只算缓存。这样做牺牲了一点实时性轮询间隔最短 200ms但换来了极佳的容错性。注意调度系统的状态机设计一定要避免分布式锁 内存状态的思路。我见过不少团队用 Redis 分布式锁包住整个调度流程锁一过期或者 Redis 抖动任务就重复派发了。状态机 数据库校验配合租约lease机制才是比较稳妥的路线。2.4 技术选型MySQL Redis 的组合存储这层最终定了 MySQL 8.0 存任务实例和日志索引Redis 主要做三件事分布式锁、延迟队列、执行进度缓存。MySQL 存任务实例的时候核心表做了分区按create_time按月分区历史数据定期归档到冷表。这主要是因为任务实例在高峰期一天能产生十几万条记录如果不分区单表索引很快会膨胀状态更新的UPDATE语句性能会急剧下降。Redis 里的延迟队列是整个触发引擎的命脉。我用了 Redis 5.0 的 Stream 结构配合消费者组的XREADGROUP命令做延迟消费。每个触发事件都有execute_at时间戳消费者判断当前时间如果没到execute_at就把事件重新投递回去。虽然这样会有空轮询的 CPU 损耗但因为 AX 的任务量级在万级以内实测下来完全可接受。3. 核心实现细节3.1 任务定义的模型设计任务定义我用一个 JSON Schema 来约束这样后端灵活前端也能基于 Schema 自动生成表单。一个典型任务的模型长这样{ task_name: order_sync, task_type: batch, owner: data_team, triggers: [ { type: cron, cron: 0 2 * * *, timezone: Asia/Shanghai } ], execution: { handler: com.ax.handler.OrderSyncHandler, timeout: 1800, retry: { max_attempts: 3, backoff: exp, initial_interval_sec: 30 } }, concurrency: { max_running_instances: 1, blocking_strategy: serial }, dispatch: { strategy: hash, tags: [data-node-a, data-node-b] } }字段不多但每一个背后都有讲究。max_running_instances控制同一任务同时运行的实例数比如订单同步这种任务绝对不允许两个实例并行跑这个字段就设成 1blocking_strategy是serial意思是上一个没跑完下一个触发就排队等着而不是直接丢弃。dispatch.tags是执行节点的分组标签。分发任务的时候调度器只会在带对应标签的执行节点里选人。这样数据同步任务可以指定跑在高配机器上而轻量的清理任务就交给普通节点成本和效率都能照顾到。3.2 触发引擎不止是 cron触发引擎初期我差点做成了cron 解析器后来仔细想了下场景发现真正需要的其实是四类触发源cron 定时触发固定时间周期用于日报生成、数据同步、日志清理这类日常任务。延迟触发事件发生后的 N 秒再执行常用于业务补偿场景比如订单支付超时关单。链式触发上一个任务成功后自动触发下一个用于报表链路这种有依赖关系的场景。事件触发外部通过 HTTP API 主动提交触发请求用于和业务系统联动。四类触发源最后都统一转成一个内部事件TaskTriggerEvent包含任务名、触发类型、计划执行时间、业务参数。这个事件进入 Redis Stream 后由触发引擎消费。cron 触发的实时性要求最高所以我在触发引擎里做了一个小优化用Ticker每分钟扫描一次最近一分钟内到期的 cron 任务精确到秒的触发由下一层延迟队列校准避免让所有 cron 任务都挤在同一秒创建实例。事件触发这里有一个细节值得提醒外部 API 提交过来的业务参数一定要做版本校验。我们的TaskTriggerEvent里带了payload_version字段执行器拿到参数后会和当前任务的schema_version比对不一致就直接拒绝执行。这个是因为线上出现过上游系统改动参数格式后下游任务用旧代码解析新字段导致批量任务全部失败的事故。3.3 完整执行链路一个任务从触发到执行完走的是下面这条链路触发引擎消费到事件校验任务状态为ENABLED在 MySQL 创建任务实例状态CREATED。分发器扫描CREATED状态的实例按任务的dispatch.strategy选一个执行节点把实例标记为DISPATCHED并把任务详情写入 Redis 的待执行列表。执行器上的 Worker 每 200ms 轮询一次待执行列表拉到任务后先向调度器注册心跳再把状态改成RUNNING开始执行业务逻辑。执行完成后Worker 异步上报结果。成功则直接把状态改成SUCCEEDED失败则判断是否还有重试次数有就回退到RETRYING按退避策略重新入队。调度器的监控模块兜底如果发现某个RUNNING任务超过timeout没上报心跳就强制标记为FAILED并按策略触发告警。Worker 注册心跳这个动作本质上是租约续期。租约默认 30 秒任务每 10 秒续一次。如果租约过期调度器就认为该执行节点已经失联会把任务重新调度到别的节点。这里有个容易踩坑的地方租约时间要比心跳间隔大得多否则网络一抖动就误判节点死亡导致任务被反复调度。真实经验是心跳 10 秒、租约 30 秒这个比例比较稳。3.4 重试机制与退避策略任务失败的自动重试是调度系统里最容易被低估的环节。我见过不少系统是失败了就立即重试结果就是下游服务还没恢复重试请求一来直接把系统打挂。AX 的重试策略是可控可配的支持固定间隔、线性递增、指数退避三种还允许设置最大重试次数。public class RetryPolicy { private int maxAttempts; private BackoffStrategy backoff; private int initialIntervalSec; private int maxIntervalSec; public long nextIntervalMs(int attempt) { return switch (backoff) { case FIXED - TimeUnit.SECONDS.toMillis(initialIntervalSec); case LINEAR - TimeUnit.SECONDS.toMillis(initialIntervalSec * attempt); case EXP - Math.min( TimeUnit.SECONDS.toMillis(initialIntervalSec) * (long) Math.pow(2, attempt), TimeUnit.SECONDS.toMillis(maxIntervalSec) ); }; } }指数退避里maxIntervalSec的封顶值非常重要。如果没有上限第 10 次重试的等待时间会膨胀到几个小时任务积压得越来越严重。我这边默认初始间隔 30 秒封顶 30 分钟。还有一点重试不能只靠调度器这一层。真正可靠的方案是调度器负责调度层面的重试执行器负责业务层面的幂等。什么意思呢调度器只保证任务会再次被派发出去而业务代码必须保证同一个任务实例即使被执行两次结果也是一致的。这个就引出下一个话题并发控制和幂等设计。4. 调度策略与生产级细节4.1 并发控制与幂等设计并发控制是 AX 里我花时间最多的部分。常见方案是分布式锁但分布式锁在调度场景下有个致命弱点锁过期时间不好定。任务跑 10 分钟锁设 30 秒肯定不行任务跑 10 分钟锁设 15 分钟又会导致故障恢复时间太长。所以 AX 里的并发控制分了两层。第一层是任务级别的max_running_instances限制。当一个新的触发事件到来时调度器会先查询当前该任务的RUNNING/DISPATCHED状态实例数如果已经达到上限就把新实例放进等待队列。这层的实现依赖 MySQL 的SELECT ... FOR UPDATE对任务行加锁确保判断和插入是原子的。第二层是执行节点上的任务级别互斥。每个执行器启动时维护一个ConcurrentTaskGuard基于 Redis 的SET NX EX命令写一个带任务 ID 的互斥键。只有抢到互斥键的 Worker 才会真正执行任务。这个互斥键的过期时间设成任务超时时间的两倍防止任务卡死导致互斥键提前过期。幂等设计的核心是业务方自己控制。AX 提供了两个辅助能力一是任务实例 ID 会透传给业务代码业务方可以用它在本地做去重二是提供了一个工具类支持基于 Redis 或数据库的唯一键来执行幂等逻辑。我强烈建议在数据同步、消息推送这类任务里业务方务必记录上次处理到的位点比如订单号、时间戳而不是依赖调度系统的重试机制来保证不重复。4.2 优先级与队列模型共享集群里跑任务最怕的就是低优先级的任务把高优先级任务的资源抢光了。AX 的调度队列做成了多级队列紧急队列、普通队列、低优队列。每个队列里任务按priority字段排序分发器优先从紧急队列取任务只有紧急队列空了才会处理普通队列。队列的存储我用的是 Redis ZSET有序集合分数就是任务的到期时间戳加优先级权重。这样既能保证按时序消费又能让高优任务插队。不过这里要注意 ZSET 的单个 key 容量问题元素太多时ZRANGEBYSCORE的耗时会上涨。我的做法是队列按小时分片比如ax:queue:urgent:20250617消费端按当前时间选择对应的分片。这种多级队列虽然能缓解优先级问题但治标不治本。真正要避免资源争抢还是得靠资源隔离——也就是 4.3 要说的。4.3 资源隔离与限流执行器侧的资源隔离我采用了每类任务独立线程池 信号量限流的组合。每个任务注册时声明自己的线程池大小比如报表任务池 4 个线程数据同步任务池 8 个线程。这样即使某个任务异常膨胀最多只能占满自己那个池子不会拖垮执行器上的其他任务。ExecutorService reportPool new ThreadPoolExecutor( 2, 4, 60L, TimeUnit.SECONDS, new ArrayBlockingQueue(200), new NamedThreadFactory(ax-report), new ThreadPoolExecutor.CallerRunsPolicy() );线程池的拒绝策略这里我踩过坑。一开始用的AbortPolicy任务一多直接抛异常导致调用链路上游也跟着报错。后来换成CallerRunsPolicy意思是线程池满了就让提交任务的线程自己去跑这样至少任务不会丢只是执行时间会被拉长。对于批处理任务来说延迟比丢失更可接受。信号量限流用在外部依赖不可控的场景。比如一个任务要调用第三方 APIQPS 上限是每秒 50 次那就给这个任务的执行逻辑外面包一层计数器超过就阻塞等待。这个限流一定要在任务内部做调度器层面做不了因为调度器不知道任务内部调用了几个外部接口。4.4 负载均衡与节点上报执行节点接入 AX 时会向调度器上报自身信息和心跳包括节点 ID、所在分组、CPU 核数、当前运行任务数、内存使用率。分发器在选择目标节点时不是简单的随机选而是按以下规则过滤掉心跳超时的节点过滤掉当前任务数超过max_running_tasks的节点有tags限制的任务只保留标签匹配的节点在剩余节点里优先选择当前运行任务数最少的节点最少连接算法。这种按需调度的方式比轮询或者随机要好很多因为它天然考虑了节点的负载差异。我之前试过一致性哈希来保证同一任务总是落在同一节点上后来发现对于多数批处理任务来说这种粘性根本没有必要反而会增加热点。最终只有任务指定了dispatch.strategy hash时才启用哈希定位。5. 常见问题与排查实录5.1 任务重复执行的根因与解法重复执行是调度系统最常见的问题AX 上线头两个月就遇到两次。第一次是数据库主从切换导致的调度器写状态时连接到了主库但读状态的时候走了从库从库延迟导致读到旧状态于是同一任务被派发了两次。解决方法是把任务实例的读写都强制走主库并且对关键操作的超时时间做了放大。第二次是任务本身的超时设置太短。一个同步任务实际要跑 25 分钟但timeout只配了 10 分钟。到了 10 分钟监控模块判定任务超时强制标记失败并重新调度但旧节点上的线程还在跑。于是新旧两个节点同时执行了同一任务。这个问题的教训是超时时间一定要结合历史执行数据来定宁大勿小同时要在业务代码里利用任务实例 ID 做二次幂等兜底。5.2 任务消失了一次延迟队列的数据丢失有个用户反馈每天凌晨的报表任务偶尔会消失。查了一圈发现是 Redis Stream 的消费者组在处理消息时ACK和业务操作不是原子的。业务操作先做了但ACK还没发出去这时候消费者崩溃消息就会被重新投递任务重复执行。或者是反过来消息先从 Stream 里POP出来但还没处理进程就崩了这个消息就永久丢失了。后来我把消费流程改成了先XCLAIM设置处理中标记再处理业务最后XACK。Stream 的XCLAIM能把消息的状态从 pending 改成正在处理配合XINFO GROUPS查看 pending 队列长度基本可以做到不丢不重。现象根因解决方案任务重复执行主从延迟导致读写状态不一致关键状态强制走主库任务重复执行超时误判导致新旧节点并行放宽超时 实例 ID 幂等任务丢失Stream 消费未 ACK 即崩溃改用 XCLAIM 处理中标记任务不触发cron 表达式时区配置错误统一使用 UTC 显式声明时区任务卡在 RUNNING租约过期但执行器存活心跳间隔与租约时间按 1:3 配置5.3 长时间任务拖垮线程池报表任务偶尔会因为上游数据库响应慢单个查询卡住 10 多分钟结果线程池被占满后面排队的任务全部堆积。这个问题的排查过程很有代表性先看线程池的活跃线程数全是RUNNING再看堆栈全部阻塞在 JDBC 的socketRead0上。显然不是代码死锁是数据库端资源问题。AX 的应对方案是三管齐下第一为数据库中慢查询配置了连接级别的超时第二每个任务的线程池都配上队列上限满了就拒绝而不是无限堆积第三引入任务级看门狗单个任务执行超过设定阈值时先尝试中断线程中断失败再强制标记失败。线程中断这个事实际执行时要注意Java 的Thread.interrupt()对阻塞在 IO 上的线程不一定有效。所以最可靠的手段还是让业务方在代码里检查Thread.currentThread().isInterrupted()主动感知中断信号。AX 的 SDK 里提供了一个TaskContext业务代码可以随时用它判断任务是否已经被外部标记为超时。5.4 监控与告警哪些指标必须盯调度系统的监控可比普通业务系统难做因为它的状态是分层的。我最后沉淀下来的核心指标有八个调度延迟从事件产生到任务实例创建的时间差正常情况下应该低于 1 秒。分发成功率分发器成功把任务发给执行器的比例低于 99% 就要查网络或节点注册问题。任务超时率单个任务超时实例数 / 总实例数这个指标最能反映任务本身的质量。排队积压数等待队列里超过 5 分钟还没被消费的任务数量。节点心跳丢失率执行节点心跳超时的频率太高说明节点负载有问题。重试率任务首次失败后重试的比例高于 20% 就要关注任务稳定性。租约过期次数误判节点死亡的次数这个指标能暴露租约配置不合理的问题。Worker 线程池使用率执行器上线程池的繁忙程度用于提前发现资源瓶颈。告警方面AX 接入了企业微信机器人按告警等级分了三档P0 是任务大面积失败或节点全部失联电话通知P1 是单任务连续失败 3 次或超时率升高群消息加 负责人P2 是排队积压或指标抖动只发群消息记录。告警文案里一定要带上任务 ID、实例 ID 和日志查询链接避免告警到了人手上还要再去翻半天日志。5.5 任务日志如何统一收口任务跑在哪台机器上日志就散落在哪台机器上这个问题不解决排查效率永远上不去。AX 的做法是执行器不直接写本地日志文件而是通过异步 logger 把结构化日志推送到统一的日志中心。日志字段固定为task_id、instance_id、worker_id、level、message、timestamp。日志推送一定要异步化并且加了本地缓冲队列。刚开始我图省事直接同步推送结果业务代码里打个日志都要等网络 IO任务执行时间直接翻倍。后来改成内存队列 批量发送日志对业务代码的性能影响基本可以忽略。日志中心按天建索引支持按实例 ID 一键查询全链路日志这个体验提升非常大。6. 说说那些从坑里爬出来的经验AX 从第一个版本上线到现在前前后后改了十几轮有几个经验我觉得比代码本身更值钱。第一调度系统的状态流转一定要做审计。任务实例每一次状态变更都记录from_state、to_state、operator、reason四个字段。这个审计日志平时不起眼但一旦出问题它是定位的唯一线索。我们有好几次线上事故都是靠审计日志逆推出问题节点的。第二永远要为调度器和执行器版本不一致留后路。业务迭代快执行器升级跟不上调度器是常有的事。AX 的任务定义里有min_worker_version字段调度器只会把任务派发给版本号达标的节点这样即使集群里存在老版本节点也不会出现老节点执行新逻辑的兼容性问题。第三也是最重要的体会调度系统本质上不是技术系统而是流程系统。技术方案再完善如果任务没有 owner、失败没有 follow-up 机制、每周没有例行巡检系统一样会烂掉。AX 上线后我强制要求每个任务都有负责人并且每个月跑一次任务健康度报表把连续失败、长期未触发、执行时间异常的任务都列出来逐个确认。最后分享一个小技巧任务上线前用 AX 的演练模式把任务跑一遍这个模式会模拟节点宕机、数据库抖动、重试风暴三种故障场景可以提前暴露大部分潜在问题。我们后面新任务接入都强制走演练效果非常明显。做调度系统永远要假设明天就会出故障然后把应对故障的能力提前焊死在系统里。
上一篇/下一篇内容由系统自动关联
返回资讯列表 →