尧图网络科技YAOTU DIGITAL 获取报价
获取报价
首页 / 资讯中心 / 文章详情

自己动手写一个分布式任务调度器ax:设计、架构与实战避坑指南

发布时间:2026/9/28 17:02:52

资讯中心
01
ARTICLE

自己动手写一个分布式任务调度器ax:设计、架构与实战避坑指南

自己动手写一个分布式任务调度器ax:设计、架构与实战避坑指南
1. 为什么我自己动手写了一个叫 ax 的调度器做数据平台的时间长了你迟早会遇到这样一个问题凌晨两点二十几个批任务同时到点结果一半在抢资源另一半在等前一天的依赖数据调度器自己先卡死了。市面上那几套开源调度框架我都试过要么太重要么调度模型和我们的业务对不上要么扩到几百台 Worker 之后心跳和锁的维护成本高得离谱。折腾了一年多我最后决定自己写一套轻量的分布式任务调度器内部代号就叫 ax核心职责就是两个字调度。我说的“ax调度”不是简单地拿 Cron 表达式把任务丢出去执行而是把任务的优先级、依赖关系、资源配额、失败重试、抢占策略全部纳入一套统一的调度模型里让每一份计算资源在任意时刻都知道自己该跑什么。开发这套系统的一线经验今天整理出来希望能给同样被任务编排折腾得够呛的人一些参考。适合谁阅读呢主要是做数据平台、后端基础架构、运维体系的人也适合那些业务系统里有一堆定时任务、异步任务想让调度更可控的开发团队。当时推动我写 ax 的导火索是某个大促的压测日。业务方一次性提交了三百多个分析任务原本用 XXL-Job 改的模型在高峰期直接超时死锁告警刷屏。那一次事故让我想明白一件事很多调度系统做得不好不是功能少而是把“什么时候跑”和“跑在哪里”“跑得是否合理”混在了一起。ax 从一开始就决定把这些问题拆开。2. ax 的核心架构与调度模型2.1 控制面与执行面彻底分离ax 的第一个设计决策是把调度器拆成两个独立的进程组ax-master 和 ax-worker。ax-master 只负责决策不碰任何业务数据ax-worker 只负责执行不维护任何全局状态。这样的好处很直接调度器挂了已经在跑的任务不会中断只是暂时无法接收新任务Worker 挂了master 能在几秒内感知并重新调度它的任务到其他节点。Master 内部又分了两层一层是 API 层负责接收任务定义、手动触发、查看状态另一层是调度核心维护着全局任务队列、依赖图、资源池状态。调度核心的数据结构必须做到无锁或者细粒度锁因为一旦调度决策本身出现竞争整个集群的吞吐量就会断崖式下跌。Worker 启动时会向 Master 注册上报自己的 IP、可用内存、CPU 核数、已分配的槽位数量。Master 给每个 Worker 维护一个“可调度容量”的计数只有计数大于 0 时才允许分配任务。这个计数不是简单的 CPU 核数而是根据每个任务声明的资源需求动态计算的。注意控制面与执行面分离的理念很多人容易理解成只是架构上的优雅实际它的价值在于故障隔离。我见过不止一个项目把调度和执行放在同一个进程里结果一个任务里用了 System.exit() 直接带走整个调度器这种事故一次就能让你下定决心拆开。2.2 任务、实例、执行单元的三层抽象ax 里最核心的概念不是任务而是实例。任务Task是静态定义描述一段逻辑要跑什么、依赖哪些上游、多久跑一次、需要多少资源实例Instance是任务在某一次调度中的动态产物比如 2024-06-01 00:00:00 这个时间窗口的实例。每个实例有一个全局唯一 IDID 的生成规则包含了任务 ID 和时间窗口这样排查问题时可以一眼看出是哪个周期。再往下拆一层每个实例真正运行起来之后会对应一个或多个执行单元Execution。为什么还要再拆一层因为同一个任务在重试时会产生多个执行单元。如果你把“实例”和“一次运行”混在一起重试历史就没法完整记录了。ax 中实例是稳定的执行单元是可变的状态都挂在执行单元上。这三层抽象让我在后面的权限控制、审计、计费中省了大量力气。比如业务方问“为什么昨天的任务跑了三次”我可以直接查实例下挂的执行单元看到第一次因数据源超时失败第二次因输出目录被占用失败第三次成功。如果没有这层抽象你只能从日志里猜。2.3 调度队列的资源配额模型ax 的资源配额仿照 YARN 的 Capacity Scheduler 做了简化版。每个业务线有一个队列队列配置了最大资源百分比和最低保障。调度时Master 会先按最低保障分配资源有剩余资源时再按权重分给高优先级队列。但跨队列抢占默认是关闭的除非你明确开启“弹性抢占”开关。这个模型在实现时最关键的地方是不能真的给 Worker 上的任务做资源隔离那样成本太高只需要在调度决策时做好资源记账。也就是说ax 是“调度层资源管控”不是“运行时资源隔离”。如果一个任务实际使用的内存超出了声明ax 不会主动 kill 它但会在监控上标记为“超用”并在下一次调度时降低该任务的可信度。我当时纠结了很久要不要引入容器来做严格隔离。后来评估了团队运维成本决定还是在无容器模式下先跑通。如果你们已经有成熟的 K8s 环境可以考虑把 ax 的 Worker 换成 Pod但调度模型保持不变。资源记账的核心代码其实就一段逻辑每个队列维护一个计数器任务分配时扣减任务结束归还超卖比例允许超过 1.0 但默认不超过 1.2。3. ax 调度策略的落地细节3.1 优先级不是数字而是一套规则很多调度系统把优先级设计成一个简单整数1 到 10越大越先跑。ax 一开始也这么干但很快就发现这个设计太粗糙了。业务方为了让自己任务先跑几乎清一色把优先级写成 10最后大家都等于没优先级。后来我把优先级拆成三个维度等级Level、提交时间SubmitTime、等待时长WaitTime三者加权综合。等级是业务方自定义的分成 0 到 5 六个档位。提交时间用于同等级下的先进先出保证公平性。等待时长是一个兜底策略如果一个任务等了超过 30 分钟权重会随着等待时间线性增长防止低等级任务被高等级任务活活饿死。这套规则看起来复杂实际计算起来就一行公式score level * 1000 (now - submitTime) / 60 max(0, waitTime - 1800) / 60。实现优先级调度时很多人会直接用优先队列PriorityQueue但 Java 的 PriorityQueue 是线程不安全的而且当你需要动态更新某个任务的 score 时得先 remove 再 add这个操作是 O(n)。ax 里我改用跳表ConcurrentSkipListMap按 score 排序key 是 score 加唯一 ID 的组合这样既支持并发读取又能实现 O(log n) 的动态调整。实测下来三千个待调度任务做一次全量重排耗时在两毫秒左右。3.2 依赖编排从 DAG 到状态机ax 的依赖编排没有用现成的工作流引擎而是自己实现了一个轻量级的 DAG 状态机。每个任务定义里包含 dependsOn 列表里面是上游任务 ID 和完成状态成功/忽略失败。当上游任务产生成功实例后Master 会向该任务的下游发送一个“依赖满足”的事件。下游任务收到所有依赖事件后才从“等待依赖”状态变为“可调度”状态。这个状态机一共有六个状态INIT、WAITING、RUNNABLE、RUNNING、SUCCESS、FAILED。WAITING 表示依赖未满足RUNNABLE 表示可以被调度器选中。开发时最需要注意的是状态流转的原子性。多线程同时收到多个下游事件时不能出现重复流转。ax 用数据库行锁来解决在 MySQL 里的 task_instance 表加一个 version 字段每次状态变更用 UPDATE ... WHERE version ? 的乐观锁失败则重试重试三次后放弃并告警。还有一个细节是“跳过依赖”功能。有时候上游任务的数据源出问题了但下游任务可以先跑空跑逻辑等数据补上再重算。ax 支持在下游任务实例上手动标记“忽略上游失败”这个操作不是简单地改状态而是会在实例上生成一条 audit log记录是谁在什么时间跳过的方便事后追责。实操心得依赖编排最容易踩坑的是循环依赖检测。你的 DAG 不能出现环否则两个任务会永远互相等待。ax 在保存任务定义时就会做一次环检测用的算法是拓扑排序的逆思维检测是否存在无法处理的节点。这个检测放在 API 层而不是调度层这样能把错误提前暴露在配置阶段而不是运行阶段。3.3 超时控制不要相信任何任务的承诺在 ax 之前我见过太多任务因为 Runnable 里的一段代码阻塞导致整个进程的线程池被打满。ax 对超时的处理分了两层调度等待超时和执行超时。调度等待超时是指一个实例从 RUNNABLE 变为 RUNNING 的等待时长默认是 30 分钟超过后触发告警并可选地自动降级为“跳过”。执行超时必须由提交任务的代码主动配合。ax 提供了统一的 TaskExecutor 包装类内部使用 Future 的 get(timeout) 机制。业务方只需要实现一个返回结果的方法ax 会自动加上超时中断逻辑。这里有个坑Future.get 超时后任务线程并不会停止它只是不再等结果了。所以 TaskExecutor 里必须在超时后强制中断线程并在线程结束前释放资源锁。我最初写的版本超时后直接丢弃 Future后来发现任务还在 worker 上继续跑占着资源不说还会重复写输出文件。改成强制中断后又遇到了中断不响应的问题——业务代码里对 InterruptedException 处理不当。所以后来我在 TaskExecutor 的文档里明确写了一条你的任务代码必须正确处理线程中断否则调度器有权强制 kill 进程。4. 从部署到跑通第一个任务ax 实操记录4.1 最小化部署拓扑与参数规划如果你只是想体验 ax三台机器就够了一台跑 MySQL存储任务定义和实例元数据一台跑 ax-master两台跑 ax-worker。Master 默认端口是 8400Worker 通过配置里的 masterAddr 来注册。部署时我先强烈建议打开一个配置项task.auto.recover true表示 Master 重启后自动找回之前 RUNNING 的实例并把它们重新标记为 UNKNOWN等待 Worker 上报心跳后再更新为真实状态。JVM 参数方面Master 的堆内存不要超过 4GB因为 Master 不处理任务数据堆里主要放的是任务定义缓存和调度队列索引。Worker 的堆内存要根据你任务的实际负载来定我一般建议 8GB 起步因为 Worker 要加载业务代码和数据源连接。GC 算法统一推荐 G1并且打印 GC 日志方便排查任务卡顿是否和 Full GC 有关。数据库连接池的配置容易被忽略。Master 和 Worker 都会频繁访问数据库尤其是 Master 的调度循环每秒钟可能会查询多次任务状态。连接池最大连接数不要低于 20最小空闲连接数不要低于 5。我在第一次压测时就是因为连接池默认配置只有 10 个连接导致 Master 在任务高峰时大量等待获取连接间接拖慢了调度延迟。4.2 提交一个定时任务并观测调度日志跑通第一个任务的流程很简单。先在 ax-web 控制台创建一个队列比如 q_demo最大资源 100%最低保障 50%。然后创建一个任务调度周期写成 Cron 表达式 */5 * * * * ?表示每五秒触发一次。任务代码就写一行日志输出。保存后等第一个调度周期到来去 Worker 日志里会看到一条关键日志[ax-8400] task10001 instance20240601000500 dispatch to worker10.0.0.12。这一条日志包含了三个关键信息任务 ID、实例 ID、目标 Worker。从这条日志开始你就能追踪这个实例的完整生命周期了。实例 ID 的时间部分对应了 Cron 表达式的触发时间而不是实际执行时间这点要注意。比如任务阻塞了十分钟这次触发的实例 ID 还是 20240601000500表示它本该在零点零五分跑实际跑到十点一刻才开始。ax 的调度日志默认是 INFO 级别但我觉得生产环境建议把 ax-scheduler 这个 logger 调到 DEBUG。DEBUG 会记录每一次调度的打分过程包括 score 计算明细、候选 Worker 列表、被过滤的原因。虽然日志量会增大很多但在排查“为什么这个任务被调度到了那台机器”这类问题时DEBUG 日志几乎是唯一可靠的依据。注意不要在日志里记录整个任务请求体。我曾经在某个版本里把任务参数全量打到了日志结果用户密码直接出现在了日志文件里。后来我加了一层脱敏工具只记录参数的长度和哈希值排查问题时再手动关联。4.3 自定义调度插件用一套规则过滤目标 Workerax 的调度核心留了一个扩展点叫 WorkerSelector。默认根据资源余量选择得分最高的 Worker但你可以实现自定义的 SelectionStrategy比如“优先选择带有 GPU 标签的 Worker”“避开某个可用区的 Worker”。我把它设计成了 SPI 接口业务方只需实现一个类并打成 jar 丢进 ax-master 的 plugins 目录重启后即可生效。这个自定义插件我实际用过一次场景是某个任务需要访问某台内网数据库数据库只允许指定 IP 段的机器访问所以需要调度器把这类任务固定分发到某个子网内的 Worker。实现时我写了一个简单的过滤器如果任务定义里的 tags 包含 db:vip则只保留 IP 前缀是 10.0.16. 的 Worker。这个接口的签名很简单入参是候选 Worker 列表和待调度实例出参是排序后的列表。实现自定义 WorkerSelector 时要特别注意它是在 Master 的调度线程里被调用的所以不能做耗时操作更不能发 HTTP 请求。一次调度循环中多个排队中的任务会连续调用这个接口如果里面有个慢操作整个集群的调度延迟会被拖垮。正确做法是提前把规则计算好缓存起来插件里只做内存匹配。5. 常见问题排查与避坑指南5.1 任务堆积如何判断是调度慢还是执行慢ax 上线后最常见的告警是“任务队列积压超过阈值”。很多人第一时间去查 Worker 的资源是不是打满了但我建议先打开 ax-master 的 API/api/v1/queue/overview它返回每个队列的等待数量、平均等待时长、最近一次调度耗时。关键指标是“调度耗时”这个值应该稳定在 50ms 以内如果超过 200ms说明 Master 内部有了瓶颈而不是 Worker 忙不过来。调度耗时的瓶颈往往出在锁竞争或 SQL 查询上。我曾经排查过一个诡异现象任务一多Master 的 CPU 不高但调度延迟直线上升。后来用 async-profiler 一看发现是 Worker 心跳接口的写库操作和调度器的读库操作死锁了原因是一个长事务没提交导致行锁一直不释放。解决方法是把心跳写成批量更新每次最多 100 条并且把事务超时时间设置为 5 秒。如果调度耗时正常那就看“等待时长”有没有持续增长。如果等待时长在涨说明 Worker 资源真的不够了。此时不是盲目加机器而是看队列里的任务是不是有相互争抢的变态任务。比如一个任务声明了需要 64GB 内存但实际只用 1GB它会把 Worker 的“可调度容量”占掉一大半其他任务只能排队。这类问题需要靠 3.1 的资源可信度机制来治本。5.2 Worker 心跳丢失脑裂与僵尸任务ax 的 Worker 默认每 5 秒心跳一次Master 连续三次没收到心跳就把 Worker 标记为 OFFLINE并把它的 RUNNING 任务全部重置为 UNKNOWN等待下一次调度重新分配。这个逻辑看起来简单但真实环境中很容易出现“假离线”。比如 Worker 因为一次 Full GC 停下来 30 秒心跳自然发不出去Master 就会把那台机器上的任务重新调度到别的机器任务就会在两个机器上重复执行。解决脑裂的常见做法是引入“Lease”机制。ax 中每个任务实例在被 Worker 接收时会附带一个租约租约默认 10 分钟。Worker 需要在租约到期前向 Master 续租否则 Master 有权把任务重新分配。这样 Full GC 导致的短暂心跳丢失不会立即触发重新调度因为还有租约兜底。租约续租不是真的要调数据库而是 Worker 每 10 秒向 Master 发一个续租请求Master 只更新内存时间戳异步异步地批量写库。如果一旦确认任务真的在两个 Worker 上重复跑了最有效的兜底是任务的“幂等键”。ax 在任务执行前会检查执行日志表里是否存在相同 execution_id 的记录如果存在且状态不是 FAILED就拒绝执行。这个幂等检查虽然增加了一次数据库查询但避免了重复执行带来的数据错乱。踩坑记录我最初把租约时间设成 30 秒结果某次线上网络抖动所有 Worker 在五分钟内陆续离线又上线Master 重新调度了大量任务数据库连接池被瞬时打满。后来我把租约时间调整为任务预估耗时的两倍并且上线了“调度风暴”保护开关如果一分钟内重新调度的任务超过 50 个自动暂停重新调度等待 Worker 恢复上报后再继续。5.3 时间轮与延迟队列处理延迟任务的两种思路ax 里有一些需要延迟再处理的任务比如“失败后 10 分钟重试”。最初我使用了 RabbitMQ 的延迟队列功能每个失败任务扔到延迟队列到期后再塞回主队列。这个方案的问题在于数量大了之后MQ 本身成了瓶颈而且重试次数多了队列里的消息堆积会非常难看。后来我把延迟逻辑移到了 Master 内部实现了一个简单的时间轮。时间轮实现并不复杂我用的是一组环形数组每个槽代表一秒指针每秒跳一格。任务到期后从延迟队列里取出并提交到调度队列。但这个方案有一个隐患Master 重启后内存中的时间轮全部丢失。所以 ax 的延迟任务在提交时同时写入了 MySQL 的延迟表Master 启动时会扫描这张表把还未到期的任务重新装进时间轮。这样一来时间轮只负责到期提醒持久化由数据库兜底两者配合。如果你不想引入时间轮另一个可选方案是直接在数据库里扫“待重试”状态的任务。比如每 10 秒执行一次 SQLSELECT ... WHERE retry_time NOW() AND status RETRY_WAIT。但这样会让数据库变成调度器的轮询引擎任务量到几千之后这种扫表方案的延迟会变得不可控。时间轮方案的实际效果是在十万个延迟任务场景下调度误差控制在 200ms 以内。6. ax 调度从定时走向按需演进的一些思考6.1 从“每五分钟跑一次”到“数据到了才跑”ax 已经稳定运行半年后业务方开始提出一个新的诉求很多任务是上一个任务跑完才能触发但触发时间不固定有的凌晨跑有的下午跑完全取决于上游数据到达时间。如果继续用 Cron 表达式就得设置一个很密的频率然后不断轮询上游数据是否就绪这样既浪费资源又不实时。所以我们需要给 ax 增加一种“事件驱动”的触发方式替代传统 Cron。事件驱动的实现思路其实不复杂ax 在 Master 里增加一个 EventListener 组件业务方通过 HTTP API 向 ax 推送事件比如POST /api/v1/tasks/{taskId}/events事件内容是 JSON。如果任务定义里的 triggerType 是 EVENT调度器就不会按 Cron 来计算而是等待事件到达后立刻生成一个实例。事件到达的一刻调度器的逻辑和定时触发完全一样后面的优先级、依赖、资源分配都是复用已有的模型。我在做这个功能时最大的体会是调度系统的核心价值不是“定时”而是“把正确的任务在正确的时间放到正确的资源上”。事件驱动只是触发来源变了调度内核完全不用改动。所以如果你也在设计类似的调度系统一定要把“触发”和“调度”两个概念解耦否则每接一个触发源就要动一遍核心代码。6.2 可观测性不要让调度器成为黑盒ax 早期版本的监控只有一个“调度失败数”指标排查问题时非常被动。后来我补了三个关键指标调度延迟从任务到达可调度队列到分配成功的耗时、队列积压数、任务实例状态分布。这三个指标分别从两个视角看一个是调度系统自身的健康度一个是业务任务的健康度。调度延迟会告诉你调度器是不是快撑不住了队列积压数和状态分布会告诉你哪些业务线在争抢资源。日志方面ax 为每个实例生成一条统一的 traceId格式是{instanceId}#{executionId}。所有与该实例相关的日志无论产生在 Master 还是 Worker都会带上这个 traceId。这样通过日志平台检索就能把一个实例从提交到执行到结束的完整链路拉出来。这一点在分布式调试中帮了大忙我甚至可以不用看代码直接看日志就能定位问题出在哪个环节。要想更进一步可以把 ax 的指标接入 Prometheus。ax 在 Master 的/metrics端口暴露了标准的 Prometheus 格式我配置了几条常用的告警规则比如“队列积压数连续 5 分钟大于 500”“调度延迟 P99 大于 200ms”“Worker 离线数量大于 3”都能在 Grafana 上看到实时告警。现在 ax 上线这么久我基本不需要再盯着业务方的工单来排查问题因为告警已经能提前暴露风险。7. 一些实践后的个人感受这套 ax 调度系统从写到上线前后大概用了一个季度的时间。我个人的感受是调度器的难点不在并发和分布式而在模型设计。只要把“任务”“实例”“执行单元”的关系理清楚把“触发”“调度”“执行”三个环节彻底分开后面的一切都会变得顺理成章。很多开源框架喜欢把能想到的功能都塞给你但真正贴合自己业务的那 20% 才是值得花精力打磨的地方。如果你也想做类似的调度器我给一个最朴实的建议先别急着写代码先用一张白纸画出你现在所有任务的生命周期然后标出哪些环节是卡住的哪些环节是你每天靠人工去干预的。那些需要人工干预的地方就是你调度系统应该优先解决的问题。我当初就是因为每天凌晨都要手动重启一堆失败任务才下决心写 ax。现在这些事ax 自己就能做完了。
02
RELATED NEWS

相关资讯

更多网站建设与数字化升级内容

03
WHY YAOTU

想打造同款高转化官网?

懂行业、懂生意,从建站到增长一站式陪跑

◈

场景化定制

不做模板站,围绕你的业务场景量身设计,小众不撞款。

◐

营销型架构

以转化目标组织内容与路径,让官网真正带来询盘。

▲

全周期服务

设计、开发、运营、运维一体,上线只是开始。

免费获取你的建站方案

留下需求,专属顾问 24 小时内为你输出方案建议。