第18章:Celery 预取、晚确认与可见性超时
0. 上一章思考题参考答案思考题 1prefork 下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_multiplier4 → 每个 Worker 一次拉 并发×4 条 → 拉多了饿死别的 Worker task_acks_lateTrue → 执行完才确认 → 崩溃可重投 → 必须幂等 visibility_timeout → 超时未确认视为死亡 → 重新可见 → 长任务易重复本章目标用「模拟扣库存中途被 kill」的对照实验把三个旋钮的代价与收益逐一验证最后输出队列级推荐配置表——以后每条队列该开什么查表说话。2. 项目设计场景库存重复扣的复盘会小周把两个 Worker 的日志并排摆开。小胖为啥要「预取」一次拿一条干完再拿多公平还省得加第 9 个 Worker 接不到活。小白我理解预取是为了减少「来回取消息」的网络往返——一次拉一批到本地省去每条消息一次的 RTT。但我算了一下-c 4×prefetch_multiplier 4 每次预取 16 条8 个 Worker 就是 128 条「在路上」如果任务都是 10 秒的长活新 Worker 确实要等这 128 条消化完才有活干。预取和公平性怎么平衡大师这就是经典矛盾。worker_prefetch_multiplier默认 4celery/app/defaults.py:357的意思是「并发 × 倍数」一次性预取进本地内存队列。预取大 → 省 RTT、吞吐高但饿死新 Worker、故障时消息跟着进程死内存里没确认的消息进程崩了就没了。三条经验法则① 长任务秒级把倍数调成 1预取 并发避免占坑② 短任务毫秒级可以 4~8吞吐优先③ 任务时长越不均匀预取越要小防止一个 Worker 预取到一堆长任务拖死其他 Worker。技术映射预取 食堂窗口一次端走 4 份盒饭到备餐台——端多了后厨其他窗口没菜炒端少了来回跑RTT浪费时间。长菜任务一次只能端 1 份。小白那acks_late呢第 16 章思考题里讨论了「崩溃代价 vs 重投代价」我还想再确认一个点早确认默认时 Worker 执行中崩溃消息会怎样晚确认时又会怎样对应到库存任务分别是什么后果大师把两个场景画清楚场景早确认默认晚确认acks_lateTrue取到消息即 ack是立即确认否执行完才确认执行中 Worker 崩溃消息已确认 →丢失不重投消息未确认 →Broker 重投对应后果库存「漏扣」静默丢任务库存「可能重复扣」必须幂等适用可容忍丢失的通知类绝不能丢的关键写操作第 11 章做过幂等所以库存任务开acks_lateTrue是「丢了最惨、重复有兜底」的理性选择。再补两个配套旋钮task_reject_on_worker_lost默认开启——子进程异常退出非正常 return时把消息拒收回队列配合晚确认才有效task_acks_on_failure_or_timeout默认 False——任务失败/超时时是否照常 ack打开可以避免「坏消息无限重投循环」。小胖那 Redis 的可见性超时呢我上次配了个visibility_timeout60结果有个任务跑了 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 appCelery(reliability,brokerredis://localhost:6379/0,backendredis://localhost:6379/1)_DBdedup.dbdefdeduct_once(order_id):connsqlite3.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(namerel.deduct,bindTrue,acks_lateTrue,max_retries2)defdeduct(self,order_id:int,sleep:float10.0)-str:time.sleep(sleep)# 执行窗口留时间给人工 killokdeduct_once(order_id)print(f[deduct]{order_id}-{ok})returndeductedifokelseskipped步骤 2实验 A——早确认 vs 晚确认Worker 崩溃对比目标亲手验证「早确认丢任务、晚确认重投」的差异。# 终端 A早确认 Workercelery-Areliability_tasks worker-c2--loglevelinfo-nearly --without-gossip --without-mingle --without-heartbeat# 终端 B晚确认 Workercelery-Areliability_tasks worker-c2--loglevelinfo-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_timeout5故意设小# 任务 rel.deduct sleep10 5 → 执行到一半消息重新可见# 两个 Worker同队列不分早/晚确认celery-Areliability_tasks worker-c1--loglevelinfo-nw1 celery-Areliability_tasks worker-c1--loglevelinfo-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 默认乘 416 条celery-Areliability_tasks worker-c4--loglevelinfo-npf-big# 灌 100 条 5 秒任务后再启动终端 B新 Workercelery-Areliability_tasks worker-c4--loglevelinfo-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短任务可丢可重4False300s吞吐优先order扣库存关键写绝不可丢1True600s可靠性优先report长任务分钟级1True3600s防重复防丢失pay支付回调强幂等2True600s晚确认幂等双保险测试验证# tests/test_reliability.pyfromreliability_tasksimportapp,deduct,deduct_once app.conf.task_always_eagerTruedeftest_late_ack_configured():assertdeduct.acks_lateisTrueassertdeduct.max_retries2deftest_dedup_guard_prevents_double():assertdeduct_once(1001)isTrueassertdeduct_once(1001)isFalsedeftest_prefetch_multiplier_default_is_4():assertapp.conf.worker_prefetch_multiplier4python-mpytest tests/test_reliability.py-v# 3 passed4. 项目总结4.1 优点 缺点旋钮开关代价acks_lateTrue崩溃可重投不丢崩溃即丢重复执行风险 ↑ → 必须幂等预取倍数调 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_timeout60s任务执行 90s超时重投无幂等兜底。对策超时调 300s 去重表。教训长任务在 Redis Broker 上的「双重奏」幂等是唯一休止符。4.5 思考题acks_lateTrue且visibility_timeout600Worker 执行到 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 实战修炼与源码剖析