news 2026/9/9 21:50:35

第9章:RabbitMQ 消费者可靠性——Ack、Nack、Prefetch 与 QoS

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
第9章:RabbitMQ 消费者可靠性——Ack、Nack、Prefetch 与 QoS

1. 项目背景

发布器已经 Confirm 了,短信仍会重复或丢失。客服截图里同一订单两条「支付成功」,另一单完全没短信。消费代码长这样:

basic_consume(..., auto_ack=True) send_sms(body) # 网关超时 3s # 进程被 k8s 杀掉

autoAck 的含义是:Broker 把消息交给 TCP 就算消费成功。短信还没发出,消息已经从队列消失。反过来,有人改成手动 Ack 却在finally里一律 Ack,业务失败也当成功。还有人 prefetch=500,慢网关下一堆积 500 条 unacked,Broker 内存涨,发布被 block,整条中台「假死」。

autoAck → 投递即删除(崩溃 = 丢) 手动 Ack → 处理成功才 basic.ack Nack/Reject → 失败回队列或丢掉(可走 DLX,第 10 章) prefetch → 通道上未确认的最大投递数 redelivered → 至少投递过一次(不保证恰好一次)

经典队列上还有:连接断开,未 Ack 会重新变 ready。测试必须能断言「杀消费者进程后消息还在」,而不是看日志「收到过」。

Consumer Timeout 对经典队列在 4.3 后不再按老方式评估(发布说明:经典队列与 Stream 不走这套超时)。不要用 3.x 文档的 consumer_timeout 解释本章实验。仲裁队列的 delivery-limit 第 19 章再讲。

测试同学还把「收到消息的日志行数」当成消费成功数。autoAck 下日志很多,库里短信记录很少——差的那一截就是崩溃窗口。验收必须对比:队列深度变化、短信发送表、redelivered 比例。三者对不上就重开事故单,而不是让开发改日志级别。


2. 项目设计

小胖把图书馆借书卡拍到白板上。

小胖:这不就是借书吗?管理员把书塞你手里就要在系统里划走,不然别人还以为架上有书。autoAck 多爽,为啥还要还书的时候再刷卡?

大师:塞你手里划走,你在楼梯上摔了书就没了,馆藏数字还显示「已借出处理完毕」。手动 Ack 是你坐到座位上打开书确认没缺页再划走。Nack 是你说这本书破了,要么放回架(requeue),要么进修复间(死信)。Prefetch 是一次允许你抱几本走,抱 50 本堵在走廊,别人借不到,前台也进不了新书。

技术映射:Ack=处理完成;Nack requeue=放回;prefetch=未 Ack 窗口;unacked=抱在手里的书。

小白:basic.rejectbasic.nack什么区别?multiple 标志会不会把别人的单子一起 Ack 掉?prefetch 是 Channel 级还是 Consumer 级?全局 QoS 还支持吗?手动 Ack 忘了写会怎样?redelivered能当幂等键吗?多消费者竞争同一队列如何公平?

大师:reject 一次一条;nack 可multiple且可requeue。multiple 按 delivery-tag本通道累计确认,不会 Ack 别的连接。经典 AMQP 的basic.qos在 RabbitMQ 里常用 prefetch_count;global语义历史坑多,推广中台规定按消费者设置 prefetch,禁止玩 global。源码上 limiter 进程按通道调解队列投递(rabbit_limiter.erl)。忘 Ack:消息一直 unacked,队列看起来有货但没人能拿走,内存涨。redelivered只是「曾经投出过」,网络重试也会真,不能当幂等键,幂等键是 orderId。多消费者是竞争消费,Broker 轮询投递,不保证同一订单始终同一实例——要粘滞用 SAC(第 23 章)。

小胖:那 prefetch 填 1 不就永远安全?大促吞吐怎么办?

大师:prefetch=1 延迟高、吞吐低,适合严格串行或处理很重;短信网关 200ms 时可 20~50。用实验画两条曲线,禁止拍脑袋 500。慢消费者 + 大 prefetch 是内存事故的标配。

技术映射:吞吐 ≈ 处理速率 × 窗口;窗口过大 = 把队列搬进消费者进程和 unacked 列表。

小白:崩溃重投会不会和 Nack requeue 打成死循环?毒消息怎么办?basic.recover 还要不要用?消费端 Confirm 吗?取消订阅basic.cancel时未 Ack 去哪?连接断了 exclusive 队列上的未 Ack 呢?

大师:会循环。requeue=true的毒消息会顶号。本章演示循环风险,第 10 章用死信+次数打断。测试要有「故意失败 N 次」用例,不能只测快乐路径 Ack。basic.recover让本通道未 Ack 重新投递,现代客户端少用,滚动发布靠断连即可。消费端没有 Confirm 这回事,Ack 就是消费侧回执。basic.cancel后未 Ack 回队列。exclusive 队列随连接删除,未 Ack 一起消失——这是 RPC 回调能「干净」的原因,也是不能把支付队列声明成 exclusive 的原因。

小胖:四枪:autoAck 杀进程丢消息、手动 Ack 杀进程消息还在、Nack 回去、prefetch 1 对 50 看 unacked。


3. 项目实战

3.1 环境准备

队列q.order.pay。先灌 20 条可识别 body(PAY-00…)。Python 3.11 + pika。

# promo-mq/ch09/seed.pyimportpika conn=pika.BlockingConnection(pika.ConnectionParameters("127.0.0.1",5672,"promo",pika.PlainCredentials("promo","promo_dev_2026")))ch=conn.channel()ch.confirm_delivery()foriinrange(20):ch.basic_publish("ex.order.direct","pay.ok",f"PAY-{i:02d}".encode(),properties=pika.BasicProperties(delivery_mode=2),mandatory=True)print("seeded 20")conn.close()

3.2 步骤一:autoAck 崩溃等于丢(反面)

步骤目标:自动确认下,进程在处理后、业务完成前退出,消息不再回到队列。

# promo-mq/ch09/autoack_crash.pyimportos,pikadefon_msg(ch,method,props,body):print("got",body,"autoacked already, now crash")os._exit(1)# 不关通道,模拟 kill -9conn=pika.BlockingConnection(pika.ConnectionParameters("127.0.0.1",5672,"promo",pika.PlainCredentials("promo","promo_dev_2026")))ch=conn.channel()ch.basic_qos(prefetch_count=1)ch.basic_consume("q.order.pay",on_msg,auto_ack=True)print("consuming autoack")ch.start_consuming()

跑之前记下messages。跑完再查。

运行结果:队列少 1 条,且不会因为崩溃回来。这就是短信丢失现场。

坑:os._exit才会跳过清理;conn.close()可能还来得及。测试要用硬退出。

3.3 步骤二:手动 Ack —— 崩溃后消息还在

步骤目标:收到后不 Ack 就退出,ready 恢复(可能带 redelivered)。

# promo-mq/ch09/manual_crash.pyimportos,pikadefon_msg(ch,method,props,body):print("got",body,"redelivered=",method.redelivered,"tag",method.delivery_tag)print("crash before ack")os._exit(1)conn=pika.BlockingConnection(pika.ConnectionParameters("127.0.0.1",5672,"promo",pika.PlainCredentials("promo","promo_dev_2026")))ch=conn.channel()ch.basic_qos(prefetch_count=1)ch.basic_consume("q.order.pay",on_msg,auto_ack=False)ch.start_consuming()

再启动一次正常消费者:

# promo-mq/ch09/manual_ack.pyimportpika,timedefon_msg(ch,method,props,body):print("process",body,"redelivered=",method.redelivered)time.sleep(0.05)# 假装调短信网关ch.basic_ack(method.delivery_tag)ifbody==b"PAY-19":ch.stop_consuming()conn=pika.BlockingConnection(pika.ConnectionParameters("127.0.0.1",5672,"promo",pika.PlainCredentials("promo","promo_dev_2026")))ch=conn.channel()ch.basic_qos(prefetch_count=1)ch.basic_consume("q.order.pay",on_msg,auto_ack=False)ch.start_consuming()conn.close()

运行结果:第一次崩溃后list_queues消息数不减(或 unacked 回 ready)。第二次同一 body 可能redelivered=True

坑:第二次处理必须幂等,否则短信双发——这解释了客服「两条短信」。
坑:Ack 了错的 tag 或多次 Ack 会通道异常(第 4 章 406 类)。

limiter 与 prefetch 的关系在模块头写得很清楚:

%% The purpose of the limiter is to stem the flow of messages from %% queues to channels ... AMQP 0-9-1's basic.qos prefetch_count %% Each channel has an associated limiter process

3.4 步骤三:Nack 放回 vs 丢掉

步骤目标:requeue=True会再拿到;False则消息从队列消失(无 DLX 时真正丢,第 10 章可接死信)。

# promo-mq/ch09/nack_demo.pyimportpika count={"n":0}defon_msg(ch,method,props,body):count["n"]+=1print("#",count["n"],body,"redelivered",method.redelivered)ifcount["n"]<=2:ch.basic_nack(method.delivery_tag,requeue=True)returnch.basic_ack(method.delivery_tag)ch.stop_consuming()conn=pika.BlockingConnection(pika.ConnectionParameters("127.0.0.1",5672,"promo",pika.PlainCredentials("promo","promo_dev_2026")))ch=conn.channel()ch.basic_qos(prefetch_count=1)ch.basic_consume("q.order.pay",on_msg,auto_ack=False)ch.start_consuming()conn.close()

运行结果:同一条至少打印 3 次,前两次 nack。这就是毒消息循环的缩影——生产必须有次数上限。

再开一次requeue=False(换一条新消息)后,深度减 1 且不再回来。

坑:单消费者 nack requeue 可能立刻拿回同一条,CPU 打满。多消费者时可能交给别人,问题变成随机。

3.5 步骤四:prefetch=1 vs 50

步骤目标:慢处理下观察messages_unacknowledged

先 seed 30 条到专用队列,避免打乱支付队列:

# promo-mq/ch09/prefetch_lab.pyimporttime,threading,pikadefconsume(prefetch,seconds=8):conn=pika.BlockingConnection(pika.ConnectionParameters("127.0.0.1",5672,"promo",pika.PlainCredentials("promo","promo_dev_2026")))ch=conn.channel()ch.queue_declare("q.lab.prefetch",durable=True)ch.basic_qos(prefetch_count=prefetch)defon_msg(ch,method,props,body):time.sleep(0.3)ch.basic_ack(method.delivery_tag)ch.basic_consume("q.lab.prefetch",on_msg,auto_ack=False)t0=time.time()whiletime.time()-t0<seconds:conn.process_data_events(time_limit=0.2)conn.close()# 先灌 30 条到 q.lab.prefetch 再分别跑 prefetch=1 和 50# 跑的同时:rabbitmqctl list_queues -p promo name messages messages_unacknowledged

另开终端每秒打一次:

dockerexecrabbit-promo-1 rabbitmqctl list_queues-ppromo name messages messages_unacknowledged

运行结果:prefetch=1 时 unacked 约为 1;=50 时 unacked 可冲到几十(不超过 50 且不超过剩余消息)。吞吐上 50 通常更高,直到网关或 Broker 内存成为瓶颈。

坑:在 BlockingConnection 里sleep会挡住心跳,实验 sleep 0.3s 可接受,生产 3s 同步 sleep 会掐连接。用线程池或异步。
坑:两个消费者同时消费同一队列时,prefetch 是每个通道的窗口,总 inflight 是相加关系。

值班口诀:ready>0 且 consumers=0 是「没人干活」;unacked 持续等于 prefetch 且 ready 仍涨,是「人慢或卡死」;unacked 长期等于消息总数且 ready=0,是「忘 Ack」。三种告警文案要分开,否则运维只会重启消费者,把忘 Ack 变成重复短信。

3.6 完整代码清单

column/samples/ch09/ seed.py autoack_crash.py manual_crash.py manual_ack.py nack_demo.py prefetch_lab.py

3.7 测试验证

编号操作期望
TC-CH09-01autoAck + 硬退出消息消失
TC-CH09-02手动未 Ack + 硬退出消息回 ready
TC-CH09-03重投redelivered true
TC-CH09-04nack requeue再次投递
TC-CH09-05prefetch=1unacked≤1
curl-s-upromo:promo_dev_2026\http://127.0.0.1:15672/api/queues/promo/q.lab.prefetch\|rg"messages_unacknowledged|messages_ready"

值班检查单:消费路径发布列车增加崩溃注入:杀掉消费 pod,断言支付队列深度不减少(手动 Ack)或明确记录「允许丢失」(仅非关键通知且书面批准)。看到重复短信先查幂等表,不要先怪 Broker。prefetch 配置必须进配置中心,禁止写死 500。unacked 告警阈值建议设为prefetch × 消费者数的 80%,持续五分钟即叫人,避免拖到内存告警才发现忘 Ack。

basic.get不受 QoS 限制,管理面「Get messages」同样会改变队列。测试与值班禁止在生产支付队列上点 Get。需要采样时复制到旁路队列或用 Tracing(第 28 章)短时打开。

消费侧还有一个组织问题:同一个队列挂了短信和「写发送记录」两个逻辑在一个回调里。短信成功但写库失败时,Ack 会丢记录,Nack 会再发短信。正确拆法是:本地事务先写「发送中」,再调网关,再更新「成功」,最后 Ack;失败则走第 10 章重试,而不是在回调里既想恰好一次又想随便 Nack。幂等表的主键建议orderId + channel(短信/邮件),不要只用 orderId,否则邮件失败会挡住短信重试。

prefetch 调参实验至少记录四列:prefetch、处理耗时、吞吐、Broker unacked。缺一列就会在评审里变成「感觉 50 比较快」。把表贴进 Wiki,第 30 章压测时作为消费侧基线,避免到了大促才把窗口从 1 改到 500。

手动 Ack 的代码审查清单可以短到三行:回调里有没有业务失败分支;失败分支有没有 Nack 或走死信而不是 Ack;成功路径是不是最后一行才 Ack。很多事故出在「日志打了成功、异常在 Ack 之后」。把 Ack 放在函数最后并用早返回处理失败,能少掉一半误 Ack。再配上集成测试杀进程,消费契约才算闭合。

滚动发布时旧消费者断连,未 Ack 会回到队列并可能带上 redelivered。新实例必须能处理「半截网关调用」:网关已成功但未 Ack 的,靠幂等跳过;网关未调用的,正常发送。这要求发送记录在调用网关之前就写入「进行中」,而不是全部成功后再写。顺序写错,滚动当天必双发。把这条写进消费脚手架 README,比口口相传可靠。评审时打开 README 对一下顺序,比只看有没有 basic_ack 更能发现双发隐患。顺序错了,再漂亮的 Ack 也救不了客服电话。把「先写进行中、再调网关、再 Ack」印成三人桌贴。小胖负责贴,小白负责抽查代码顺序,大师负责卡住不按顺序的合并。


4. 项目总结

优点与缺点

策略优点缺点
autoAck代码少、快崩溃即丢
手动 Ack可对齐业务成功忘 Ack 会堵死
Nack requeue暂时故障可恢复毒消息死循环
prefetch 大吞吐高unacked 吃内存
prefetch=1简单、压力平滑延迟差

优点:1)语义可测。2)redelivered 提示至少一次。3)limiter 把 QoS 从 Channel 抽出去避免打爆 channel 进程。
缺点:1)至少一次 ≠ 恰好一次。2)经典队列无 delivery-limit。3)sleep 式消费害心跳。

消费侧口诀:先做事,再 Ack;失败就分类(再试 / 死信 / 丢);窗口按慢速环节设,而不是按 QPS 设。QPS 是结果,prefetch 是约束。用 QPS 反推窗口可以,用「感觉卡」直接把窗口加到 500 不行。

适用场景

  • 短信/邮件等必须手动 Ack。
  • CPU 很重或要严格串行时 prefetch=1。
  • 网关稳定时适度加大窗口。
  • 崩溃注入与重复投递的测试训练。

不适用:autoAck 用于支付;用 redelivered 当去重 ID;用 Nack 循环当重试退避(第 10 章 TTL)。

注意事项

  • 4.3 经典队列不再按旧 consumer_timeout 那套评估。
  • 多线程不要共享 Channel 去 Ack。
  • 安全:消费者账号只要 read,不要 configure 删队列。
  • basic.get不受 prefetch 限制(limiter 注释写明),监控脚本乱 get 会捣乱。

常见踩坑(生产)

  1. autoAck + 网关超时,K8s 杀 pod,短信丢失。根因:投递即删除。
  2. prefetch=1000,消费者 Full GC,unacked 占满内存触发告警。根因:窗口当缓冲。
  3. finally 里一律 Ack,业务失败也消消息。根因:把「通道还活着」当「业务成功」。

思考题

  1. 两个消费者 prefetch 各 50,队列 10 条,unacked 最大可能多少?会不会一个吃完另一个饿死?
  2. 手动 Ack 成功后应用仍崩溃,用户已看到短信,补偿流程应靠什么,而不是 Nack?

附录 C:第 8 章思考题参考答案

题 1:Confirm 后 kill -9。
经典队列仍可能丢尾部。quorum 多数派提交后更稳。发布器不能承诺「Confirm=永存」。

题 2:SENT 后用户再点支付。
靠订单状态机与短信发送记录幂等,不靠 MQ 去重。消费者 Ack 只表示这一次投递处理完。

延伸阅读与资源

SQLAlchemy 2.0从入门到进阶的实战之旅
Dify 从入门到进阶:LLM 应用平台实战修炼
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/9 21:48:58

C++ STL容器详解:stack、queue与deque的底层原理及实战应用

C里最容易上手、也最容易被误用的容器&#xff0c;我觉得就是这三个&#xff1a;stack、queue、deque。说它们容易上手&#xff0c;是因为接口少到可以两分钟全记住&#xff1b;说它们容易被误用&#xff0c;是因为很多人不清楚deque到底是干什么的&#xff0c;也不知道stack和…

作者头像 李华
网站建设 2026/9/9 21:48:49

如何配置 MinIO 桶事件通知发布到 Apache Kafka 主题?

如何配置 MinIO 桶事件通知发布到 Apache Kafka 主题&#xff1f; 【免费下载链接】minio MinIO is a high-performance, S3 compatible object store, open sourced under GNU AGPLv3 license. 项目地址: https://gitcode.com/GitHub_Trending/mi/minio 如果你的应用需…

作者头像 李华
网站建设 2026/9/9 21:48:20

Matlab GUI数字均衡器实战:从滤波器设计到实时音频处理

简介&#xff1a;面向音频处理与数字信号处理初学者的Matlab GUI数字均衡器设计资源&#xff0c;以IIR滤波器为核心&#xff0c;解决如何在图形界面中直观调整音频频响、完成均衡处理的问题&#xff0c;涵盖滑动条、按钮、文本框等常见控件设计思路&#xff0c;适合课程设计、毕…

作者头像 李华
网站建设 2026/9/9 21:48:19

EC20 4G模块TCP透传模式配置全流程:从AT指令到串口数据桥接实战

简介&#xff1a;面向STM32F4系列开发者的EC20模块TCP透传模式通信工程示例&#xff0c;解决微控制器通过AT指令建立套接字连接并透明收发数据的常见需求。资源适合有一定嵌入式基础、正在调试无线通信模块的开发者&#xff0c;也适合作为物联网终端联网功能的学习参考。压缩包…

作者头像 李华
网站建设 2026/9/9 21:48:12

基于STM32F103ZET6的示波器设计:从ADC采样到波形显示全解析

简介&#xff1a;基于STM32F103ZET6的简易示波器程序包&#xff0c;面向单片机学习者和嵌入式开发入门者&#xff0c;演示如何利用Cortex-M3内核芯片的ADC采集、定时器控制与LCD显示实现正弦波、方波等波形可视化。配套工程完整&#xff0c;可直接用于学习信号采集、数据处理、…

作者头像 李华