ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

ax调度两级模型实战:队列限流重试与幂等设计全解析

ax调度两级模型实战:队列限流重试与幂等设计全解析 最近好几个群里都在聊 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 里 BLPOPRedis 单实例下所有请求都打到同一个 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_HOSTredis - REDIS_PORT6379 worker1: image: python:3.10-slim volumes: - ./worker:/app working_dir: /app command: python worker.py worker1 environment: - REDIS_HOSTredis worker2: image: python:3.10-slim volumes: - ./worker:/app working_dir: /app command: python worker.py worker2 environment: - REDIS_HOSTredis这里 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, prioritylow, delay0, max_retry3): 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_size20): now time.time() due_tasks self.r.zrangebyscore(ax:zset:schedule, 0, now, start0, numbatch_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], timeout1) 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 每取完一个就要重新调用一次 brpopbrpop 的 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 tokenax 调度在分发任务时会用分布式锁保证同一时间只有一个调度器在处理同一个任务。这个锁本身能挡住多数并发问题但遇到长任务还是可能翻车。我遇到过一次典型故障执行器处理一个任务花了 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它反映存活任务是否经常被误杀。三个指标配合起来看能快速定位瓶颈在入口、调度器还是执行器不用再凭感觉猜了。
返回列表