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 的语义

三个动作,语义要分清:

用 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 提高吞吐的前提是所有消费者都差不多快;不一致时会加剧不公平

生产边界

动手

  1. 起两个消费者连同一队列,一个 sleep 0.3s、一个不 sleep,不设 prefetch 跑 5 秒,记录各自消费条数与中途的 messages_unacknowledged。
  2. 加上 basic_qos(prefetch_count=1) 重跑,对比两次的分发比例。
  3. 再试 prefetch_count=5,观察慢消费者一次囤积多少条未确认消息。
  4. 让消费者处理时抛异常并在 except 里 basic_nack(requeue=True),观察这条消息被反复取到的次数;改成 requeue=False 并配 DLX,确认它只被处理一次就进死信。

自测

  1. requeue=True 与 requeue=False 分别把消息送到哪里?各自适合什么场景?
  2. 为什么 prefetch_count=0(不设)会让慢消费者把消息全吸走?机制上它比消费者「主动拉」多了什么?
  3. messages_unacknowledged 长期不为零说明什么?它和 messages_ready 的区别在哪?
  4. 先 ack 再处理业务,在处理中进程崩溃会发生什么?和先处理后 ack 有什么不同?
  5. 同一个 channel 上加 prefetch 和给每个消费者单独加,RabbitMQ 4.x 更推荐哪种?为什么?

↓ 下一步:04 章 · 死信与延迟

进入 keel 阅读