第9章:RabbitMQ 消费者可靠性——Ack、Nack、Prefetch 与 QoS 1. 项目背景发布器已经 Confirm 了短信仍会重复或丢失。客服截图里同一订单两条「支付成功」另一单完全没短信。消费代码长这样basic_consume(..., auto_ackTrue) send_sms(body) # 网关超时 3s # 进程被 k8s 杀掉autoAck 的含义是Broker 把消息交给 TCP 就算消费成功。短信还没发出消息已经从队列消失。反过来有人改成手动 Ack 却在finally里一律 Ack业务失败也当成功。还有人 prefetch500慢网关下一堆积 500 条 unackedBroker 内存涨发布被 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_countglobal语义历史坑多推广中台规定按消费者设置 prefetch禁止玩 global。源码上 limiter 进程按通道调解队列投递rabbit_limiter.erl。忘 Ack消息一直 unacked队列看起来有货但没人能拿走内存涨。redelivered只是「曾经投出过」网络重试也会真不能当幂等键幂等键是 orderId。多消费者是竞争消费Broker 轮询投递不保证同一订单始终同一实例——要粘滞用 SAC第 23 章。小胖那 prefetch 填 1 不就永远安全大促吞吐怎么办大师prefetch1 延迟高、吞吐低适合严格串行或处理很重短信网关 200ms 时可 2050。用实验画两条曲线禁止拍脑袋 500。慢消费者 大 prefetch 是内存事故的标配。技术映射吞吐 ≈ 处理速率 × 窗口窗口过大 把队列搬进消费者进程和 unacked 列表。小白崩溃重投会不会和 Nack requeue 打成死循环毒消息怎么办basic.recover 还要不要用消费端 Confirm 吗取消订阅basic.cancel时未 Ack 去哪连接断了 exclusive 队列上的未 Ack 呢大师会循环。requeuetrue的毒消息会顶号。本章演示循环风险第 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 条可识别 bodyPAY-00…。Python 3.11 pika。# promo-mq/ch09/seed.pyimportpika connpika.BlockingConnection(pika.ConnectionParameters(127.0.0.1,5672,promo,pika.PlainCredentials(promo,promo_dev_2026)))chconn.channel()ch.confirm_delivery()foriinrange(20):ch.basic_publish(ex.order.direct,pay.ok,fPAY-{i:02d}.encode(),propertiespika.BasicProperties(delivery_mode2),mandatoryTrue)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 -9connpika.BlockingConnection(pika.ConnectionParameters(127.0.0.1,5672,promo,pika.PlainCredentials(promo,promo_dev_2026)))chconn.channel()ch.basic_qos(prefetch_count1)ch.basic_consume(q.order.pay,on_msg,auto_ackTrue)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)connpika.BlockingConnection(pika.ConnectionParameters(127.0.0.1,5672,promo,pika.PlainCredentials(promo,promo_dev_2026)))chconn.channel()ch.basic_qos(prefetch_count1)ch.basic_consume(q.order.pay,on_msg,auto_ackFalse)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)ifbodybPAY-19:ch.stop_consuming()connpika.BlockingConnection(pika.ConnectionParameters(127.0.0.1,5672,promo,pika.PlainCredentials(promo,promo_dev_2026)))chconn.channel()ch.basic_qos(prefetch_count1)ch.basic_consume(q.order.pay,on_msg,auto_ackFalse)ch.start_consuming()conn.close()运行结果第一次崩溃后list_queues消息数不减或 unacked 回 ready。第二次同一 body 可能redeliveredTrue。坑第二次处理必须幂等否则短信双发——这解释了客服「两条短信」。坑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-1s basic.qos prefetch_count %% Each channel has an associated limiter process3.4 步骤三Nack 放回 vs 丢掉步骤目标requeueTrue会再拿到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,requeueTrue)returnch.basic_ack(method.delivery_tag)ch.stop_consuming()connpika.BlockingConnection(pika.ConnectionParameters(127.0.0.1,5672,promo,pika.PlainCredentials(promo,promo_dev_2026)))chconn.channel()ch.basic_qos(prefetch_count1)ch.basic_consume(q.order.pay,on_msg,auto_ackFalse)ch.start_consuming()conn.close()运行结果同一条至少打印 3 次前两次 nack。这就是毒消息循环的缩影——生产必须有次数上限。再开一次requeueFalse换一条新消息后深度减 1 且不再回来。坑单消费者 nack requeue 可能立刻拿回同一条CPU 打满。多消费者时可能交给别人问题变成随机。3.5 步骤四prefetch1 vs 50步骤目标慢处理下观察messages_unacknowledged。先 seed 30 条到专用队列避免打乱支付队列# promo-mq/ch09/prefetch_lab.pyimporttime,threading,pikadefconsume(prefetch,seconds8):connpika.BlockingConnection(pika.ConnectionParameters(127.0.0.1,5672,promo,pika.PlainCredentials(promo,promo_dev_2026)))chconn.channel()ch.queue_declare(q.lab.prefetch,durableTrue)ch.basic_qos(prefetch_countprefetch)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_ackFalse)t0time.time()whiletime.time()-t0seconds:conn.process_data_events(time_limit0.2)conn.close()# 先灌 30 条到 q.lab.prefetch 再分别跑 prefetch1 和 50# 跑的同时rabbitmqctl list_queues -p promo name messages messages_unacknowledged另开终端每秒打一次dockerexecrabbit-promo-1 rabbitmqctl list_queues-ppromo name messages messages_unacknowledged运行结果prefetch1 时 unacked 约为 150 时 unacked 可冲到几十不超过 50 且不超过剩余消息。吞吐上 50 通常更高直到网关或 Broker 内存成为瓶颈。坑在 BlockingConnection 里sleep会挡住心跳实验 sleep 0.3s 可接受生产 3s 同步 sleep 会掐连接。用线程池或异步。坑两个消费者同时消费同一队列时prefetch 是每个通道的窗口总 inflight 是相加关系。值班口诀ready0 且 consumers0 是「没人干活」unacked 持续等于 prefetch 且 ready 仍涨是「人慢或卡死」unacked 长期等于消息总数且 ready0是「忘 Ack」。三种告警文案要分开否则运维只会重启消费者把忘 Ack 变成重复短信。3.6 完整代码清单column/samples/ch09/ seed.py autoack_crash.py manual_crash.py manual_ack.py nack_demo.py prefetch_lab.py3.7 测试验证编号操作期望TC-CH09-01autoAck 硬退出消息消失TC-CH09-02手动未 Ack 硬退出消息回 readyTC-CH09-03重投redelivered trueTC-CH09-04nack requeue再次投递TC-CH09-05prefetch1unacked≤1curl-s-upromo:promo_dev_2026\http://127.0.0.1:15672/api/queues/promo/q.lab.prefetch\|rgmessages_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 吃内存prefetch1简单、压力平滑延迟差优点1语义可测。2redelivered 提示至少一次。3limiter 把 QoS 从 Channel 抽出去避免打爆 channel 进程。缺点1至少一次 ≠ 恰好一次。2经典队列无 delivery-limit。3sleep 式消费害心跳。消费侧口诀先做事再 Ack失败就分类再试 / 死信 / 丢窗口按慢速环节设而不是按 QPS 设。QPS 是结果prefetch 是约束。用 QPS 反推窗口可以用「感觉卡」直接把窗口加到 500 不行。适用场景短信/邮件等必须手动 Ack。CPU 很重或要严格串行时 prefetch1。网关稳定时适度加大窗口。崩溃注入与重复投递的测试训练。不适用autoAck 用于支付用 redelivered 当去重 ID用 Nack 循环当重试退避第 10 章 TTL。注意事项4.3 经典队列不再按旧 consumer_timeout 那套评估。多线程不要共享 Channel 去 Ack。安全消费者账号只要 read不要 configure 删队列。basic.get不受 prefetch 限制limiter 注释写明监控脚本乱 get 会捣乱。常见踩坑生产autoAck 网关超时K8s 杀 pod短信丢失。根因投递即删除。prefetch1000消费者 Full GCunacked 占满内存触发告警。根因窗口当缓冲。finally 里一律 Ack业务失败也消消息。根因把「通道还活着」当「业务成功」。思考题两个消费者 prefetch 各 50队列 10 条unacked 最大可能多少会不会一个吃完另一个饿死手动 Ack 成功后应用仍崩溃用户已看到短信补偿流程应靠什么而不是 Nack附录 C第 8 章思考题参考答案题 1Confirm 后 kill -9。经典队列仍可能丢尾部。quorum 多数派提交后更稳。发布器不能承诺「Confirm永存」。题 2SENT 后用户再点支付。靠订单状态机与短信发送记录幂等不靠 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 实战修炼与源码剖析