0. 上一章思考题参考答案
思考题 1:prefork 下time.sleep(0.3)是真阻塞——子进程进入内核睡眠,进程模型下一个子进程同一时刻只能跑一个任务,槽位被占死;gevent 下time.sleep已被 monkey patch 成协程调度器的让出点——协程睡眠时调度器切到其他可运行的协程,进程不闲着。本质区别:一个是「操作系统线程/进程级睡眠」,一个是「用户态协程调度」,这也是两种池子吞吐差异的根源。
思考题 2:一个 Worker 一个池子模型,混订两类任务必然互相拖累(CPU 任务占住子进程、IO 任务排队)。解法:按队列拆两个 Worker——Worker-A -P prefork -Q cpu_queue、Worker-B -P gevent -Q io_queue,各池各队(第 9 章多 Worker 隔离 + 本章池子选型组合)。拆之前先在契约表里给每条队列标注「池子类型」。
1. 项目背景
大促压测时,小周遇到了「薛定谔的库存」:同一条deduct_stock任务,有时候扣一次,有时候扣两次,查日志发现两个 Worker 都打印了「首次扣减成功」——两条消息、同一个 task_id。运维补充了另一个怪象:8 个 Worker 都在跑,但inspect reserved显示每个 Worker 手里都攥着几十条没执行的任务,新加的第 9 个 Worker 启动后前 5 分钟几乎接不到任务。
这两个问题背后是消息可靠性最精密的三个旋钮:预取(prefetch)、晚确认(acks_late)、可见性超时(visibility_timeout)。它们单独看都不难,组合起来却像「刹车、离合、油门」——踩错了方向,要么丢任务,要么重复执行,要么饿死同伴。
三个旋钮的因果链 worker_prefetch_multiplier=4 → 每个 Worker 一次拉 并发×4 条 → 拉多了饿死别的 Worker task_acks_late=True → 执行完才确认 → 崩溃可重投 → 必须幂等 visibility_timeout → 超时未确认视为死亡 → 重新可见 → 长任务易重复本章目标:用「模拟扣库存中途被 kill」的对照实验,把三个旋钮的代价与收益逐一验证,最后输出队列级推荐配置表——以后每条队列该开什么,查表说话。
2. 项目设计
场景:库存重复扣的复盘会,小周把两个 Worker 的日志并排摆开。
小胖:为啥要「预取」?一次拿一条,干完再拿,多公平!还省得加第 9 个 Worker 接不到活。
小白:我理解预取是为了减少「来回取消息」的网络往返——一次拉一批到本地,省去每条消息一次的 RTT。但我算了一下:-c 4×prefetch_multiplier 4= 每次预取 16 条,8 个 Worker 就是 128 条「在路上」;如果任务都是 10 秒的长活,新 Worker 确实要等这 128 条消化完才有活干。预取和公平性怎么平衡?
大师:这就是经典矛盾。worker_prefetch_multiplier(默认 4,celery/app/defaults.py:357)的意思是「并发 × 倍数」一次性预取进本地内存队列。预取大 → 省 RTT、吞吐高,但饿死新 Worker、故障时消息跟着进程死(内存里没确认的消息,进程崩了就没了)。三条经验法则:① 长任务(秒级)把倍数调成 1(预取 = 并发,避免占坑);② 短任务(毫秒级)可以 4~8(吞吐优先);③ 任务时长越不均匀,预取越要小(防止一个 Worker 预取到一堆长任务拖死其他 Worker)。
技术映射:预取 = 食堂窗口一次端走 4 份盒饭到备餐台——端多了后厨(其他窗口)没菜炒;端少了来回跑(RTT)浪费时间。长菜(任务)一次只能端 1 份。
小白:那acks_late呢?第 16 章思考题里讨论了「崩溃代价 vs 重投代价」,我还想再确认一个点:早确认(默认)时 Worker 执行中崩溃,消息会怎样?晚确认时又会怎样?对应到库存任务,分别是什么后果?
大师:把两个场景画清楚:
| 场景 | 早确认(默认) | 晚确认(acks_late=True) |
|---|---|---|
| 取到消息即 ack? | 是(立即确认) | 否(执行完才确认) |
| 执行中 Worker 崩溃 | 消息已确认 →丢失,不重投 | 消息未确认 →Broker 重投 |
| 对应后果 | 库存「漏扣」(静默丢任务) | 库存「可能重复扣」(必须幂等) |
| 适用 | 可容忍丢失的通知类 | 绝不能丢的关键写操作 |
第 11 章做过幂等,所以库存任务开acks_late=True是「丢了最惨、重复有兜底」的理性选择。再补两个配套旋钮:task_reject_on_worker_lost(默认开启)——子进程异常退出(非正常 return)时把消息拒收回队列,配合晚确认才有效;task_acks_on_failure_or_timeout(默认 False)——任务失败/超时时是否照常 ack,打开可以避免「坏消息无限重投循环」。
小胖:那 Redis 的可见性超时呢?我上次配了个visibility_timeout=60,结果有个任务跑了 90 秒,被两个 Worker 各执行了一遍。这跟 acks_late 是两回事吗?
大师:是两回事但会叠加。可见性超时是 Redis Broker「模拟 ack」的机制(第 7 章讲过):消息被取走后进入「不可见」状态,N 秒内 Worker 没确认,Redis 把它重新变回可见 → 另一个 Worker 又能取到。它和 acks_late 的关系:晚确认 + 可见性超时= 执行中崩溃要靠「超时窗口」重投;早确认 + 可见性超时= 取走即确认,崩溃后消息已「不可见但未确认」?——不,早确认时 Kombu 直接删除消息,不存在重投。所以 Redis 上晚确认的重投延迟 ≈ 可见性超时剩余时间,超时设太短,长任务被误判「死亡」重复投递;设太长,真崩溃时重投变慢。规则:可见性超时 ≥ 任务最长执行时间 × 2。
技术映射:可见性超时 = 「外卖超时未送达自动重新派单」的计时器——菜品(消息)在骑手(Worker)手里超过 N 分钟没标记送达,平台就派新骑手;菜本来就慢(长任务),超时设太短就会一菜两骑手(重复执行)。
3. 项目实战
3.1 环境准备
沿用环境(Redis Broker + Backend)。本章用「扣库存」任务做对照实验,需要两个 Worker 终端。
3.2 分步实现
步骤 1:定义扣库存任务(可模拟执行中被杀)
目标:任务里留一个「执行中观察窗口」,方便手动 kill。
# reliability_tasks.pyimporttime,sqlite3fromceleryimportCelery app=Celery('reliability',broker='redis://localhost:6379/0',backend='redis://localhost:6379/1')_DB="dedup.db"defdeduct_once(order_id):conn=sqlite3.connect(_DB)conn.execute("CREATE TABLE IF NOT EXISTS log(order_id INTEGER PRIMARY KEY)")try:conn.execute("INSERT INTO log VALUES (?)",(order_id,))conn.commit();returnTrueexceptsqlite3.IntegrityError:returnFalsefinally:conn.close()@app.task(name='rel.deduct',bind=True,acks_late=True,max_retries=2)defdeduct(self,order_id:int,sleep:float=10.0)->str:time.sleep(sleep)# 执行窗口:留时间给人工 killok=deduct_once(order_id)print(f"[deduct]{order_id}->{ok}")return"deducted"ifokelse"skipped"步骤 2:实验 A——早确认 vs 晚确认(Worker 崩溃对比)
目标:亲手验证「早确认丢任务、晚确认重投」的差异。
# 终端 A:早确认 Workercelery-Areliability_tasks worker-c2--loglevel=info-nearly --without-gossip --without-mingle --without-heartbeat# 终端 B:晚确认 Workercelery-Areliability_tasks worker-c2--loglevel=info-nlate --without-gossip --without-mingle --without-heartbeat# 实验 A-1:早确认场景——投递后 3 秒 kill Workercelery-Areliability_tasks call rel.deduct--args='[101, 10]'# 立刻 kill 终端 A 的 Worker 进程(模拟崩溃)# 重启 Worker,观察 101 是否被再次执行运行结果(文字描述):早确认 Worker 被 kill 时消息已 ack——重启后 101 任务消失,dedup 表里没有 101(任务静默丢失);晚确认 Worker 同样被 kill,重启后101 被重新投递并执行(dedup 表出现 101)——但如果没有幂等,第二次执行就会重复扣。一丢一重,这就是两个旋钮的代价,二选一是逃不掉的。
步骤 3:实验 B——可见性超时引发的「一菜两骑手」
目标:复现「任务时长 > 可见性超时」的重复投递。
# 配置:visibility_timeout=5(故意设小)# 任务 rel.deduct sleep=10 > 5 → 执行到一半消息重新可见# 两个 Worker(同队列,不分早/晚确认)celery-Areliability_tasks worker-c1--loglevel=info-nw1 celery-Areliability_tasks worker-c1--loglevel=info-nw2 celery-Areliability_tasks call rel.deduct--args='[202, 10]'运行结果(文字描述):w1 取到 202 开始执行;约 5 秒后(可见性超时到期)w2 也收到 202 并开始执行——同一个 order_id 被两个 Worker 同时扣。dedup 表靠唯一键兜底只留一条,但如果任务没有幂等,这就是真实世界的「重复扣款」事故。
步骤 4:实验 C——预取与「新 Worker 饿死」
目标:验证预取对公平性的影响。
# 终端 A:预取倍数 8(-c 4 默认乘 4=16 条)celery-Areliability_tasks worker-c4--loglevel=info-npf-big# 灌 100 条 5 秒任务后,再启动终端 B(新 Worker)celery-Areliability_tasks worker-c4--loglevel=info-npf-new# 观察 pf-new 的收到第一条任务的耗时运行结果(文字描述):pf-big 一次性预取 16×N 条,pf-new 启动后5~30 秒收不到任何任务(消息都被老 Worker 预取光了);把倍数调成 1(--prefetch-multiplier 1)重跑,pf-new 几秒内开始接活。结论:预取倍数越大,集群「新成员」的饥饿期越长。
三旋钮总结:晚确认(防丢)+ 幂等(防重)+ 可见性超时(防长任务误判)——三者必须成组设计、成组评审,单独调任何一个都会把代价转嫁给另外两个(第 4.3 节注意事项的落地案例)。
3.3 可能遇到的坑及解决方法
| 坑 | 现象 | 解决 |
|---|---|---|
| 长任务被「当成死亡」重复执行 | 任务时长 > visibility_timeout | 可见性超时 ≥ 最长任务 × 2;或换 RabbitMQ |
| 开了 acks_late 后重复执行 | 崩溃重投 + 无幂等 | 幂等键必须与 acks_late 同生命周期评审 |
| 新 Worker 半天不干活 | 预取倍数过大 | --prefetch-multiplier 1(长任务)或按队列调 |
| 消息「凭空消失」 | 早确认 + Worker 崩溃 | 评估任务丢失代价,关键任务改晚确认 |
| reject_on_worker_lost 不生效 | 子进程正常 return 不算异常 | 该旋钮只管「异常退出」(被 kill/异常)场景 |
3.4 完整代码清单与测试验证
清单:reliability_tasks.py+ 三个实验的启动与 kill 步骤。队列级推荐配置表(沉淀 Wiki,第 9 章队列规划表的配套):
| 队列 | 任务特征 | 预取倍数 | acks_late | 可见性超时 | 依据 |
|---|---|---|---|---|---|
| sms | 短任务,可丢可重 | 4 | False | 300s | 吞吐优先 |
| order(扣库存) | 关键写,绝不可丢 | 1 | True | 600s | 可靠性优先 |
| report | 长任务(分钟级) | 1 | True | 3600s | 防重复防丢失 |
| pay | 支付回调,强幂等 | 2 | True | 600s | 晚确认+幂等双保险 |
测试验证:
# tests/test_reliability.pyfromreliability_tasksimportapp,deduct,deduct_once app.conf.task_always_eager=Truedeftest_late_ack_configured():assertdeduct.acks_lateisTrueassertdeduct.max_retries==2deftest_dedup_guard_prevents_double():assertdeduct_once(1001)isTrueassertdeduct_once(1001)isFalsedeftest_prefetch_multiplier_default_is_4():assertapp.conf.worker_prefetch_multiplier==4python-mpytest tests/test_reliability.py-v# 3 passed4. 项目总结
4.1 优点 & 缺点
| 旋钮 | 开 | 关 | 代价 |
|---|---|---|---|
| acks_late=True | 崩溃可重投(不丢) | 崩溃即丢 | 重复执行风险 ↑ → 必须幂等 |
| 预取倍数调 1 | 公平、新 Worker 即时接活 | 吞吐略降 | RTT 增加 |
| 预取倍数调大 | 吞吐高 | 饥饿其他 Worker | 故障窗口内消息随进程丢失 |
| visibility_timeout 调大 | 长任务安全 | 真崩溃时重投变慢 | 失败感知延迟 |
4.2 适用场景
- 适用:① 关键写操作(库存/金额)→ acks_late + 幂等 + 小预取;② 长任务(报表/导出)→ 小预取 + 大可见性超时;③ 通知类短任务 → 大预取高吞吐。
- 不适用:① 幂等做不了的资源型任务(如物理信号),别用晚确认;② 消息量极小且必须顺序消费的场景,预取/ack 配置意义有限,优先保证单消费者;③ 尚未做幂等的存量任务,先补幂等再开 acks_late——顺序反了就是给生产埋重复执行的地雷。
4.3 注意事项
- 三个旋钮必须成组评审:acks_late、幂等键、可见性超时是同一决策的三面,别单改一个。
- 可见性超时按「排队等待时间 + 执行时间」计算(第 7 章故障 3 的教训),不是只看执行时间。
- 修改预取倍数后要重新压测吞吐,防止「公平了但慢了一半」。
- Redis 与 RabbitMQ 的确认语义不同(超时重投 vs 显式 Ack),迁移 Broker 时配置表要整体重评。
4.4 常见踩坑经验(3 个生产故障)
- 故障:大促加 Worker 后吞吐不升反降。根因:预取倍数 4,新 Worker 抢不到消息,集群空转。对策:短任务倍数调 2 + 长任务倍数调 1。教训:加机器前先看预取水位。
- 故障:扣库存任务执行中服务重启,库存「凭空少了一件」。根因:早确认 + 执行中崩溃,消息已 ack 丢失。对策:关键任务改 acks_late + 幂等。教训:早确认的「省事」是用「静默丢失」买单。
- 故障:同一订单被扣两次,客诉。根因:visibility_timeout=60s,任务执行 90s,超时重投;无幂等兜底。对策:超时调 300s + 去重表。教训:长任务在 Redis Broker 上的「双重奏」,幂等是唯一休止符。
4.5 思考题
acks_late=True且visibility_timeout=600,Worker 执行到 700 秒时崩溃——消息最终会怎样?(提示:确认了吗?重新可见吗?)- 预取倍数 = 并发 × 倍数。如果
-c 8、倍数 8(预取 64 条),其中 63 条是 10 秒长任务——第 9 个 Worker 要等多久?这个场景推荐什么配置?
答案见第 19 章开头的「上一章思考题参考答案」。
延伸阅读与资源
Java 工程师进阶:从 JVM 生产排障到OpenJDK原理
NumPy 从入门到生产落地:全链路实战指南(科学计算/向量化)
Redis 8 实战精讲:从 CRUD 到源码,构建高可用缓存系统
Redis 实战修炼与原理进阶
Python 3实战精进:从脚本到高并发订单引擎
python入门:Rquests从菜鸟脚本到企业级SDK的网络实战圣经
Milvus向量数据库实战修炼:从 0 到 1精通向量检索与生产落地
MongoDB 实战进阶与内核修炼
后端工程师的 AI 转型第一课:Ollama 与私有化大模型实战
10倍开发者的 Dify 魔法书:从零构建全栈 AI 应用
后端工程师转型AI第一课-Ollama 与私有化大模型实战
大型语言模型(LLM) vLLM 高性能推理落地实战
Agent开发之LlamaIndex 实战修炼与源码进阶
大语言模型Transformers 实战修炼与源码剖析