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.reject 和 basic.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.py
import pika
conn = pika.BlockingConnection(pika.ConnectionParameters(
"127.0.0.1", 5672, "promo", pika.PlainCredentials("promo", "promo_dev_2026")))
ch = conn.channel()
ch.confirm_delivery()
for i in range(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.py
import os, pika
def on_msg(ch, method, props, body):
print("got", body, "autoacked already, now crash")
os._exit(1) # 不关通道,模拟 kill -9
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=True)
print("consuming autoack")
ch.start_consuming()
跑之前记下 messages。跑完再查。
运行结果: 队列少 1 条,且 不会 因为崩溃回来。这就是短信丢失现场。
坑: os._exit 才会跳过清理;conn.close() 可能还来得及。测试要用硬退出。
3.3 步骤二:手动 Ack —— 崩溃后消息还在
步骤目标: 收到后不 Ack 就退出,ready 恢复(可能带 redelivered)。
# promo-mq/ch09/manual_crash.py
import os, pika
def on_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.py
import pika, time
def on_msg(ch, method, props, body):
print("process", body, "redelivered=", method.redelivered)
time.sleep(0.05) # 假装调短信网关
ch.basic_ack(method.delivery_tag)
if body == 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.py
import pika
count = {"n": 0}
def on_msg(ch, method, props, body):
count["n"] += 1
print("#", count["n"], body, "redelivered", method.redelivered)
if count["n"] <= 2:
ch.basic_nack(method.delivery_tag, requeue=True)
return
ch.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.py
import time, threading, pika
def consume(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)
def on_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()
while time.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
另开终端每秒打一次:
docker exec rabbit-promo-1 rabbitmqctl list_queues -p promo 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-01 | autoAck + 硬退出 | 消息消失 |
| TC-CH09-02 | 手动未 Ack + 硬退出 | 消息回 ready |
| TC-CH09-03 | 重投 | redelivered true |
| TC-CH09-04 | nack requeue | 再次投递 |
| TC-CH09-05 | prefetch=1 | unacked≤1 |
curl -s -u promo: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 会捣乱。
常见踩坑(生产)
思考题
附录 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 实战修炼与源码剖析



