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

分布式任务调度实践:基于时间轮与延迟队列的ax调度全解析

发布时间:2026/9/28 16:16:27

资讯中心
01
ARTICLE

分布式任务调度实践:基于时间轮与延迟队列的ax调度全解析

分布式任务调度实践:基于时间轮与延迟队列的ax调度全解析
做后端这几年任务调度一直是个绕不开的坎。最近团队内部在推动一套轻量级调度组件代号就叫“ax”所以大家平时聊的“ax调度”其实不是某个开源框架的版本号而是我们自己封装的一套基于时间轮加延迟队列的分布式任务调度方案。今天把这套东西的设计思路、落地过程和踩坑记录完整梳理一遍如果你也在纠结“该不该自己写调度器”“用框架总差点意思”这类问题这篇应该能给你一些参考。ax调度主要解决的是业务里那些“定时要跑”、“延迟要触发”、“按条件要重试”的任务统一管理问题。上到订单超时关闭、离线数据补算下到消息补偿、报表生成都可以往这套调度器里塞。它不追求大而全而是强调轻量、可控、好排查适合准备自研调度组件的中小型团队也适合想在调度这个技术点上彻底搞懂原理的同学。我会把核心设计、关键代码、参数计算和线上问题都拆开了讲。1. 为什么我要自己写一个调度器ax调度的诞生背景1.1 先盘点一下市面上的调度方案在动手之前我们其实已经把常见的调度模式看了一圈。第一类是大家最熟悉的定时任务框架比如Quartz、Elastic-Job、XXL-JOB这类工具胜在功能完整cron表达式、失败重试、控制台都有。但问题是它们服务端和调度模型相对偏重部署一套就要引入额外的注册中心和管理端对我们这种“只想把任务跑好”的团队来说有点杀鸡用牛刀。第二类是基于消息中间件的延迟消息方案比如RabbitMQ的延迟插件、RocketMQ的定时消息。这种方案解决“延迟触发”很顺手但处理不了复杂的cron表达式也不方便做任务状态的统一查询。第三类是最朴素的数据库轮询通过扫描任务表里下一个执行时间来判断该跑哪些任务实现起来确实简单但一旦任务量上来要么加索引要么分表不然这个轮询本身就会变成新的瓶颈。我举这些例子不是想分个高下而是想说每个方案都有明确的能力边界。我们当时的需求其实很聚焦要支持cron定时、支持延迟任务、要能横向扩展执行节点、要能快速定位“任务到底跑没跑”。市面上的开源框架可以做但要配上各种外围组件去改造反而比从零写一个更费劲。ax调度就是在这种背景下诞生的。1.2 ax调度的设计目标与适用范围我们自己写的这个东西目标非常明确调度中枢要轻量执行节点要简单任务模型要统一。所谓轻量指的是调度器本身不依赖重量级外部组件只用Redis保存任务上下文和分布式锁用Zookeeper或数据库表都能做节点注册这样整个部署拓扑就是“调度器节点加执行节点加Redis”。执行节点简单指的是业务方接入时只需要提供一个普通的执行函数不需要关心任务是怎么触发的。统一指的是不管你是定时任务、延迟任务还是失败重试任务最后都抽象为同一个Task对象调度器负责在合适的时机把它投递出去执行器负责幂等消费。适用范围我觉得可以总结成三个特征第一任务总量在几万到几十万这个量级不需要像大数据平台那样百万级调度第二任务类型以RPC调用、SQL补数、消息推送为主没有特别重的流式依赖第三团队有很强的排障诉求希望每个任务的完整生命周期从产生、触发到执行、重试都能在一两个页面里看清楚。如果你的业务正处于这个阶段ax这种自研轻量调度方案的性价比会非常高。2. ax调度的核心模型任务、队列、触发器和执行器2.1 四个基础对象怎么定义ax调度里没有搞花哨的抽象核心只有四个概念Task、TaskQueue、Trigger、Executor。Task就是任务本身它不保存具体的业务参数太长的内容只保存一个taskKey和payload地址业务方执行时需要的数据存在Redis或数据库里避免Redis存储过大的事件内容。TaskQueue是任务队列它不是一个单独的数据结构而是按照任务类型分成多组队列。比如“定时任务队列”、“延迟队列”、“死信队列”。每组队列底层用Redis的zset或list实现。我们之所以不用RabbitMQ那一套是因为调度器需要频繁做“取下一批任务”的操作Redis zset天然支持按分数范围扫描配合Lua脚本能保证原子性比消费消息更直接。Trigger是触发器概念。它有两种cron触发器和延迟触发器。cron触发器保存表达式和下个执行时间调度线程每秒扫一次“时间到期的任务”把它们推进执行队列延迟触发器则干脆就是一条带上时间戳的记录时间戳到了就释放。Executor是执行器每个执行节点启动时会往Redis里注册一个心跳key调度器分配任务时按负载最低的节点优先投递。从模型上你可以看到ax调度其实把“调度决策”和“执行逻辑”完全解耦了。调度器只决定“谁、什么时候、去哪个机器跑”不再关心业务代码长什么样。这个解耦在线上非常重要因为只要模型清晰后续加新任务类型就不用改调度器。2.2 调度语义cron到延迟队列的取舍在调度语义上设计过程中反复犹豫了三轮。第一轮的想法是统一用cron表达式延迟任务就转成“当前时间加N秒”的cron结果cron最小粒度只有秒而且每次都要解析计算低效不说还要额外处理“今天不跑明天跑”这种跨天逻辑。第二轮改成统一用延迟队列但定时任务每天凌晨两点跑一次如果也做成延迟记录那意味着每次执行完又要重新插入一条延迟记录逻辑绕了一圈而且一旦任务失败重试很容易把定时语义弄混乱。最后定下来的方案是双向支持。cron任务走时间轮管理器每次任务执行完成后调度器根据cron表达式计算下一次触发时间并更新延迟任务直接写入Redis zsetscore就是执行时间戳。两者最后都会进入同一个调度循环从对应数据结构中取出当前时间属于待执行范围的任务放到投递队列里。这里有一个容易忽略的点我们叫它“时间片对齐”。定时任务如果每秒扫描一次可能在极端情况下把同一秒内新增和到期的任务混在一起导致重复投递。我们最终把调度动作拆成了“取任务”、“加锁”、“投递”三个步骤取任务时用Lua脚本按到期时间范围批量取出并且同时更新nextTime从源头减少重复。延迟队列同样如此release操作和删除操作必须原子完成。2.3 为什么选时间轮而不是每次都轮询数据库很多自研调度器都死于数据库轮询瓶颈。ax调度在设计之初就定了一条规则调度器节点尽量不接触业务数据库。时间轮是一个固定桶数的环形数组每个桶挂着一个双向链表链表节点就是到期任务。系统启动时我们初始化512个时间槽每个槽代表一个时间精度区间比如100毫秒一格。这样调度线程每隔100毫秒走一格把这一格上的任务全部释放出来。为什么不用Quartz那种按nextFireTime排序并select for update的方式因为数据库行锁在任务密集时会导致大量互相等待而我们用内存时间轮加Redis分布式锁的组合调度压力只落在内存遍历上Redis只负责最终投递和锁状态。实测单调度节点每秒可以处理近千次任务触发这已经远超我们现有业务峰值。不过时间轮也不是没有代价它需要消耗一定的内存维持时间槽。我们的场景里任务总量不算极端所以直接把任务对象放在轮子里如果你的任务是百万级我建议轮子里只放taskId任务详情留在Redis里取的时候再回查这样内存模型会更稳一点。3. 实操一个ax调度实例从0到13.1 环境准备与最小化部署我直接讲一套可以落地的最小化环境。我们用的组件是三台应用服务器一台Redis 6.x一台MySQL。调度器本身是一个Spring Boot应用但为了让你更好理解核心代码我这里提供一个纯Python版本的调度核心示例语言不重要逻辑是通用的。你需要安装redis-py、apscheduler中的cron解析部分但不用apscheduler的调度框架我们自己来控制时间轮。服务启动时调度器会先做三件事加载所有注册的task模板到内存连接Redis创建对应的队列结构注册当前节点到“调度节点列表”。这三个步骤任意一个失败节点都不能进入工作状态避免出现“以为自己活着但实际不在调度列表里”的诡异问题。下面是调度循环的核心骨架import redis import time import threading r redis.Redis(host127.0.0.1, port6379) def scan_due_tasks(): # 从zset中取出score小于等于当前时间的所有任务 # 同时删除这些记录保证同一批任务不会再次被取出 now int(time.time() * 1000) lua local tasks redis.call(ZRANGEBYSCORE, KEYS[1], 0, ARGV[1], LIMIT, 0, 500) if #tasks 0 then redis.call(ZREM, KEYS[1], unpack(tasks)) end return tasks tasks r.eval(lua, 1, ax:delay:queue, now) return tasks这段代码用Lua脚本完成了“取任务、删任务”的原子操作。很多同学会问为什么不先ZRANGE再ZREM因为这两步之间只要有别的调度线程插入同一个任务就可能被两个线程同时取出来导致后续重复派发。用Lua一次搞定是ax调度中最关键的一个细节。3.2 配置一个真实的调度任务假设我们现在要做一件事情每5分钟同步一次订单数据到报表库同时要求如果同步失败30秒后进行最多3次重试。在ax调度里我们把这个任务定义为一个TaskTemplate注册到调度器里。task_template { task_key: order_sync_5min, trigger_type: cron, cron: */5 * * * *, handler: orderSyncHandler, retry: { max_retries: 3, delay: 30 } }这里最重要的字段是handler它对应执行节点上注册的函数名。调度器不关心handler内部逻辑它只知道把这个taskKey投递到“order_sync_5min_logical_queue”里。执行节点通过订阅这个逻辑队列拿到任务ID再从Redis里取payload最后执行真正的同步逻辑。第一次实现的时候我们踩了一个不小的坑把payload直接放进了调度消息里。订单同步的payload是一批订单ID可能有几千个全部塞到Redis list里既占内存又让消息体变得很大。后来改成payload只存一个Redis的key执行节点使用前反序列化拿到真正的数据。这个调整让调度器到执行器的网络消耗下降了80%我建议任何调度系统都遵循“调度消息只传引用不传大对象”的原则。3.3 控制台与监控指标怎么看调度系统最怕黑盒ax调度上线第一天就标配了几个监控指标。第一个是“到期任务数”反映当前时间应该执行的任务总量如果这个数字持续上涨说明调度能力不够。第二个是“投递成功率”计算从调度器到执行节点之间的成功率这里主要反映网络问题和消费连接问题。第三个是“任务在线时长”展示每个worker节点的心跳延迟。我们的控制台其实很简单就是三个指标卡片加一张任务列表。任务列表按状态分等待中、执行中、成功、失败、重试中。这里要强烈建议你在自己的系统里也加一个“下次执行时间”的列不然排查“任务好像没跑”的时候很难判断到底是没到时间还是真丢了。为了直观还可以按队列加一组mini报表每天看一次延迟队列积压曲线基本就能把握调度系统健康度。关于监控指标我补充一句不要一开始就上几十个Grafana面板太多指标反而看不出问题。我们最有用的是两个任务执行延迟分布和worker空转率。任务执行延迟分布能告诉你延迟任务是不是准点worker空转率能告诉你目前机器是否冗余。4. 常见问题与排查实录4.1 问题一任务积压不执行调度日志却显示已投递这个问题在刚上线的第一个月暴露得特别明显。某天早上收到告警延迟队列任务积压了快两万条但是调度器日志里每秒钟都在打印“投递成功”。一开始我们怀疑是Redis网络波动后来发现是执行节点的消费线程池满了。因为大量业务任务都在同一个逻辑队列里消费某个慢任务把线程池所有线程都占住了其他任务虽然被投递到队列但根本没有可用的消费者线程去拉取。排查方法是看两个数字逻辑队列的剩余量以及执行节点线程池的activeCount。如果activeCount长期等于corePoolSize且队列待消费数不为0那就是典型的“消费者饿死”。解决思路是把不同优先级的任务拆到不同执行组每个组有独立的线程池。我们在任务定义里增加了一个group字段调度器按group生成不同的队列名执行节点也会为每个group创建独立线程池。调整后同一个慢任务只会堵死自己的分组不再波及全局。另外线程池参数不能拍脑袋。我们这边有个比较实用的计算方式单任务平均执行时间是t预期的任务到达速率是每秒n个那核心线程数至少要满足core n * t * 1.5多出来的是为了应对峰值。比如每秒到达10个任务每个任务平均0.5秒那核心线程数就是7.5取整加缓冲后设成10。队列长度则建议等于peakRate * maxBlockSeconds保证在消费者短暂不可用时调度器投递不会立刻被拒绝。4.2 问题二执行器重复消费同一个任务重复消费是分布式调度里最让人头大的问题之一。ax调度虽然用Lua保证了一个任务不会被同一个调度器同时取出但在执行端仍然可能出现重复。典型场景是执行节点处理完任务后更新状态的请求还没有写回Redis节点正好发生full GC导致消费超时。调度器按照超时重试策略又把同一个任务投递给了另一个节点最终两边都跑了业务代码。对这个问题的处理我们的核心思路是“业务层幂等兜底”。调度器确实会尽量做到at-most-once但任何调度器都不能保证网络完全可靠所以执行层必须带防重。具体做法是任务开始时往Redis写入一个带唯一任务ID的SETNX锁设置过期时间等于“预估最大执行时间加5秒”。只有SETNX返回1的节点才能继续执行其他节点直接放弃。等到业务执行完成删除这个锁。如果执行失败或机器宕机锁到期自动消失重试就不会被卡死。当年我们曾经因为锁过期时间设太短出现过“任务还在执行锁却提前释放”的情况导致同一个任务并发跑了两遍。后来把所有锁的过期时间改成动态计算estimated_millis * 2 5000并且每次执行业务代码前都会先刷新锁的过期时间。这是一个非常小的细节但能避免好几类线上事故。4.3 问题三任务触发时间漂移明明设置的整点却跑了4秒这个现象很多人会忽略但它会实实在在影响业务。我们的延迟队列本来是按当前机器时间作为score写入Redis的但线上服务器偶发时钟跳跃有一台机器时钟快了3秒它写入的延迟任务就会提前3秒被调度器取出。如果业务方在整点后立刻查数据可能就会看到“数据未更新”的假告警。处理这个问题要做两层。第一层是基础设施层所有服务器都必须接入NTP时钟同步并且要监测时钟偏移量超过200毫秒就要告警。第二层是调度器逻辑层我们在执行时间判断时引入了一个“统一时钟源”。简单说调度器不以本地时间为唯一标准而是每次调度循环开始时从Redis读取一条由主调度节点周期写入的时间记录用这个时间作为本次扫描的基准。如果本地时间和Redis时间偏差较大调度日志里会主动打印一条“clock drift detected”方便排查。实际操作中我们还会为最终用户提供一个“容忍度配置”。比如某些任务允许延迟5秒内触发那么扫描范围就从(0, now]调整为(0, now tolerance]但任务真正下发前还要再校验一次业务要求的执行时间避免所有任务都被提前释放。4.4 问题四任务重试风暴一个失败任务把整个队列打爆最后说一个比较隐蔽的问题。某个下游接口故障时队列里大量任务都进入重试逻辑。如果重试延迟设为常量比如10秒那每10秒就会有一批任务同时压向下游下游刚恢复又被重新打瘫。我们一开始也踩过这个坑后来把重试策略改成了指数退避加随机抖动。重试延迟不再是固定值而是min(initial_delay * 2^retry_count, max_delay) random(0, jitter_ms)。比如第一次重试延迟30秒第二次60秒第三次120秒再搭配最大延迟300秒和不超过30秒的随机抖动。这样能让重试流量在时间上均匀散开而不是整齐划一地冲击下游。还有一个细节延迟队列里我们额外加了一个“重试次数”字段任务被投递前先判断重试次数是否已达上限达到上限的任务会进入死信队列由第二天的人工巡检脚本统一处理。这样既不让业务静默失败也不让失败任务无限消耗调度资源。5. 基于ax调度扩展出来的几个实用技巧5.1 让cron任务跑出“秒级准时”的感觉刚才提到时间轮以100毫秒为单位扫描但cron任务在执行时还要经过网络投递、worker拉取等环节整体延迟大概在几十毫秒到几百毫秒不等。对于日志报表这类任务无所谓但如果你是支撑秒杀、抢券等场景最好把“到期触发时间”作为业务参数传给执行器让业务方自己根据这个时间判断是否需要补偿。我们在ax调度里把scheduleTime放进了每个task消息里。执行节点拿到后发现当前时间和scheduleTime差距超过阈值比如1秒会记录一条耗时日志。这比单纯靠外部监控要灵敏得多因为外部监控只能告诉你“任务开始时间”告诉不了你“任务本该开始的时间”。这个字段也是后面做任务耗时分析的基础。5.2 分布式节点扩容的免重启方案只要理解ax调度的模型扩容就不会有太大压力。执行节点是无状态的每次启动都会带着自己的分组标识来注册调度器向某个队列投递时会通过一个简单的“当前空闲worker列表”完成选择。所以扩容一台机器只需要启动一个新的执行节点它在注册到Redis的时候调度器自然会把任务分给它。唯一要注意的是执行节点注册后不要立刻参与消费应该等心跳信号积累两秒以上再进消费者线程。否则节点刚启动内部线程池还没完全就绪任务一到反而会堆积在队列里。我们在每个节点启动时加了一个warmup阶段用两秒做本地预热这个细节虽然小却避免了好几次上线后的短暂任务抖动。5.3 和业务系统的集成姿势最后建议一下集成姿势。ax调度本身不要和业务系统共享数据库事务尤其是执行任务的时候尽量通过RPC去调业务服务。调度器负责触发业务服务负责实际动作。这样业务代码哪怕定义得糟糕导致事务很长也不会拖住调度器本身的扫描线程。我们的业务接入基本是“两步走”第一步在ax后台注册一个taskKey和对应的handler名称第二步在执行节点代码里实现handler并用一个装饰器标注让它自动注册。执行节点启动时会把所有handler注册到Redis的哈希表里调度器投递任务时只需要查找这个哈希表如果找不到就直接记录“handler missing”。这样即使某个业务模块下线也只影响对应handler不会让调度主流程崩溃。6. 我在实际使用中的几点体会从ax调度立项到线上跑稳前前后后大概有三个月中间重构过两轮。第一轮重构是把所有延迟任务从统一队列改成按group拆分第二轮重构是引入时间轮替代数据库轮询。每次重构的诱因都不是“看某个技术不顺眼”而是线上具体的问题逼着我们必须改。所以我建议你也别一上来就追求完美架构先让核心链路跑通再根据业务反馈迭代。如果让我给准备自研调度器的团队提一个最实际的建议那就是一定要把任务的生命周期日志埋全。每个任务从注册到调度器到被时间轮取出再到投递执行节点再到执行完成每一步都要有带taskId的日志。因为调度系统本身就是为了解决分布式环境的不确定性如果它自己都不透明线上任何抖动都会变成灾难现场。最后分享一个我们在运维上养成的小习惯每周一次延迟队列积压巡检看着积压曲线慢慢往下走比看任何业务指标都有成就感。ax调度目前支撑了内部几十个核心业务的定时、延迟和补偿任务规模不算大但每一条任务都能说清楚“刚才发生了什么”这让我夜里睡得踏实不少。
02
RELATED NEWS

相关资讯

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

03
WHY YAOTU

想打造同款高转化官网?

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

◈

场景化定制

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

◐

营销型架构

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

▲

全周期服务

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

免费获取你的建站方案

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