3.1 消息确认机制与消费者限流 本节摘要:确认机制决定"消息什么时候算送达":自动确认在投递瞬间记账,手动确认在业务处理完成后签收。限流(prefetch)决定"一次预支多少条"。两者共同构成消费端的防丢与防压双闸。本节从深夜追查的第一现场说起,讲透签收的三种结果与限流的调参逻辑。 第 3 章的追查从消费端开始不是偶然——统计意义上,消息丢失的头号原因不在服务器,而在消费者代码里那个没有关掉的 autoack。本节把第一现场的证据链完整摆一遍。 自动确认为什么危险 消费者订阅队列时,autoack(即 AMQP 的 auto-ack 模式)决定签收语义。开启时,Broker 把消息写出 TCP 连接的瞬间就将其标记为已送达并从队列删除。
本节摘要:确认机制决定"消息什么时候算送达":自动确认在投递瞬间记账,手动确认在业务处理完成后签收。限流(prefetch)决定"一次预支多少条"。两者共同构成消费端的防丢与防压双闸。本节从深夜追查的第一现场说起,讲透签收的三种结果与限流的调参逻辑。
第 3 章的追查从消费端开始不是偶然——统计意义上,消息丢失的头号原因不在服务器,而在消费者代码里那个没有关掉的 auto_ack。本节把第一现场的证据链完整摆一遍。
消费者订阅队列时,auto_ack(即 AMQP 的 auto-ack 模式)决定签收语义。开启时,Broker 把消息写出 TCP 连接的瞬间就将其标记为已送达并从队列删除。之后发生什么,Broker 一概不管:消费者进程被杀、回调抛异常、机器断电——消息已经不在队列里了,它只存在于那个再也无法回调的进程内存里。
第一现场的真相正是如此:值班时某节点 OOM 被系统杀掉,自动确认模式下已投递未处理的三百条消息随之蒸发,队列自然查无此物,消费日志自然没有记录。证据闭环。
手动确认把记账时点挪到业务完成之后:回调处理成功才发 basic_ack,Broker 收到签收才删消息。处理中崩溃的消息会在连接断开后重新入队,等待下一次投递——"至少一次"送达语义由此而来。
三种签收动作的语义要分清:
| 动作 | 语义 | 消息去向 |
|---|---|---|
| basic_ack | 处理成功,签收 | 从队列删除 |
| basic_nack requeue=true | 处理失败,退回 | 重新入队,可再投 |
| basic_nack requeue=false | 处理失败,不退 | 进入死信流程(若有配置) |
nack 的 requeue 参数是设计决策点:一律退回会造成"毒消息"循环——一条解析必然失败的消息被无限重投,把整条队列堵死。工程惯例是有限重试:失败先退回重试两三次,超过次数后 requeue=false 转投死信队列,由人工或补偿任务接手。
把一条消息在消费端的完整状态旅程画出来,签收时点的分野一目了然:

图底那句话值得再读一遍:自动确认不是"另一种签收方式",而是绕过整个状态机——失败退回、重试上限、死信接管统统不生效。这也是为什么本册把"关掉 auto_ack"列为可靠链路的第一刀。
背景:把 3.1 的消费端改造成手动签收加有限重试,验证异常场景下的消息不丢。
操作:
import pika, json connection = pika.BlockingConnection(pika.ConnectionParameters(host="localhost")) channel = connection.channel() channel.queue_declare(queue="trace.orders", durable=True) fail_count = {} def on_message(ch, method, properties, body): event = json.loads(body) try: process(event) # 业务处理 ch.basic_ack(delivery_tag=method.delivery_tag) except TransientError: # 瞬时错误:退回重试,累计三次后转死信 n = fail_count.get(event["order_id"], 0) + 1 fail_count[event["order_id"]] = n requeue = n < 3 ch.basic_nack(delivery_tag=method.delivery_tag, requeue=requeue) except Exception: # 业务性错误:重试无意义,直接转死信人工排查 ch.basic_nack(delivery_tag=method.delivery_tag, requeue=False) channel.basic_consume( queue="trace.orders", on_message_callback=on_message, auto_ack=False) # 显式手动签收 channel.start_consuming()
验证不丢的实验:消费回调里故意抛异常且 requeue=True,随后 kill 掉消费者进程再重启,消息依然能被重新投递处理——连接断开时未签收的消息自动回到队列。
结果与解读:改为手动签收后,"处理成功"与"消息删除"被绑定成一个动作,蒸发窗口消失。但注意代价:未签收的消息占着 unacked 名额,若消费端堆积大量慢处理,队列会呈现 messages_ready 为零而 messages_unacknowledged 巨大的形态——这是 1.3 节三列指标的实战读法,也是引入限流的直接理由。
变式:把 auto_ack=True 与手动版各跑一次 kill 实验,对比队列残留:自动版队列清空但业务未完成,手动版消息完整归队。这个对比实验建议在本地亲手做一遍,五分钟胜过十页文档。
Broker 默认尽可能多地把消息推向消费者(prefetch 无上限时甚至一次推空队列)。这对快消费者是好事,对慢消费者是灾难:消息全堆在消费者本地的未签收缓冲里,既占内存,又让其他消费者抢不到任务,还让"重新入队"的失败代价变高。
prefetch_count 限定"每条通道上最多多少条未签收消息":达到上限后 Broker 停止投递,直到有签收回流。它因此成为消费端唯一的流控阀门。调参经验:prefetch 约等于单条处理耗时乘以吞吐目标,再留一点余量。处理耗时五十毫秒、希望单消费者每秒二十条,prefetch 设一到五足够;处理耗时三秒的慢任务,prefetch 设一反而最稳——一条一签,节奏完全受控。
# 限流声明:同一通道最多 10 条未签收 channel.basic_qos(prefetch_count=10) channel.basic_consume(queue="trace.orders", on_message_callback=on_message, auto_ack=False)
💡 关键直觉:prefetch 是用吞吐换公平与安全的阀门。自动确认模式下 prefetch 完全失效——没有签收动作,就没有计数依据。这也是两者必须配套的原因。
⚠️ 常见坑:回调里先处理再签收时忘记 try,异常直接冒泡导致通道与连接被客户端库关闭,未签收消息全部重新入队,看似"没丢"实则引发了整通道重投——异常必须兜在回调内部。
消费端堵住了。下一现场深入存储层:持久化三件套如何让消息在服务器重启中幸存。