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

自研分布式调度器ax:从状态机设计到生产级实践

发布时间:2026/9/28 17:06:46

资讯中心
01
ARTICLE

自研分布式调度器ax:从状态机设计到生产级实践

自研分布式调度器ax:从状态机设计到生产级实践
调度这件事最怕的不是任务多而是看起来一切正常突发时全盘崩溃。大约从压测环境里三分钟收到一批告警开始我意识到手头这套常规任务执行体系已经走到非重写不可的地步上游流量一抖动一批批任务挤在同一时间点触发部分节点负载冲高数据库连接被争抢队列越堆越长随后触发无效重试把系统二次打瘫。后来我用一个代号ax的自研调度器解决了这些问题它负责延迟任务、周期任务、批量任务编排与分布式下发。这篇文章想把这个项目从设计到落地的完整过程记录下来包括核心模型、触发机制、分布式一致性处理、执行节点的资源控制、可观测性和几轮真实踩坑给正在做任务调度中间件或想自研一套调度体系的同学做个参考。先说一下ax这个名字的由来我在代码仓库里给调度模块取名ax原本是 async executor 的意思后来团队里口口相传就成了调度器 ax。它本质上不是一个框架而是一个可以独立部署的调度服务替代了原先散落在业务代码里的各种Scheduled、while(true)扫表轮询和单机定时器让所有需要按时间、按条件、按批次触发的工作统一收敛到一个地方管理。1. 为什么从零写一个调度器自研的边界与底气1.1 商业调度产品的共性痛点市面上不是没有成熟的分布式调度方案开源的有 xxl-job、ElasticJob商业的也有各类分布式任务平台。早先我也不是没调研过最后没有直接拿一套开源框架来做主要原因有三类一是触发模型不够灵活。很多开源调度器强绑定 Cron 表达式适合固定周期但不太适合延迟 30 秒再执行依赖上一个任务成功后再触发按用户指定的时间窗口批量补跑这类场景。虽然可以在封装层强行拼但拼出来的代码往往是另一种形式的临时方案。二是和内部的执行环境耦合太重。我们这边任务执行器分布在各业务应用里有的在容器集群有的在裸机服务网络环境也不完全互通。开源调度器通常假设执行器能主动注册、心跳上报到一个中心调度端这会导致部分隔离网络里的执行器需要额外打隧道维护成本很高。我更想要一种调度端只负责任务状态推进执行端通过标准接口反编的模式。三是可观测性不足。调度不是一个孤立系统它旁边连着业务库、缓存、对象存储和消息队列。一旦某个环节慢调度器需要快速定位是触发慢了、派发丢了、执行超时了还是任务本身的问题。开源工具一般只给到任务开始/结束两级日志想做到链路级别的归因很费劲。所以自研边界其实很清楚我们不是要做一个通用调度平台去和成熟产品竞争而是要做一个贴合内部部署形态、能完全控制状态流转和可观测细节的调度底座。1.2 自研意味着要接管哪些事自研之前先列了接管清单避免做一半才发现责任边界模糊任务元数据存储任务定义、触发配置、执行结果、失败原因、重试次数。触发引擎支持 Cron、延迟触发、事件驱动触发三类能力。派发机制调度端如何通知执行端该跑哪个任务了。执行节点的资源控制并发线程数、单机最大同时执行任务数、队列缓冲。失败与重试语义失败几次算终态重试间隔如何退避是否需要人工介入。幂等保障同一个任务被重复下发时执行端按什么规则去重。审计与补偿如果任务半分半不遂例如已派给节点 A节点 A 返回超时但实际执行完了如何保证最终一致。列完这个清单心里就盘算好了不做秒级以下的实时调度不做通用工作流编排不做一个需要单独配置中心的规则引擎。守住这个边界才能把核心路径做扎实。后续所有设计都围绕这七件事展开。2. ax调度的核心模型状态机与调度分层2.1 任务状态机设计设计调度系统我建议第一步先把任务状态图画清楚。状态不清晰后面所有并发控制、重试、告警逻辑都会变成打补丁。ax里任务状态精简成 7 个PENDING任务已创建等待触发条件满足。TRIGGERED触发条件已满足等待派发。DISPATCHED已派发给执行节点等待确认。RUNNING执行节点确认开始执行。SUCCEEDED任务正常完成。FAILED任务最终失败进入终态。RETRY_WAIT执行失败且满足重试条件等待下一次触发。这里有个关键细节重试不是从任务开头重新走一遍而是在FAILED和TRIGGERED之间插入一个等待节点重试时的任务 ID 不变但触发版本号递增。每次重试都会生成新的执行上下文execution_context里面记录了第几次重试、上次失败原因、上次执行节点 IP这样排查问题时能直接看到同一任务 ID 下的多次尝试轨迹。下面的代码是状态机里最核心的流转判定我把它写成一个很朴素的advance方法public OptionalState advance(Job job, Event event) { State current job.getState(); switch (current) { case PENDING: if (event Event.TRIGGER_CONDITION_MET) { return Optional.of(State.TRIGGERED); } break; case TRIGGERED: if (event Event.DISPATCH_SUCCESS) { return Optional.of(State.DISPATCHED); } if (event Event.DISPATCH_FAILED) { return Optional.of(State.TRIGGERED); // 重新派发不增加重试计数 } break; case DISPATCHED: if (event Event.EXECUTOR_ACKED) { return Optional.of(State.RUNNING); } if (event Event.DISPATCH_TIMEOUT) { // 基于版本号重新入队 return Optional.of(State.TRIGGERED); } break; case RUNNING: if (event Event.EXECUTE_SUCCESS) { return Optional.of(State.SUCCEEDED); } if (event Event.EXECUTE_FAILED job.getRetryCount() job.getMaxRetryCount()) { return Optional.of(State.RETRY_WAIT); } if (event Event.EXECUTE_FAILED) { return Optional.of(State.FAILED); } break; case RETRY_WAIT: if (event Event.RETRY_TIME_ARRIVED) { return Optional.of(State.TRIGGERED); } break; default: break; } return Optional.empty(); }实际运行里这个状态机被包裹在各处事务中执行每三次成功流转都会写一条状态变更日志方便事后审计。2.2 三层调度结构接入层、调度层、执行层ax在逻辑上拆成了三层各层职责单一出问题时能快速定位边界接入层面向业务方提供 API接收任务创建、暂停、取消、查询请求。对接入层的要求是最好无状态可以任意横向扩容。它只负责把任务元数据写入存储层不参与任何时间相关的计算。调度层核心决策层。扫描到期的任务做触发决策然后通过派发通道将任务交给某个执行节点。多个调度节点可并行运行相同任务在同一时刻只能被一个调度器抢到。执行层运行在业务应用或独立执行服务里。它的职责是拉取任务、构建运行时上下文、执行真实的业务代码并把结果回传给调度层。接入层和调度层之间没有直接网络调用接入层写入任务之后调度层通过事件通知或者轮询感知新任务。这个异步过程一开始很多人不理解觉得为什么要绕一圈。这么设计是为了应对高频创建和低频触发的差异可能某秒创建了一万个任务但真正在下一秒需要触发的只有几十个。让接入层直接阻塞去触发既浪费资源也可能引发调用方超时。2.3 核心数据结构与 API任务表job的长相大致如下省略了部分索引字段我把注释写在字段旁边方便理解用途CREATE TABLE job ( id BIGINT NOT NULL AUTO_INCREMENT COMMENT 任务ID, job_code VARCHAR(64) NOT NULL COMMENT 业务方定义的唯一编码, state VARCHAR(16) NOT NULL COMMENT 当前状态, trigger_type VARCHAR(16) NOT NULL COMMENT cron/delay/event, trigger_conf VARCHAR(512) NOT NULL COMMENT 触发配置JSONcron表达式或延迟秒数, handler VARCHAR(128) NOT NULL COMMENT 执行器标识对应某个handlerBean, payload TEXT COMMENT 任务参数执行器透传使用, retry_count INT NOT NULL DEFAULT 0, max_retry INT NOT NULL DEFAULT 3, next_trigger_time BIGINT NOT NULL COMMENT 下次触发时间戳(ms), status_version BIGINT NOT NULL DEFAULT 0 COMMENT 版本号用于乐观锁和重试, biz_time BIGINT COMMENT 业务时间用于幂等判断, created_at DATETIME NOT NULL, updated_at DATETIME NOT NULL, PRIMARY KEY (id), KEY idx_state_next_trigger (state, next_trigger_time), KEY idx_job_code (job_code) );对外 API 我最常被问到的是任务挂在执行器上还是执行器挂在任务上。ax的答案是任务指定handler执行器应用里注册 handler调度层根据任务上的handler字段找到可用的执行节点列表。这种设计让同一类任务可以跑在不同的节点上也方便后续做灰度执行。3. 触发引擎与队列时间轮、延迟队列与优先级处理3.1 时间轮的实现与 N 轮降级触发引擎是调度器的心脏。ax第一版用了优先队列共享扫描的方式每个调度节点有一个next_trigger_time的索引每秒扫一次到期任务。数据量小的时候没太大问题但任务量上来后单纯依赖数据库轮询库的负载会成为瓶颈。后来引入了时间轮 数据库兜底的双层触发机制。时间轮放在内存里用于近一两分钟内即将触发的任务这一层承担每秒高频率的触发检查避免每秒把大量任务从数据库捞出来比对。更远时间的任务仍然放在数据库由调度节点定期把未来一个时间窗口内的任务预加载到时间轮。这个窗口值在ax里默认是 90 秒可按任务密集程度调整。时间轮本身是环形数组每个槽位放一个任务集合。指针每秒走一格落到当前槽的任务一次性取出并投递到触发线程池。当任务量特别大、时间轮槽位溢出时ax的做法不是强行堆内存而是将溢出的任务放回数据库并设置一个稍晚的next_trigger_time让数据库轮询兜底。这个近端内存、远端数据库的组合模式在任务规模从几万涨到百万级别时依然能保持稳定的触发延迟。3.2 延迟队列的消费模型对于延迟任务我没有直接使用常见的 Redis 延迟队列方案因为内部网络偶尔抖动Redis 不可用会造成任务丢失还需要额外的补偿逻辑。ax在数据库层实现了延迟队列每个任务写入时的next_trigger_time now delay调度节点按时间轮加载这本质上是把延迟队列退化成按触发时间排序的任务表省掉了一层中间件依赖。但有延迟任务时要特别注意批量加载的效率。我曾经写了一个按时段批量加载的逻辑前缀是WHERE statePENDING AND next_trigger_time ?执行计划偶尔走错索引扫了全表。后来强制加了trigger_type字段作为前缀条件才稳定下来。一个看似简单的时间字段一旦命中大批量历史数据优化器很容易犯迷糊这个坑后面展开讲。3.3 优先级与资源配额如何叠加ax引入了优先级概念但优先级不是单独排序而是与触发时间结合priority_score next_trigger_time * 100 priority_level。时间越早、优先级越高越靠前触发。这种设计的好处是系统不会因为一个高优任务无限等待也不会因为一堆低优定时任务抢占即将触发的延迟任务。所有任务都围绕最迟完成时间这唯一的锚点排序。光有排序还不够执行端必须配合资源配额。每个执行节点限制了最大并发执行数max_concurrency调度层在派发前会查询目标节点的实时负载和排队任务数超过阈值时暂缓派发。这解决了调度层疯狂下发执行层瞬间堵死的问题。实现时执行层用一个信号量控制并发核心参数如下public class ExecutorLimiter { private final Semaphore semaphore; private final int maxWaitMs; public ExecutorLimiter(int maxConcurrency, int maxWaitMs) { this.semaphore new Semaphore(maxConcurrency); this.maxWaitMs maxWaitMs; } public boolean tryAcquire() { try { return semaphore.tryAcquire(maxWaitMs, TimeUnit.MILLISECONDS); } catch (InterruptedException e) { Thread.currentThread().interrupt(); return false; } } public void release() { semaphore.release(); } }试想一个场景某个执行节点上所有 worker 线程都卡在等待数据库连接池新的任务还在不断派过来如果不加信号量限制最终就是连接池报错、线程堆满、节点假死。加了ExecutorLimiter之后调度层派发前就能感知到节点饱和任务会停留在TRIGGERED状态等节点恢复后再继续派发。4. 分布式一致性抢单、续约与脑裂排查4.1 基于数据库行锁的抢单实现多个调度节点同时运行时必须保证同一个任务不会被两个节点同时触发。ax没有引入额外中间件而是用数据库行锁当调度节点扫描到到期任务后先执行一个抢占更新将任务的status_version加一并要求status_version不变。UPDATE job SET status_version status_version 1, state TRIGGERED, updated_at NOW() WHERE id ? AND status_version ? AND state PENDING;只要UPDATE影响行数为 1就说明当前节点抢到了这个任务的任务所有权可以安全进入派发流程。这个操作的粒度非常大同一时刻同一个任务只有一条行锁生效从根上避免了双节点同时触发的问题。相比引入分布式锁基于行锁的实现更简单、更贴近数据本身的一致性边界配合status_version也天然支持了重试隔离。4.2 任务续约体系与过载保护任务被派发出去后执行端需要回报状态。ax没有设计派发后不闻不问的死信模式而是采用续约机制执行端开始执行后每 10 秒向调度端上报一次心跳续约调度端在等待执行结果时会记录每个任务的last_heartbeat_time。如果一个任务超过heartbeat_timeout默认 30 秒没有收到续约调度端会认为该任务可能跑飞或者节点宕了主动将任务重新置为TRIGGERED并重新派发。这个机制特别适合那些执行时间长的任务比如大量数据清洗、脚本编译、文件处理这类任务无法快速判定结果只能靠活着来证明没死。续约本身也会产大量请求所以ax将续约请求做成批量接口执行端把执行中的任务 ID 数组打包每 10 秒上报一次调度端一次更新多个任务的心跳时间。接口返回时顺带告知哪些任务已被标记为超时执行端会收到放弃指令停止运行对应的业务逻辑。4.3 脑裂与双跑排查实践脑裂场景在调度系统里是最隐蔽的调度端以为执行端失联把任务重新派给另一个节点此时原节点可能还在运行于是同一个任务在两个地方同时执行。最常见的原因是网络瞬断超过续约超时或者执行节点 GC 停顿超过了超时阈值。ax处理双跑的核心思路是幂等键 执行前仲裁。每个任务在开始执行前执行端会往结果表写入一条执行记录以job_id status_version作为唯一键。先插入成功的一方继续执行插入失败的说明已有更新版本在执行主动放弃。这个逻辑避免了依赖网络一定没问题的假设比靠心跳判断可靠得多。INSERT INTO job_execution (job_id, status_version, executor_ip, start_time) VALUES (?, ?, ?, ?) ON DUPLICATE KEY UPDATE status_version status_version;另外ax会把双跑事件记录为告警而不是静默忽略。即使幂等逻辑兜住了也要让运维知道发生过网络分区分裂因为这通常是更深层次网络问题的前兆。5. 执行节点与资源控制并发、超时、重试与优雅停机5.1 并发分片与线程隔离执行层接到的任务类型五花八门既有跑一秒就结束的缓存预热也有执行半小时的数据批处理任务。把它们全部丢进同一个线程池会相互挤占。ax的执行层按handler类型分池每个 handler 对应一个独立线程池池子的核心线程数和最大线程数都由该类型任务的日均量、执行时长和资源开销决定。分片同样发生在执行层。大任务进来后执行节点根据自身分片配置唤醒多个 worker 并发处理分片信息会写入任务上下文。设计时的一条经验分片数不要和执行节点数强绑定尽量允许运行时动态调整否则节点扩缩容时又要停任务、改配置、重新派发十分折腾。线程池的拒绝策略也很关键。默认使用CallerRunsPolicy会让调度线程直接执行任务阻塞派发链路非常危险用AbortPolicy直接丢任何任务更是不可取。ax最终采用DiscardOldestPolicy的变体队列满时不丢弃现有任务而是返回一个节点繁忙的信号调度层收到信号后对此任务延后 2 秒重试。这个方案不太激进又不会让调度线程去跑业务代码。5.2 超时判定与主动取消任务粒度上的超时常被分为两类排队超时和执行超时。排队超时指任务被派发后一直没有拿到执行线程持续占用调度层的派发名额执行超时指任务已经开始执行但超过了预设的时间上限。处理排队超时比较简单执行端在收到派发请求时记录dispatch_arrive_time若任务在队列中滞留超过queue_timeout直接向调度端回报排队超时调度端重新排队。执行超时则复杂一些因为任务超时了要不要 kill 正在跑的线程是个问题。ax默认不物理终止线程只通过Future.cancel(true)标记中断若执行代码没有响应中断则任务进入超时标记状态等实际跑完后由执行端主动上报 result。同时调度端会基于续约机制在超时未恢复时将任务重新置为可重试保证整体进度能继续推进。5.3 重试策略与幂等保障重试策略在ax里是可配置的但有个通用默认值最大重试 3 次间隔采用指数退避10s、30s、90s。退避不是固定值而是叠加了随机抖动避免大量失败任务同时进入重试潮把本就脆弱的下游打挂。重试语义上ax区分了可重试异常和不可重试异常。相当于自己定义了一组异常类型例如RetryableException和NonRetryableException。执行器抛出NonRetryableException时直接进入终态FAILED抛出RetryableException时才走重试逻辑。这比所有异常都重试要安全得多数据不满足条件重试一百次也是白搭下游依赖短暂不可用重试才可能救回来。幂等保障除了前面提到的唯一键插入之外还有一层业务幂等键。任务可以携带biz_time之类的业务时间戳执行器在处理时将job_code biz_time作为落库唯一键天然保证重复派发时不会把同一笔数据写入多次。6. 可观测性建设链路、指标与止血手段6.1 埋点设计与指标口径调度系统如果没有指标就像蒙眼开车。ax的指标分三层触发层每秒触发任务数、触发延迟 P99、时间轮槽负载、数据库扫描耗时。派发层派发成功数、派发失败数、派发超时数、执行节点负载分布。执行层任务执行成功数、失败数、重试数、执行时长分布、续约次数。指标口径要尽早定死否则后期很容易出现谁的数字对不上。比如执行时长是从执行端收到请求开始算还是从任务真正开始执行业务代码开始算ax统一从handler方法入口开始计算执行前准备和资源获取耗时另算一个overhead指标。分开统计之后才能知道到底是业务算法慢还是资源分配慢。6.2 全链路追踪方案任务调度往往横跨接入层、调度层、执行层做好链路追踪能大幅缩短定位时间。ax在任务创建之初就生成一个trace_id写进任务上下文。派发时trace_id随派发请求带到执行端执行端上报结果时再原样带回。各处日志都带着trace_id通过日志系统一次搜索就能把一条任务从产生到结束的完整轨迹拉出来。这里要提一个细节很多系统会忽略执行端的运行时日志只记录调度层的操作日志结果排查问题时发现任务根本没有在业务代码里输出一行。ax的执行端统一做了一个JobLogContext把任务 ID、版本号、trace_id 都塞进MDC业务代码里只需要写普通日志日志平台就能自动聚合出某个任务的全部运行时上下文。6.3 三个止血开关线上出问题时最需要的是瞬时止血、快速恢复而不是深入分析。ax的运维控制台常备三个开关我称之为止血三件套任务暂停开关把某个任务的调度状态置为暂停一旦置为暂停调度层不会再扫描它。适合上游临时故障时先把下游任务暂停防止雪崩。派发黑洞开关如果执行节点整体异常比如发布过程中被注册中心剔除了但调度层还不知道可以先打开黑洞开关让调度层不再派发新任务。重试熔断开关当某个执行节点的重试成功率明显低于阈值时自动熔断这个节点上的所有任务不再让它们快速重试。观察一段时间后再关闭。这三个开关都要求在毫秒级生效因此不能依赖任务扫描周期而是通过配置中心推送调度层监听配置变更后实时更新内存中的开关状态。7. 部署与版本迭代中的真实踩坑7.1 跨时区与夏令时引发的延迟ax第一版上线后产品提了个需求业务活动在特定时区的早上九点开始任务需要在本地时间九点整触发。我第一次实现时直接用了系统默认时区结果部署在 UTC 环境的节点果然出了问题活动开始时间整体偏移了八个小时。排查后发现不仅时区重要夏令时切换也隐藏着雷。比如有的时区在每年三月某日凌晨 2 点跳变Cron 表达式如果用0 30 2 * * *会存在一个根本不存在的凌晨 2:30。ax的修复方式是存储触发配置时统一使用带时区的完整 cron 描述内部计算触发生效时通过ZonedDateTime解析出精确的next_trigger_time避免使用裸的LocalDateTime。每次换算都明确指定ZoneId并在配置中心里维护一份时区白名单上线前对所有任务做一次未来 24 小时的触发时间预览直接暴露异常。7.2 数据库连接池被打爆的教训分布式调度系统最大的隐忧就是依赖数据库连接池。某个版本上线后业务方集中创建了大量短周期定时任务调度节点疯狂扫描导致数据库连接池被打满连带业务主应用报连接超时。这不是任务数量本身多而是每个调度节点都新建了太多长连接且没有控制扫描频率。后来做了三个改进所有扫描任务统一走一个共享只读事务减少连接占用。每批扫描最多带出 200 个任务处理完再拉下一批避免单次查询占用大量行锁。数据库连接池最小空闲连接数调低最大活跃连接数按调度节点数量动态调整。经历过这一步之后我再也不敢轻视扫描频率的设置有句糙理不糙的话调度器是踩在数据库肩膀上的它不珍惜连接库就给它脸色看。7.3 滚动升级的暂停机制调度服务部署升级时最怕的是旧节点还在派发任务新节点已经接管了一部分双节点状态不一致。ax的解决办法是引入滚动升级暂停窗口升级时先把节点标记为draining不再参与新的任务派发等待当前执行中的任务全部结束或超时再真正下线旧节点同时新节点注册后并不立刻接管全部任务而是先进入warming状态等时间轮预加载完成、数据库扫描正常后再开始处理。这套机制看起来简单但对老节点的等待时间要做上限控制。如果某个长任务执行 40 分钟你不能干等它结束而是设置一个最长 2 分钟的 drain 超时超了就把剩余任务置为可重派状态让新节点接手。这也是ax调度层里唯一允许的状态滞后场景通过后续补偿来保证最终一致。写在最后的一段经验这套ax调度器从设计到稳定运行前后迭代了大半年。我最大的体会是调度系统的难点并不多在触发任务这个动作上而在状态边界、异常恢复和资源保护这三个地方。任务状态机越早固定越好重试语义越早明确越好派发路径上的每一个故障环节都要有对应的补偿动作。很多问题看起来是并发导致实际上是对状态的假设不够严谨。如果你也准备自研调度体系建议先从最小闭环开始单机、单库、三个状态、两种触发跑通后再谈分布式和可靠性。把状态机整理清楚把超时和重试的口径统一掉至少能规避大半线上事故。
02
RELATED NEWS

相关资讯

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

03
WHY YAOTU

想打造同款高转化官网?

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

◈

场景化定制

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

◐

营销型架构

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

▲

全周期服务

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

免费获取你的建站方案

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