☰
ax调度两级模型实战:队列限流重试与幂等设计全解析
2026/9/28 17:22:07 网站建设 项目流程

最近好几个群里都在聊 ax 调度这个词,一开始我以为是哪个框架又出了新特性,翻了一圈代码仓库和历史笔记才确认,大家说的 ax,其实就是 Accelerated eXecution,一种把排队、限流、分发、重试揉在一起的两级调度模型。核心思路特别朴素:让任务在正确的时间、以可控的并发度、落到真正有能力消费它的机器上。

它不是什么灵丹妙药,但确实解决了很多后台系统“平时安静如鸡、一到高峰期就连环超时、任务一多就死给你看”的老毛病。这篇不是官方文档,是我自己从零搭过、压过、也排过不少障之后的落地记录,适合正在折腾定时任务、批处理、工作流调度,或者被 cron 和 Redis 队列反复折腾的同学参考。

1. 先把 ax 调度的设计思路拆开看:两级模型到底解决了什么

1.1 ax 调度是什么:前台分诊加科室排班的两级协作

ax 调度的核心,是把“任务接收”和“任务执行”彻底分开,而不是像老式任务系统那样,一个队列接一个执行器组,从头到尾一条流水线。

一级调度负责任务的“到达管理”。所有任务统一从入口进来,做合法性检查、幂等去重、按优先级排队、判定是否允许立刻投递。这一层不关心任务具体怎么跑,只负责把任务“接住”并“整理好”。

二级调度负责任务的“执行管理”。它根据当前执行器的负载情况、令牌桶余量、队列优先级,把任务从等待区拉出来,分发给执行器,并跟进任务的完成确认。

我常拿一个类比来解释:一级调度是医院前台分诊台,先把病人分科、挂号、排号,急诊优先;二级调度是科室里的排班护士,知道哪个医生手头有病人、哪个现在空着,再把号叫进去。两者如果混在一起做,就会变成“医生既要站在门口接诊,又要自己刷号叫号”,医院早就乱套了。

这个设计最值钱的地方在于,两级可以各自独立扩容。任务量暴涨时,只需要把一级调度的入队能力和队列容量扩上去;执行能力不够时,单独给二级调度挂更多执行器。瓶颈在哪里,就扩哪里,不会整个系统跟着一起动。

1.2 为什么不能直接搬老一套:单队列模式到底卡在哪

很多人问过我,之前用 Redis 的 LPUSH / BLPOP 做个队列,多个 worker 一起消费,不也能跑吗?为什么要费劲搞 ax 调度这种两级模型?

老实说,任务量小的时候,Redis 队列完全够用,我也不是让大家一上来就推翻旧系统。但一旦任务形态复杂起来,单队列模式有几个绕不开的问题。

第一个问题是没法表达优先级。普通任务队列里的元素没有区分度,所有任务都是平等的。但现实业务里,数据补偿任务可以等 10 分钟,用户主动触发的同步任务最好在 1 秒内执行。放到同一个队列,要么大家都排队,要么你用多个队列自己造轮子。ax 调度里用的是按优先级分桶的多级队列,高优任务进来可以直接插到前面,低优任务再满也不至于堵死核心链路。

第二个问题是热键竞争。多个 worker 同时从同一个 key 里 BLPOP,Redis 单实例下所有请求都打到同一个 key,连接一旦多起来,Redis 自己在等待锁和网络回包上就消耗不少。尤其执行器数量上到几十个之后,队列本身反而成了瓶颈。

第三个问题是没有回执机制。任务从队列里 pop 出来,worker 开始处理,如果 worker 在任务执行到一半时宕机了,这个任务就永远丢了。因为没有“我拿了但我没做完”这种中间状态。ax 调度在两级之间引入了一个隐藏的暂存区,任务从队列里取出来不算结束,要给调度器回执确认,才算真正消费完成。单队列模式想做到这一点,就要自己额外维护 in-flight 状态,代码很快会写成一团乱麻。

2. 落地前必须抠清楚的四个核心细节

2.1 队列该不该有界?这是我被压测教做人的一次

ax 调度里的队列,建议做成有界队列。很多人一听,第一反应是“有界会不会丢任务?”其实是不会的,因为一级调度在任务入队之前还有一道入口,队列满了之后,不是把新任务丢掉,而是把新任务挡在入口层,可以先持久化到落盘存储里,也可以直接告诉上游“现在繁忙,稍后再试”。

为什么要加这个限制?我早期做过一个无界版本,任务高峰期时队列积压了上百万条,内存直接飙上去,GC 时间变长,然后整个调度引擎的响应速度开始线性恶化。更麻烦的是,Redis 里的 list 越来越大,后续清理、扫描、迁移都变成灾难。加了一个 max_queue_size 参数之后,系统反而变得健康:队列满了就触发背压,上游感知到压力会主动降频,而不是闷头一直灌。

具体的做法是按优先级分桶,每个优先级队列有独立的容量上限。高优队列可以相对短一些,比如 1000,因为高优任务本来就应该快速消化掉;低优队列可以给更大空间,比如 10000。这样即使低优任务大量堆积,也不会抢占高优任务的队列位置。

2.2 令牌桶限流到底在保护什么

很多人做调度系统,第一版只关心“任务有没有派出去”,不太愿意做限流。直到某一次压测,我把 2000 个任务同时丢给了执行器,执行器调用的老数据库接口直接被打到连接池耗尽,我才意识到:调度器的职责不只是把任务派出去,更要保护下游系统不被冲垮。

ax 调度里用的是令牌桶,而不是简单的信号量。信号量只能限制“同时在跑的任务数”,但没法控制“单位时间内的启动速率”。比如有 1000 个任务,每个任务只需 10 毫秒,在信号量限制下,结果就是 10 毫秒内瞬间启动了一大波任务,下游还是会被突刺打懵。令牌桶控制的是速率,每秒恒定放出多少张令牌,任务想执行,必须先拿到令牌。

每个执行器拉取任务的时候,会从调度器拿一批令牌配额。调度器根据执行器上报的压力量级,动态调整发令牌的速率。不像静态配置那样死板,每个机器能承受的 QPS 不一样,按各自实际能力拉取,整体吞吐反而更好。

2.3 任务重试是常态,幂等才是保命符

分布式系统里,任务重试是必然的。网络闪断、执行器重启、超时判定,都会导致同一个任务被投递多次。所以 ax 调度默认采用 at-least-once 的投递语义,但这个语义必须搭配幂等消费,否则重试一次,数据就被重复处理一次。

幂等设计的关键,是给每个任务安排一个全局唯一键。这个唯一键不能是数据库自增主键,它必须来自业务本身。比如“用户 12345 在 2025-06-01 的下单积分补偿”,可以根据用户 ID 加业务日期做 hash,生成一个稳定的业务幂等键。执行器处理前先查一下幂等表,如果已经处理过,直接丢弃或返回成功。

这个方案的坑在于幂等表的存储和清理。处理记录不能膨胀得太快,我在实际项目里用的是 Redis 加布隆过滤器双层结构:先过布隆过滤器过滤掉肯定没见过的任务,再查 Redis 里的近期处理记录。同时给 TTL 设置到任务的“最大可能重复窗口”以上,比如重试最长持续 3 天,TTL 就设置 7 天,留足余量。

2.4 一组可以直接抄的默认参数

以下是我调过几个项目之后觉得比较适合中等规模系统的初始参数,抄回去之后再根据自己的业务形态调整:

参数默认值说明
max_queue_size10000每个优先级队列的最大积压条数,超限触发背压拒绝新任务
high_priority_ratio3:1高优队列与低优队列的轮询出队权重
worker_pool_sizemin(CPU*2, 16)单机执行器并发线程数,保守一点没错
dispatch_batch_size20二级调度单次拉取的任务批量大小
dispatch_timeout_ms3000调度器等待执行器回执的超时时间
retry_max3单任务最大重试次数,超过进入死信队列
idempotent_key_ttl7d幂等键在 Redis 里的保留时间
heartbeat_timeout_ms15000执行器连续两次心跳超过该值即判定失联

这几个参数的核心逻辑是:宁可保守,不要激进。dispatch_batch_size 设小一点,任务分批拉,避免一次性拉太多导致执行端内存波动;heartbeat_timeout_ms 设大一点,避免网络抖动造成大量任务误判超时、重复投递。

3. 实操记录:从零复现一个可运行的 ax 调度器

3.1 最小化环境:一个 Redis 加两台执行器就够了

这个示例不需要云平台,本地 Docker 就能跑起来。我用的版本是 Python 3.10、Redis 6.2,执行器用最简单的 FastAPI 对外暴露任务处理接口。

先写一份 docker-compose,把 Redis 和调度器跑起来:

services: redis: image: redis:6.2-alpine ports: - "6379:6379" command: redis-server --appendonly yes scheduler: image: python:3.10-slim volumes: - ./scheduler:/app working_dir: /app command: python main.py environment: - REDIS_HOST=redis - REDIS_PORT=6379 worker1: image: python:3.10-slim volumes: - ./worker:/app working_dir: /app command: python worker.py worker1 environment: - REDIS_HOST=redis worker2: image: python:3.10-slim volumes: - ./worker:/app working_dir: /app command: python worker.py worker2 environment: - REDIS_HOST=redis

这里 Redis 开启了 appendonly,主要是为了保证任务和回执数据不丢。调度器本身是个常驻进程,跑一个事件循环,从 Redis 里不停地读队列、派任务、收回执。

3.2 核心代码:延迟队列为什么要用 zset

我在示例里用 zset 做任务触发队列,每个任务的 score 是计划触发时间戳。调度循环每秒做一次轮询,用 ZRANGEBYSCORE 取出当前时间之前的任务,再批量推入就绪队列。

关键代码如下:

import redis import time import uuid class AxScheduler: def __init__(self, redis_client): self.r = redis_client self.queue_names = { "high": "ax:q:high", "low": "ax:q:low", } self.dispatch_lock_key = "ax:dispatch:lock" def submit(self, task_id, biz_key, body, priority="low", delay=0, max_retry=3): score = time.time() + delay task = { "task_id": task_id, "biz_key": biz_key, "body": body, "priority": priority, "max_retry": max_retry, "retry_count": 0, "create_at": time.time(), } self.r.zadd("ax:zset:schedule", {json.dumps(task): score}) def dispatch_loop(self, batch_size=20): now = time.time() due_tasks = self.r.zrangebyscore("ax:zset:schedule", 0, now, start=0, num=batch_size) for task_json in due_tasks: task = json.loads(task_json) qname = self.queue_names[task["priority"]] # 入队前检查队列长度,超限就背压 if self.r.llen(qname) >= self.max_queue_size: continue self.r.lpush(qname, json.dumps(task)) self.r.zrem("ax:zset:schedule", task_json)

这里我觉得最需要注意的一点是,任务从 zset 里取出来到放入就绪队列,整个过程要保证原子性。上面这段简化代码没做事务,真实落地时必须用 Lua 脚本或者 Redis 事务把zrangebyscore + lpush + zrem包成一个原子操作,不然两个调度实例同时运行时,同一个任务会被重复取出、重复入队。

执行器侧的逻辑更简单,worker 启动后循环执行:

def run_worker(worker_id): while True: # 从高优队列优先获取,超时等待 1 秒 task_json = r.brpop(["ax:q:high", "ax:q:low"], timeout=1) if not task_json: continue key, value = task_json task = json.loads(value) # 幂等检查 if r.sismember("ax:set:done", task["biz_key"]): continue try: process_task(task) r.sadd("ax:set:done", task["biz_key"]) r.expire("ax:set:done", 3600 * 24 * 7) except Exception as e: retry_count = task.get("retry_count", 0) if retry_count < task.get("max_retry", 3): task["retry_count"] = retry_count + 1 # 指数退避: 2^retry_count * 30秒 r.zadd("ax:zset:schedule", {json.dumps(task): time.time() + 2 ** retry_count * 30})

worker 用 brpop 的时候把高优队列放在参数列表最前面,这样只要有高优任务,优先拿到的一定是它。很多半路出家的调度器最容易在这里犯错,把低优放在前边,高优任务的响应时间就平白无故被拉高了。

幂等检查放在任务真正处理之前,处理成功后才写入完成集合。这里有个细节是r.expire("ax:set:done", ...)每次成功都会刷新整个 key 的过期时间,好处是活跃任务的幂等记录不会半路被清掉,坏处是如果同一业务键持续高频出现,这个 set 会一直膨胀。我的处理方式是给幂等键补一层带 TTL 的单独 key,而不是只用一个大 set,过期的自动消失,不会无限增长。

3.3 上线后的第一波压测:结果和预期不一致时怎么定位

第一次压测我用了最简单的模型:10000 个任务,每个任务处理耗时平均 200 毫秒,4 个 worker。理论上每台 worker 每秒能处理 5 个任务,4 台就是 20 QPS,整个批次理论耗时应该在 500 秒左右。

实际跑出来的结果,总耗时超过 1000 秒,吞吐打了对折。我当时第一反应是 worker 性能不行,后来单测每个任务处理确实稳定在 200 毫秒左右,排除了处理函数本身的问题。

接着我盯了调度器日志,发现一个奇怪现象:worker 进程经常处于等待状态,但就绪队列里几乎看不到积压任务。理论上任务应该早就被拉走了,为什么没有数据在跑?

问题出在 brpop 的超时时间和调度循环批次太大之间的不匹配。调度器每次从 zset 里取 20 个任务入队,但 worker 每取完一个就要重新调用一次 brpop,brpop 的 timeout 设置得又比较长,导致下一个任务要等好几秒才开始处理。简单说,任务不是被处理得太慢,而是被“取”得太慢。

这个问题的解法是把 brpop 的超时调短到 300 毫秒,如果等不到任务就立刻循环,重新发起下一次拉取。同时 worker 端改为批量拉取,一次从队列里 pop 出 5 到 10 个任务,放进本地内存队列,再逐条处理。调整之后,整体吞吐很快达到了 18 到 19 QPS,基本贴近理论值。

还有一次压测是调度器单实例部署,任务量大时调度进程 CPU 冲高到 90%,日志显示大部分时间花在了 zrangebyscore 的扫描上。后来我把 zset 里的 score 改成了更粗粒度的分钟级时间戳,并给调度循环加了一个自适应的 sleep,任务少时循环频率降低,任务多时频率提高。CPU 占用降到了 30% 左右,响应时间也稳定下来。

4. 上线后最容易踩的五个坑:排查实录

4.1 worker 明明闲着,任务却在排队区堆着:问题出在调度节奏上

现象:从监控面板看,就绪队列一直在涨,但 worker 的 CPU 使用率很低,像在摸鱼。

排查思路先看 worker 拉取日志,确认它是不是真的在循环里跑。如果日志每隔很久才刷一条,多半是拉取节奏过慢,brpop 或 block 取任务的超时时间设得太长,或者调度器隔很久才往队列里放一批任务。

我遇到最隐蔽的一版问题,是调度器有个 bug,当队列长度超过 max_queue_size 后直接 continue,但这个 continue 跳过了zrem,导致同一个任务卡在 zset 里反复被取出来、反复入队失败。从外面看任务一直堆积,其实是调度器在空转。所以队列长度检查一定要跟延迟队列的移除操作放在同一个事务里,要么一起成功,要么一起不执行。

4.2 超时任务被重复触发:这个阈值别硬凑好看的数字

现象:同一个任务明明很耗资源,结果每隔几十秒就重复执行,数据库里出现大量重复记录。

这类问题多半是调度器对“任务超时”的判断太激进。executing 状态的任务,执行器会周期性上报心跳,一旦调度器超过一定时间没收到心跳,就判定任务挂掉,重新投递。

这里的核心矛盾在于,任务执行耗时天然有波动,不能按均值设超时。我见过有人按任务平均耗时 300 毫秒,把超时定在 500 毫秒,结果某个 JVM 触发 GC 停顿 1.2 秒,任务就被重复调度了。超时阈值至少要按“P99 耗时乘以 2”来设,也就是正常情况下绝大多数任务都能在阈值内完成,只有真正死掉的任务才会触发重投。

另一个隐藏细节是,重复执行不一定来自调度器的超时重投,也可能来自执行器内部的重试框架。你处理完了任务,但在准备上报 ack 的瞬间网络闪断了,调度器没收到 ack,任务被重新分发。这个时候任务的幂等键设计就显得极其重要,否则重试一次就是一次灾难。

4.3 节点宕机后任务“凭空消失”:ack 的顺序不能错

现象:凌晨一台 worker 宕机,重启后整个团队发现好几个小时前提交的任务不见了,没有任何失败日志。

原因基本都在“先处理后确认”还是“先确认后处理”的顺序选择上。我的建议是,业务处理成功之后,再向调度器发送 ack;没有 ack 的任务,在等待一段时间后会被自动重新入队。

但有个反常识的细节:ack 不应该放在业务代码的 return 语句之前,也不应该放在 finally 块里。正确位置是业务数据落库、拿到数据库事务提交成功之后。因为你 ack 太早,宕机时数据库事务还没提交,任务其实没真正“做完”,但调度器已经认为它完成了,这个任务才是真的凭空消失。

反过来说,如果确认过晚,会稍微增加重复执行的概率,但配合幂等键,影响可控。两害相权取其轻:宁可多跑一次,也不要少跑一次。

4.4 分布式锁失效导致的重入问题:降级到 fencing token

ax 调度在分发任务时会用分布式锁保证同一时间只有一个调度器在处理同一个任务。这个锁本身能挡住多数并发问题,但遇到长任务还是可能翻车。

我遇到过一次典型故障:执行器处理一个任务花了 40 秒,当初给锁设置的过期时间是 20 秒。执行到一半锁过期了,另一个调度器实例发现这个任务“无人认领”,立刻重新派发了一个副本。两台机器同时跑同一个任务,一头在写 A 表,另一头也在写 A 表,数据直接乱套。

只把锁过期时间调大治标不治本,因为 GC 停顿或网络卡死,锁过期时间再大也可能被突破。正确的防重入方案是引入 fencing token:调度器在发锁时生成一个单调递增的编号,执行器在处理任务时把这个编号带上;如果发现有更新的 token 出现,旧的那个处理过程要主动自我终止。这个机制跟数据库的版本号乐观锁是一个道理,防止的不是锁被抢走,而是旧执行者不知道锁已经不属于自己了。

4.5 时钟回拨引发的触发乱序:zset 的 score 别裸用服务器时间

ax 调度大量依赖时间戳排序,比如延迟队列的 score。单独一台机器还好,一旦调度器多实例部署,各台机器的时钟不可能完全同步,甚至同一台机器也可能因为 NTP 同步出现时钟回拨。

时钟回拨的典型症状是:本该延迟 10 分钟执行的任务,当场就被触发了;或者已经执行完的任务,突然又延迟几分钟后重新出现。因为 score 是时间戳,回拨后新写入的任务 score 反而比旧任务更小,触发了排序错乱。

我在生产环境里的处理方式是:调度器内部不直接用墙上时钟,而是用单调递增的本地序列号,配合时间戳共同组成排序值。跨实例的场景则引入混合逻辑时钟,这种时钟可以保证事件之间的先后顺序不会因为回拨而混乱。

如果项目不想引入额外组件,最简单也实用的兜底方案是:检测到系统时间发生回拨(比如当前时间比上一次记录的时间小),就暂时拒绝接收新的延迟任务,直到时钟恢复到原先的位置。同时把延迟任务统一延后 10 秒再入队,用时间缓冲消化掉误差。

还有一个容易忽略的点:执行器的任务触发时间记录也要谨慎。执行器本地时间跟调度器对不上,可能导致任务实际执行时间比预定时间早或晚很多。我对执行器做过处理,所有时间统一以调度器返回的时间戳为准,执行器只负责执行,不负责自己判断“到没到时间”。

5. 个人体会:ax 调度真正值钱的部分是什么

折腾完这一整套,我最深刻的体会是,ax 调度本质上不是某个具体的算法,而是一种把复杂问题拆成两级、然后守着每一个并发边界做控制的工程思路。

一级接入层帮你挡住流量尖峰,二级分发层帮你合理分配容量,中间的排队、限流、重试、幂等,全都是为了让系统在异常场景下也能按预期运转。你不一定需要一个完整的 ax 框架,但完全可以在现有系统里逐步引入这些设计:先加幂等键,再加背压,再加超时重投,一件一件来,系统会变得结实很多。

最后分享一个我自己调调度系统时的核心监控思路:三个指标盯死了,系统一般不会出大乱子。第一是 pending_dispatch 积压数,它反映入口承受的压力;第二是 worker_busy_rate,它反映执行端是否真正饱和;第三是 task_timeout_rate,它反映存活任务是否经常被误杀。三个指标配合起来看,能快速定位瓶颈在入口、调度器还是执行器,不用再凭感觉猜了。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询