🤖
AI审核中

从Ready、Unacked到Prefetch的排障实战

中间件 18分钟 101浏览 0评论

遇到 RabbitMQ 消息积压,最直接的处理方式往往是增加消费者。

两个实例不够,就加到四个;四个不够,再加到八个。与此同时,把 Prefetch 从几十调到几百,希望消费者一次多拿一点消息,尽快把队列处理完。

但有时会出现一种反直觉的现象:消费者增加了,Ready 消息减少了,业务完成速度却没有明显提升。部分实例持续忙碌,另一些实例几乎没有任务,甚至应用内存和下游数据库压力还进一步上升。

理解这种问题,首先要把三个概念分开:

消息被投递了,不等于消息正在执行;消息正在执行,也不等于业务已经完成。

本文围绕 AMQP 0-9-1、basic.consume 推送消费和手动确认模式展开,重点讨论普通工作队列,不将这些结论直接套用到 RabbitMQ Stream Protocol。

一、Ready 下降,不一定意味着积压正在消失

1. 先看懂三个队列指标

RabbitMQ 的队列指标中,最容易混淆的是下面三项:

指标 含义 不能据此得出的结论
messages_ready 等待投递给消费者的消息数量 不能代表系统全部未完成任务
messages_unacknowledged 已投递、但尚未收到消费者确认的消息数量 不能代表真正同时执行的任务数量
messages Ready 与 Unacked 之和 不能直接代表业务系统中全部未完成任务

这些指标描述的是消息在 Broker 视角下的状态,而不是业务状态。

例如,假设某次调整前:

  • Ready 为 98,000。
  • Unacked 为 2,000。

调整 Prefetch 后,Ready 变成 20,000,Unacked 变成 80,000。

如果只观察 Ready,似乎已经消化了大量积压。但两项相加,仍然是 100,000 条。

在这个示例中,变化的是消息所处的阶段,不是已经完成的工作量。

2. Unacked 里面可能有大量“尚未开始”的任务

手动确认模式下,一条消息可能经过如下路径:

graph LR
    A["生产者发布"] --> B["Broker 中等待投递"]
    B --> C["网络传输与客户端缓冲"]
    C --> D["应用线程池等待"]
    D --> E["执行业务逻辑"]
    E --> F["完成事务或可靠接管"]
    F --> G["发送 Ack"]
    G --> H["Broker 收到 Ack"]

从 Broker 投递消息,到收到确认之间,消息都可能处于 Unacked 状态。其中既包括正在执行业务的消息,也包括仍在客户端缓冲区等待的消息。消费者确认与生产者 Publisher Confirm 则属于不同方向的机制,后者不是业务处理完成的回执。

因此,排障时至少应同时回答两个问题:

Broker 还保留着多少未确认消息?业务每秒真正完成了多少任务?

这两个问题没有建立联系之前,单独优化队列曲线很容易产生误判。

二、Prefetch 是未确认窗口,不是处理线程数

1. Prefetch 到底限制什么

Prefetch 限制的是允许存在的未确认投递数量。

例如,消费者使用 basicQos(20) 并采用手动确认,意味着在相应限额范围内,Broker 可以先投递一批消息;当未确认数量达到限制后,需要等待确认释放额度,才能继续投递。

RabbitMQ 常用的非全局配置按 Consumer 分别应用,而不是天然按整个应用实例共享。0 表示不设置该 QoS 限制,并不是暂停消费;具体队列类型仍可能施加额外限制。

所以,下面三件事完全不同:

实例数决定部署规模,业务执行并发决定同时处理能力,Prefetch 决定允许提前投递多少未确认消息。

把 Prefetch 从 20 改成 1,000,不会自动创建更多处理线程。

2. 一个简单模型就能看出问题

假设某个消费者串行执行任务,每条任务耗时约 100 毫秒。

忽略网络和其他开销,其处理能力约为:

μ≈10.1=10 条/秒\mu \approx \frac{1}{0.1}=10\ \text{条/秒}

现在把 Prefetch 设置为 1,000。

如果这 1,000 条消息被提前投递到该消费者,排在末尾的任务可能要等待接近 100 秒才能完成。

这里的 100 秒只是串行、等耗时条件下的估算,不是 RabbitMQ 的固定行为,更不是实测结果。但它说明了一个问题:

扩大预取窗口,可能只是让更多任务提前进入一个处理能力没有变化的等待区。

3. 内存预算要按整个应用计算

假设一个应用有四个实例,每个实例注册四个 Consumer,每个 Consumer 的 Prefetch 为 500。

在独立限额、没有其他更小限制的估算条件下,总未确认窗口可以达到:

4×4×500=8,0004 \times 4 \times 500 = 8{,}000

若每条消息体为 256 KiB,且这些消息体都进入客户端,仅消息体数据就接近 1.95 GiB。

这还没有计算反序列化对象、业务上下文和其他缓冲区。

这个数字是容量预算,不代表这些数据必然同时驻留在堆内存中。它提醒我们:评估 Prefetch 时,不能只看一个配置值,还必须结合 Consumer 数量和消息大小。官方文档也指出,提高预取量可能增加消费者侧的内存消耗。

三、为什么加消费者没有明显效果

1. 消息已经被老消费者提前拿走

这是大 Prefetch 下很容易被忽略的情况。

假设队列里有 1,000 条任务,消费者 A 先启动,并设置较大的预取窗口。它很快拿到了这批任务,但实际处理速度很慢。

随后启动消费者 B。

B 虽然有空闲处理能力,却不意味着已经交给 A、仍未确认的任务会立即重新分配给它。RabbitMQ 的工作队列教程也正是通过较小的预取窗口,避免忙碌消费者提前领取过多任务。

此时可能同时出现:

Ready 很低、Unacked 很高、老实例忙碌、新实例空闲。

不过,这只是排查方向,不能仅凭指标组合就下结论,还要确认消息在各个 Channel 上的分布。

另一个容易踩坑的地方是:取消消费订阅不等于把已经投递的消息全部放回队列。basic.cancel 主要停止后续投递,已有的未确认消息不会因此自动重新入队。

2. 真正的限制在下游

考虑另一种假设:消费者的主要工作是访问一个下游服务,而该服务在当前业务条件下最多能稳定完成每秒 200 次请求。

把消费者从四个增加到二十个,并不会自动让下游变成每秒 1,000 次。

新增并发可能表现为更多等待、更多超时,或者更高的重试比例。因此,在决定扩容前,应验证扩容能否增加“有效完成能力”,而不是只增加同时发起的请求数。

可以把每条任务的耗时拆成:

T任务=T本地等待+T资源获取+T业务执行+T下游等待T{\text{任务}} = T{\text{本地等待}} + T{\text{资源获取}} + T{\text{业务执行}} + T_{\text{下游等待}}

这是一种排障建模方式。它要求应用补充观测数据,而不是由 RabbitMQ 指标直接推导结果。

例如,如果耗时主要花在数据库连接获取上,首先需要验证连接池和数据库的约束;如果主要卡在同一业务对象的串行修改上,则应分析业务串行点,而不是默认消费者不足。

3. 增加的是配置并发,不是实际执行并发

在 Java 客户端中,还需要区分 Consumer、Channel 和执行线程。

官方 Java 客户端会保持同一 Channel 上的投递回调顺序。多个 Consumer 共用一个 Channel 时,一个耗时回调可能拖住同一 Channel 上的其他回调。仅仅扩大回调线程池,并不能保证同一 Channel 的消息开始并行处理。

因此,应检查的是:

到底有多少个任务正在同时执行业务逻辑,而不是配置文件里写了多少个线程。

排查时,可以把应用的活动任务数、线程池等待数,与 RabbitMQ 的 Consumer 数量放在一起观察。两者不一致,往往比单独增加线程更值得研究。

四、先建立证据链,再修改参数

1. 查看队列与消费者

以下命令在能够管理目标 RabbitMQ 节点的环境中执行。示例虚拟主机为 /,实际使用时应替换为目标虚拟主机。

rabbitmqctl list_queues -p / \
  name messages_ready messages_unacknowledged consumers

rabbitmqctl list_consumers -p /

rabbitmqctl list_channels \
  pid connection number vhost \
  consumer_count messages_unacknowledged prefetch_count

第一条命令观察队列状态,第二条核对消费订阅、确认要求和预取限制,第三条查看未确认消息集中在哪些 Channel 上。

需要特别注意:list_channelsprefetch_count 描述的是该 Channel 对新消费者的 QoS 设置,核对具体订阅时仍应结合 list_consumers,不能只看一列数值。

大型集群中,不要为了观察一个队列而高频扫描全部对象。RabbitMQ 官方监控指南明确提醒,过量请求监控数据本身也会增加节点负担。

2. 用指标组合提出假设

下面这张表是排查分支,不是自动诊断规则:

观察到的现象 优先验证的方向
Ready 持续增加,Consumer 为零 消费订阅是否建立,是否连错虚拟主机或队列
Ready 很低,Unacked 很高 过度预取、客户端排队、处理阻塞、确认遗漏
Ready 和 Unacked 都高 输入超过完成能力,或消费者已达到未确认窗口
老实例忙,新实例空闲 未确认消息分布是否高度集中
投递很活跃,业务成功数不增长 重复投递、重试循环、业务异常
Broker 消息减少,应用任务仍未完成 是否提前确认,或已经转交其他处理系统

这些方向需要继续用 Channel 分布、应用日志、业务成功数和下游耗时交叉验证。队列指标和投递行为的基础语义来自 RabbitMQ 的监控与确认机制。

3. 不要忽略 Broker 资源告警

如果观察到发布端阻塞,还应检查连接状态和节点资源告警。

RabbitMQ 的内存、磁盘告警主要通过阻塞发布连接来保护节点;仅用于消费的连接不会因此被同样阻塞。生产和消费混用一个连接时,影响边界会更复杂。

所以,“发布开始阻塞”和“消费者拿不到更多消息”不能直接当成同一种问题。

五、用对照实验观察“积压转移”

下面给出一个隔离环境中的实验方案,观察结果以实际运行记录为准,不将理论估算当作实测吞吐。

实验使用 Linux Shell、Docker、JDK 17 或 21,以及 RabbitMQ 官方 PerfTest。PerfTest 支持控制生产者数量、消费者数量、预取量和模拟处理耗时。

1. 准备独立实验环境

创建一个只向本机开放端口的 RabbitMQ 容器。下面的账号仅用于本地实验,不要复用于生产环境。

docker run -d \
  --name mq-prefetch-lab \
  --hostname mq-prefetch-lab \
  -p 127.0.0.1:5673:5672 \
  -p 127.0.0.1:15673:15672 \
  -e RABBITMQ_DEFAULT_USER=lab \
  -e RABBITMQ_DEFAULT_PASS=lab-only \
  rabbitmq:4.2-management

查看启动状态,确认节点已就绪后再进行实验:

docker logs --tail 50 mq-prefetch-lab

docker exec mq-prefetch-lab rabbitmq-diagnostics -q ping

下载固定版本的 PerfTest。这里使用官方发布的 2.25.0 JAR,而不是跟随浮动版本。

curl -fL \
  -o perf-test.jar \
  https://github.com/rabbitmq/rabbitmq-perf-test/releases/download/v2.25.0/perf-test-2.25.0.jar

java -jar perf-test.jar --help

2. 先放入 1,000 条消息

创建一个独立实验队列,只发布消息,不启动消费者:

java -jar perf-test.jar \
  --uri amqp://lab:lab-only@127.0.0.1:5673/%2f \
  --producers 1 \
  --consumers 0 \
  --queue prefetch-lab-a \
  --queue-args x-queue-type=classic \
  --auto-delete false \
  --flag persistent \
  --size 1024 \
  --rate 500 \
  -C 1000 \
  -c 100 \
  --use-millis

确认消息已经进入队列:

docker exec mq-prefetch-lab rabbitmqctl list_queues -p / \
  name messages_ready messages_unacknowledged consumers

3. 启动大预取的慢消费者

在一个终端中运行:

java -jar perf-test.jar \
  --uri amqp://lab:lab-only@127.0.0.1:5673/%2f \
  --producers 0 \
  --consumers 1 \
  --predeclared \
  --queue prefetch-lab-a \
  --qos 1000 \
  --consumer-latency 100000 \
  --use-millis \
  --id slow-a \
  -z 180

这里模拟每条消息约 100 毫秒的处理耗时。--consumer-latency 的单位是微秒,不是毫秒;实验没有开启自动确认。

再次观察队列,等待出现 Ready 接近零、Unacked 仍明显大于零 的状态。

4. 再启动更快的消费者

在另一个终端中运行:

java -jar perf-test.jar \
  --uri amqp://lab:lab-only@127.0.0.1:5673/%2f \
  --producers 0 \
  --consumers 1 \
  --predeclared \
  --queue prefetch-lab-a \
  --qos 10 \
  --consumer-latency 10000 \
  --use-millis \
  --id fast-a \
  -z 180

这个消费者模拟每条消息约 10 毫秒的处理耗时。

预期观察重点不是“它一定能达到多少吞吐”,而是:

当旧消费者已经持有大部分未确认消息时,新消费者能否拿到足够多的任务?

如果没有新消息进入,而且存量消息都已投递给旧消费者,新消费者即使处理能力更强,也可能没有任务可做。这是根据预取和未确认消息机制得出的实验预期。

5. 做第二轮对照

记录第一轮结果,结束本轮两个消费者进程。

随后使用全新队列 prefetch-lab-b 重复上述过程:仍然先发布 1,000 条消息,但把慢消费者的 --qos1000 改成 10,其他条件保持一致,再启动快消费者。

对照时记录 Ready、Unacked、两个消费者的处理分布和整体清空时间。

这个实验验证的是任务分配和预取窗口的影响,并不等价于真实数据库、外部接口或复杂业务的容量测试。

实验结束后,删除的仅应是本节新建的实验容器及其匿名数据卷:

docker rm -f -v mq-prefetch-lab

六、修复积压,不能以牺牲可靠性为代价

1. 不要提前 Ack,让问题从监控里消失

一种危险的“优化”是:收到消息后立即放进本地线程池,然后马上确认。

这样 Broker 指标可能迅速下降,但任务仍然在应用内等待。如果进程此时退出,已经确认的消息不能再指望由 Broker 按未确认消息重新投递。

确认边界应该与应用承担的责任一致:业务已经完成,或者任务已被另一个可靠存储系统接管,而不是仅仅进入了一个内存队列。

需要区分两种设计:

“提交到本地线程池后确认”只是转移到易失内存;“可靠写入任务表后确认”则是把后续执行责任转交给持久化任务系统。

后一种设计可以成立,但必须继续监控任务表积压,不能用 RabbitMQ 清空来证明业务已经处理完。

2. 正确确认,仍然需要幂等

即使坚持业务完成后才确认,仍然存在一个窗口:

业务事务已经提交,但确认未能被 Broker 成功接收。

之后消息可能再次投递,因此消费者需要能够安全处理重复任务。RabbitMQ 的可靠性指南明确要求应用考虑重复投递与幂等处理。

工程上,可以把业务幂等记录和业务变更放在同一个数据库事务中。涉及外部系统时,还需要设计对方认可的幂等键或可查询的操作结果。

这里最重要的边界是:

不能把“Ack 失败”直接解释成“刚才的业务操作没有成功”。

3. 不要把所有失败都立即重新入队

basicNackbasicReject 可以要求消息重新入队,但重新入队的消息可能很快再次被投递。

如果所有消费者都因为同一个下游故障不断失败,再立即重新入队,就可能形成高频重试。

此时应区分可恢复故障、永久性错误和重试耗尽,给出有上限、有间隔的处理策略,而不是只追求“消息不能离开原队列”。

同时,配置死信交换机不等于完成了可靠重试设计。死信转发本身也存在失败边界;默认死信重新发布并不提供所有情况下的可靠转移保证,Quorum Queue 的至少一次死信机制则需要按对应条件配置。

七、Prefetch 应该怎样调

1. 先保持业务并发不变,只改变预取量

可以把较小、中等和较大的预取窗口作为实验组,例如 1、8、32、128。

这不是推荐所有系统使用这些值,而是为了观察变化趋势。每轮只改变一个因素,并保持消息大小、任务耗时分布和输入流量尽量一致。

重点观察:

维度 需要验证的问题
有效完成速度 真正成功完成的任务是否增加
端到端延迟 尾部任务是否等待更久
资源消耗 内存和下游压力是否上升
分配情况 慢消费者是否持有过多未确认消息

如果增大 Prefetch 后,完成速度几乎不变,但等待时间和资源占用明显增加,那么继续扩大窗口就缺乏收益。

目标不是最大 Prefetch,而是找到能维持目标吞吐、又不过度占用资源的窗口。

2. 把“清空时间”算出来

假设待处理存量为 BB,持续输入速率为 λ\lambda,有效完成速率为 μ\mu。

在消息耗时相对稳定、没有大量过期、丢弃或重试干扰的简化条件下:

T清空≈Bμ−λT_{\text{清空}} \approx \frac{B}{\mu-\lambda}

这个估算只有在 μ>λ\mu>\lambda 时才有意义。

例如,假设积压 120,000 条,每秒输入 800 条,每秒有效完成 1,000 条:

T清空≈120,0001,000−800=600 秒T_{\text{清空}} \approx \frac{120{,}000}{1{,}000-800} = 600\ \text{秒}

也就是约 10 分钟。

如果有效完成速率只有每秒 700 条,就不存在正向清空能力。此时继续修改队列展示方式或扩大客户端缓冲,并不能改变这个容量缺口。

3. 最后用业务时间线验收

建议在应用中记录任务产生、消费回调进入、业务开始、业务完成几个时点。

其中,“消费回调进入时间减去任务产生时间”不能直接命名为纯 Broker 排队时间,因为它还可能包含发布前等待、网络传输和客户端库内部等待;跨机器比较时间戳,也需要考虑时钟偏差。

而“业务开始减去回调进入”,可以帮助观察应用自身的排队;“业务完成减去业务开始”,则帮助定位实际执行阶段。

验收的核心应是:

在不降低可靠性、不恶化资源风险的前提下,有效完成速度是否提高,端到端延迟是否下降。

八、总结

RabbitMQ 消息积压不是一个只能通过“增加消费者”解决的问题。

Ready 降低,可能只是消息变成了 Unacked;Prefetch 增大,可能只是客户端等待区变长;实例增加,也可能仍然受制于同一个下游瓶颈或执行串行点。

正确的排障顺序,是先识别消息状态,再核对实际处理并发和下游能力,最后通过受控实验调整预取窗口。

真正值得优化的,不是消息离开 Ready 状态的速度,而是任务被可靠完成的速度。

0 条评论
如果你觉得文章对你有帮助,那就请作者喝杯咖啡吧☕
微信
支付宝
  0 条评论