news 2026/9/5 7:33:37

第18章:Celery 预取、晚确认与可见性超时

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
第18章:Celery 预取、晚确认与可见性超时

0. 上一章思考题参考答案

思考题 1:prefork 下time.sleep(0.3)真阻塞——子进程进入内核睡眠,进程模型下一个子进程同一时刻只能跑一个任务,槽位被占死;gevent 下time.sleep已被 monkey patch 成协程调度器的让出点——协程睡眠时调度器切到其他可运行的协程,进程不闲着。本质区别:一个是「操作系统线程/进程级睡眠」,一个是「用户态协程调度」,这也是两种池子吞吐差异的根源。

思考题 2:一个 Worker 一个池子模型,混订两类任务必然互相拖累(CPU 任务占住子进程、IO 任务排队)。解法:按队列拆两个 Worker——Worker-A -P prefork -Q cpu_queueWorker-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短任务,可丢可重4False300s吞吐优先
order(扣库存)关键写,绝不可丢1True600s可靠性优先
report长任务(分钟级)1True3600s防重复防丢失
pay支付回调,强幂等2True600s晚确认+幂等双保险

测试验证:

# 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==4
python-mpytest tests/test_reliability.py-v# 3 passed

4. 项目总结

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 个生产故障)

  1. 故障:大促加 Worker 后吞吐不升反降。根因:预取倍数 4,新 Worker 抢不到消息,集群空转。对策:短任务倍数调 2 + 长任务倍数调 1。教训:加机器前先看预取水位
  2. 故障:扣库存任务执行中服务重启,库存「凭空少了一件」。根因:早确认 + 执行中崩溃,消息已 ack 丢失。对策:关键任务改 acks_late + 幂等。教训:早确认的「省事」是用「静默丢失」买单
  3. 故障:同一订单被扣两次,客诉。根因:visibility_timeout=60s,任务执行 90s,超时重投;无幂等兜底。对策:超时调 300s + 去重表。教训:长任务在 Redis Broker 上的「双重奏」,幂等是唯一休止符

4.5 思考题

  1. acks_late=Truevisibility_timeout=600,Worker 执行到 700 秒时崩溃——消息最终会怎样?(提示:确认了吗?重新可见吗?)
  2. 预取倍数 = 并发 × 倍数。如果-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 实战修炼与源码剖析

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/5 7:33:37

AI编程如何评审达标?从订单代码看星级代码的五大标准

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/5 7:28:51

MyEMS接入Modbus TCP智能电表全流程:从配置到排障的实践指南

1. 从项目背景聊起:为什么偏偏是 MyEMS 加 Modbus TCP先交代一下我这边的情况。手头负责的厂区有好几栋楼,配电房里电表品牌杂得很,有施耐德、有正泰、有安科瑞,还有一些国产小厂出的表。以前想统一看数据,只能每个品牌…

作者头像 李华
网站建设 2026/9/5 7:26:44

PCB艺术创作指南:用EDA软件两小时完成电路板画设计

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/5 7:25:15

Linux内核页表深度解析:多级结构、TLB与缺页异常全链路

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华