KEEL · 龙骨 · A CURRICULUM FOR THE AI ERA
03 · 下行可靠性:ack、nack、requeue 与 prefetch — keel 龙骨
上行解决了「消息进队列」,下行要解决的是「消息被正确处理、且处理能力被正确分配」。这里有两个独立的坑:一是 ack 时机错了导致消息丢或重复,二是没设 prefetch 导致 broker 把整队列消息全推给一个慢消费者、其它消费者空转。第二个坑在监控上表现为「消费速率正常、但队列迟迟不空」,很容易误判。本章的输出出自本机 RabbitMQ 4.3.6,节点名已脱敏。
上行解决了「消息进队列」,下行要解决的是「消息被正确处理、且处理能力被正确分配」。这里有两个独立的坑:一是 ack 时机错了导致消息丢或重复,二是没设
prefetch导致 broker 把整队列消息全推给一个慢消费者、其它消费者空转。第二个坑在监控上表现为「消费速率正常、但队列迟迟不空」,很容易误判。本章的输出出自本机 RabbitMQ 4.3.6,节点名已脱敏。
一、现场:两台消费者,一台忙死一台闲死
订单队列起了两台消费者,代码一样。压测发现吞吐只有单台的水平,另一台的 CPU 贴着地板。同时 list_queues 显示这个队列有大量 messages_unacknowledged,但两台机器「看起来都在处理」。
另一个现场:某个消费逻辑抛异常后没调用 basic_nack(requeue=False),而是让异常往上冒、连接断开。消息被 requeue,又被同一台机器取到,又失败——日志里同一条消息刷了上千遍。
两个现象都指向同一层:消费者确认机制。
二、下行链路的确认与限流
flowchart TD
Q["queue lab.prefetch.q"] --> G{"消费者未确认数<br/>达到 prefetch 上限?"}
G -->|"未达上限"| DEL["deliver 到 on_message"]
G -->|"已达上限"| HOLD["留在队列,不投递"]
DEL --> PROC["业务处理"]
PROC --> DEC{"处理结果"}
DEC -->|"成功"| ACK["basic_ack<br/>从 unacked 移除"]
DEC -->|"失败、可重试"| REQ["basic_nack requeue=True"]
DEC -->|"失败、不可重试"| DEAD["basic_nack requeue=False"]
REQ -.->|"重新入队,可能原地打转"| Q
DEAD -.->|"死信或直接丢弃"| DLX["lab.ex.dlx -> queue lab.dlq"]
ACK --> G
HOLD -.->|"消费者卡住又不 ack"| STUCK["消息长期占用 unacked<br/>只有连接断开才释放"]
style Q fill:#e8f5e9,color:#1b5e20
style G fill:#fff3e0,color:#8a4b00
style DEL fill:#e3f2fd,color:#0d3b66
style PROC fill:#e3f2fd,color:#0d3b66
style DEC fill:#fff3e0,color:#8a4b00
style ACK fill:#e8f5e9,color:#1b5e20
style REQ fill:#fff8e1,color:#8a4b00
style DEAD fill:#ffebee,color:#b71c1c
style DLX fill:#ffebee,color:#b71c1c
style HOLD fill:#f1f8e9,color:#33691e
style STUCK fill:#ffebee,color:#b71c1c
三、ack 与 nack 的语义
三个动作,语义要分清:
basic_ack:这条消息处理完了,broker 可以删掉它。必须在业务处理成功之后调用。basic_nack(requeue=True):拒绝,但请你把它放回队列。消息回到队列头部(RabbitMQ 会尽量放回原位置),可被任何消费者(包括自己)再次取到。basic_nack(requeue=False):拒绝且不要了。消息按队列的x-dead-letter-exchange规则走死信,没配 DLX 就直接丢弃。
用 lab.dl.main(配了 DLX 指向 lab.ex.dlx,lab.dlq 绑在 DLX 上)实测两种 nack,输出(05-manual-ack-dlx.txt):
收到消息(auto_ack=False 手动确认) body=被拒绝的消息 delivery_tag=1
nack(requeue=False) 后 lab.dl.main 0 条,lab.dlq 1 条
nack(requeue=True) 后 lab.dl.main message_count=1
requeue=False 让消息离开主队列、进了死信队列;requeue=True 让消息回到主队列(passive 查询显示深度 1)。这就是两个动作的差别,auto_ack=False 则保证了「消息在 ack 之前一直算 unacked,进程崩了会被重新投递」。
有一个容易忽略的点:requeue=True 如果用在「确定会失败」的消息上,就是死循环。 上面第二个现场里的消息就是这类——处理逻辑对某种 payload 永远抛异常,requeue 之后又被取到,永远失败。正确做法有二:一是业务侧识别出「重试也没用」的异常,用 requeue=False 进死信;二是用重试次数(x-death 的 count,第 04 章)设上限。
四、prefetch:把「限流」从 broker 记忆里拿回来
AMQP 的 basic.qos 里的 prefetch_count 控制的是一个消费者最多能同时持有多少条未确认消息。不设时,broker 会尽可能快地把消息推给消费者——这本来是为了吞吐,但遇到「一个消费者慢、其它消费者快」时就出事。
实测用两台消费者连到同一个队列 lab.prefetch.q:慢消费者每条处理 0.3 秒,快消费者每条 0.05 秒,队列里预置 20 条。先不设 prefetch(等于不限):
[不设 prefetch(0=不限)] 运行中(约 1s)消费条数:慢=11,快=0
[不设 prefetch(0=不限)] 同一时刻 lab.prefetch.q 行:messages=10 messages_unacknowledged=10
[不设 prefetch(0=不限)] 运行 4s 结束后累计:慢=20,快=0
慢消费者把 20 条全拿走了,快消费者一条没接到。中途快照更能说明机制:队列里剩 10 条 ready、10 条 unacked——broker 已经把一半推给了慢消费者并标记为未确认,慢消费者正在一条条处理,而快消费者因为抢不到而空转。
再设 basic_qos(prefetch_count=1):
[prefetch_count=1] 运行中(约 1s)消费条数:慢=3,快=17
[prefetch_count=1] 运行 4s 结束后累计:慢=3,快=17
消费被均分了:快消费者拿到 17 条,慢消费者 3 条。prefetch=1 的含义是「一个消费者手上最多 1 条未确认」,慢消费者在慢慢处理那条时,broker 就把新消息投给空闲的快消费者。
prefetch_count 的取值是个吞吐与公平性的权衡:设 1 最公平、但每条消息都要一次往返,吞吐会掉;设得大(比如 10~100)吞吐高,但慢消费者会一次性囤积很多未确认消息,公平性变差;设成 0 等于不限。这里没有通用阈值——要根据单条处理耗时、消息体大小、消费者数量实测。RabbitMQ 官方已把 channel 级的 global QoS 标记弃用(global_qos,2021 年 8 月,检索于 2026-10-05),现在应按消费者(per-consumer)设置,也就是在消费前对该 channel 调 basic_qos。
五、常见误判
| 常见误解 | 对着哪条输出核对 | 结论 |
|---|---|---|
| 消费速率正常就说明分发正常 | 不设 prefetch 时中途 messages=10, unacked=10,快消费者拿到 0 条 |
「在处理」不等于「在并行」;一台在慢慢处理、其余空转,速率一样看起来正常 |
requeue=True 是重试的标准做法 |
nack(requeue=True) 后主队列深度回到 1 |
requeue 无限循环只适合「可重试且很快会成功」的场景,否则要上进死信 |
| 收到消息就该先 ack,处理失败再补 | nack(requeue=False) 后消息从主队列消失、进了 lab.dlq |
先 ack 再处理,处理失败的消息就真丢了;应当处理成功才 ack |
| prefetch 越大吞吐越高 | prefetch=1 时快消费者反而吃下 17 条 | 大 prefetch 提高吞吐的前提是所有消费者都差不多快;不一致时会加剧不公平 |
生产边界
- 教学替身 vs 真实依赖:实验用两台同进程、只靠 sleep 模拟快慢的消费者。生产的消费者快慢差异来自下游依赖(数据库、外部 API),而且会随负载变化;prefetch 值需要在真实下游延迟下校准。
- 上线要盯的指标:
list_queues的messages_unacknowledged(长期不降说明有消费者卡住或处理慢)、每个消费者的实际消费速率、消息从入队到 ack 的端到端延迟(不只是消费速率)。 - 失败策略:明确区分「可重试异常」(网络抖动、下游暂时不可用)与「不可重试异常」(payload 格式错、业务规则拒绝);前者用带次数上限的重试,后者直接
requeue=False进死信;所有消费者必须用auto_ack=False,否则处理中崩溃必丢。
动手
- 起两个消费者连同一队列,一个 sleep 0.3s、一个不 sleep,不设 prefetch 跑 5 秒,记录各自消费条数与中途的
messages_unacknowledged。 - 加上
basic_qos(prefetch_count=1)重跑,对比两次的分发比例。 - 再试
prefetch_count=5,观察慢消费者一次囤积多少条未确认消息。 - 让消费者处理时抛异常并在
except里basic_nack(requeue=True),观察这条消息被反复取到的次数;改成requeue=False并配 DLX,确认它只被处理一次就进死信。
自测
requeue=True与requeue=False分别把消息送到哪里?各自适合什么场景?- 为什么
prefetch_count=0(不设)会让慢消费者把消息全吸走?机制上它比消费者「主动拉」多了什么? messages_unacknowledged长期不为零说明什么?它和messages_ready的区别在哪?- 先 ack 再处理业务,在处理中进程崩溃会发生什么?和先处理后 ack 有什么不同?
- 同一个 channel 上加 prefetch 和给每个消费者单独加,RabbitMQ 4.x 更推荐哪种?为什么?
↓ 下一步:04 章 · 死信与延迟