前面我们已经把 RabbitMQ 的主要功能基本讲完了:
Producer
↓
Exchange
↓
Queue
↓
Consumer
以及消息可靠性:
Confirm
↓
持久化
↓
ACK
↓
Retry
↓
DLQ
↓
幂等
但真正把 RabbitMQ 用到业务中以后,还会遇到一个非常现实的问题:
生产者发送消息太快,消费者处理不过来怎么办?
比如:
Producer:
每秒发送 5000 条消息
Consumer:
每秒只能处理 1000 条
那么每秒就会多出来:
5000 – 1000
= 4000 条
这些消息只能暂时留在 Queue 中。
于是:
第 1 秒:积压 4000
第 2 秒:积压 8000
第 3 秒:积压 12000
…
这就是:
消息积压
这一篇就来看看 RabbitMQ 消息积压以后,到底应该怎么处理。
1. 什么是消息积压?
正常情况下:
Producer
↓
Queue
↓
Consumer
生产速度和消费速度基本平衡:
生产:1000 条/s
消费:1000 条/s
Queue 中不会长期保存大量消息。
但是如果:
生产速度
>
消费速度
就会:
Producer
↓↓↓↓↓↓↓↓↓↓↓
┌────────────────┐
│ Message │
│ Message │
│ Message │
│ Message │
│ Message │
└────────────────┘
Queue
↓
Consumer
Queue 越堆越多。
本质上:
消息积压并不是 RabbitMQ 突然出现了问题,而是生产速度长期大于消费速度。
可以简单写成:
积压增长速度
=
生产速度 – 消费速度
例如:
生产:3000/s
消费:2000/s
那么:
每秒积压 1000 条
一分钟:
1000 × 60
= 60000 条
一个小时:
60000 × 60
= 3600000 条
也就是:
360 万条
所以如果消费能力长期跟不上,积压增长会非常快。
2. 消息积压有什么影响?
可能有人会想:
RabbitMQ 本来不就是用来存消息的吗?堆着不就行了?
短时间积压:
确实很正常
比如秒杀活动:
瞬间大量请求
↓
RabbitMQ
↓
慢慢消费
这本来就是我们前面讲过的:
削峰
真正的问题在于:
消息一直进来,而消费者长期追不上。
这时会出现:
消息延迟越来越高
内存和磁盘压力增大
业务处理越来越晚
RabbitMQ 节点压力增加
例如一条:
发送短信
的消息本来应该:
1 秒后处理
但前面已经积压了:
100 万条
那消费者真正处理到它的时候可能已经过去:
几分钟
甚至更久
所以很多时候:
消息积压首先表现出来的不是“消息丢失”,而是业务延迟越来越严重。
3. 怎么判断 RabbitMQ 是否出现消息积压?
RabbitMQ Management UI 中可以看到 Queue 的各种指标。
比较重要的包括:
Ready
Unacked
Consumers
Publish Rate
Deliver / Ack Rate
RabbitMQ 官方当前文档也说明,Management UI 可以监控 Queue Length、消息进入和离开速率、Consumer 数量以及消息不同状态等指标。
4. Ready 是什么?
Ready
表示:
还在 Queue 中等待发送给 Consumer 的消息。
例如:
Ready = 100000
说明现在还有:
10 万条消息
等待被消费。
如果这个数字:
持续增加
通常就说明:
消费速度
<
生产速度
已经开始产生积压。
5. Unacked 又是什么?
假设 RabbitMQ 已经把消息发送给 Consumer:
Queue
↓
Consumer
但是 Consumer 还没有:
ACK
这些消息就属于:
Unacked
也就是:
已经发送给 Consumer
但是还没有确认完成
所以:
Ready
和:
Unacked
不要混在一起。
可以理解成:
Ready
还在 RabbitMQ Queue 里排队
而:
Unacked
已经交给 Consumer
但是 Consumer 还没处理完
RabbitMQ 官方同样把 Queue 中消息区分为 Ready 和 Delivered-but-not-yet-acknowledged 两种主要状态。
6. Ready 很多和 Unacked 很多分别说明什么?
这个判断在排查问题时非常有用。
Ready 一直很多
例如:
Ready = 500000
Unacked = 10
说明:
大量消息还在 Queue 里
消费者拿消息的能力不足。
可能需要考虑:
增加 Consumer
提高并发
调整 Prefetch
提高单条消息处理速度
Unacked 特别多
例如:
Ready = 1000
Unacked = 100000
说明大量消息已经被 Consumer 拿走:
RabbitMQ
↓↓↓↓↓
Consumer
但是迟迟没有 ACK。
这时候就应该重点检查:
业务处理是不是很慢?
数据库是不是卡了?
第三方接口是不是很慢?
Prefetch 是不是设置得太大?
所以:
Ready
更像:
RabbitMQ 中排队的任务
而:
Unacked
更像:
Consumer 手里还没做完的任务
7. 最直接的解决办法:提高消费速度
既然积压产生的根本原因是:
生产速度 > 消费速度
那么解决方向其实非常明确:
提高消费速度
例如原来:
Producer:5000/s
Consumer:1000/s
如果把 Consumer 提高到:
6000/s
那么:
生产:5000/s
消费:6000/s
不但不会继续积压,还可以逐渐消化之前的历史积压。
所以第一件事情不是:
疯狂修改 RabbitMQ 参数
而应该先看:
Consumer 为什么这么慢?
8. 先优化 Consumer 自己
例如消费者:
@RabbitListener(queues = "order.queue")
public void consume(OrderMessage message) {
queryDatabase();
callRemoteService();
updateDatabase();
sendHttpRequest();
}
假设:
查数据库:100ms
调用第三方:500ms
修改数据库:100ms
HTTP 请求:300ms
一条消息可能需要:
1 秒
一个 Consumer:
1 秒只能处理 1 条
那消息当然很容易积压。
所以首先应该看看:
SQL 是否太慢?
有没有大量重复查询?
第三方接口是否阻塞?
能否批量处理?
有没有不必要的同步操作?
如果:
1000ms / 条
优化成:
100ms / 条
那么单个 Consumer 的吞吐量理论上就可能提高很多。
所以:
解决消息积压最优先的方案,通常还是减少每条消息的处理时间。
9. 第二种方案:增加消费者并发
假设业务代码暂时没有明显优化空间。
一个 Consumer:
100 条/s
那可以增加消费者数量:
Consumer 1:100/s
Consumer 2:100/s
Consumer 3:100/s
Consumer 4:100/s
理论总消费能力就可能提高到:
400 条/s
结构:
┌→ Consumer 1
│
order.queue ───┼→ Consumer 2
│
├→ Consumer 3
│
└→ Consumer 4
同一个 Queue 的多个 Consumer 会竞争消费消息。
一条消息正常只会交给其中一个 Consumer。
所以:
增加 Consumer
就是最常见的横向提高消费能力的方法之一。
10. Spring Boot 怎么增加消费者并发?
Spring Boot 可以直接配置:
spring:
rabbitmq:
listener:
simple:
concurrency: 3
max-concurrency: 10
其中:
concurrency = 3
表示:
初始最少创建 3 个消费者
而:
max-concurrency = 10
表示:
负载较高时最多可以扩展到 10 个消费者
Spring AMQP 当前文档中,SimpleMessageListenerContainer 支持 concurrentConsumers 和 maxConcurrentConsumers,并能够根据负载在两者之间动态增加或减少消费者;Spring Boot 也提供了对应的 spring.rabbitmq.listener.simple.concurrency 和 max-concurrency 配置。
可以简单理解成:
平时:
Consumer
Consumer
Consumer
消息变多:
Consumer
Consumer
Consumer
Consumer
Consumer
…
11. Consumer 越多越好吗?
当然不是。
比如 Consumer 最终都需要操作:
MySQL
数据库最多只能稳定处理:
2000 次请求/s
你原来:
5 个 Consumer
已经把数据库跑到:
80% CPU
这时候突然增加到:
50 个 Consumer
结果可能不是消费速度提高 10 倍。
而是:
数据库连接池耗尽
SQL 变慢
大量请求超时
Consumer 反而越来越慢
最终:
RabbitMQ 没挂
Consumer 没挂
MySQL 先挂了
因此增加 Consumer 的前提是:
下游资源还能承受更多并发。
需要一起考虑:
CPU
数据库
Redis
第三方接口
连接池
网络带宽
所以消息积压不能只看 RabbitMQ。
12. Prefetch 是什么?
讲完消费者并发,就到了 RabbitMQ 中非常重要的参数:
Prefetch
它特别容易被理解错。
先假设:
prefetch = 10
它并不是说:
Consumer 一次批量执行业务 10 条
也不是:
创建 10 个线程
而是:
RabbitMQ 最多允许这个 Consumer 同时持有一定数量的、已经投递但尚未 ACK 的消息。
RabbitMQ 当前官方文档把 Consumer Prefetch 定义为限制一个 Consumer 可以拥有多少条 outstanding / unacknowledged deliveries。RabbitMQ 会把 prefetch count 分别应用到各个 Consumer。
例如:
prefetch = 3
那么:
RabbitMQ
↓
Message 1
Message 2
Message 3
↓
Consumer
如果这 3 条都还没有 ACK:
Unacked = 3
RabbitMQ 就不会继续无限地往这个 Consumer 手里塞消息。
等 Consumer:
ACK Message 1
腾出一个位置:
Unacked = 2
RabbitMQ 才可以继续发送下一条。
13. 为什么需要 Prefetch?
假设完全不限制。
RabbitMQ 中有:
100000 条消息
Consumer A 网络很快:
RabbitMQ
↓↓↓↓↓↓↓↓↓↓↓↓↓
Consumer A
大量消息全部被推过去。
结果:
Consumer A 手里堆了几万条
Consumer B 却没分到多少
而且大量:
Unacked
消息会占用:
Consumer 内存
RabbitMQ 资源
所以 Prefetch 相当于告诉 RabbitMQ:
别一次给我太多,我处理完一些,你再继续给。
这其实是一种:
流量控制
14. Prefetch 太小会怎么样?
假设:
prefetch = 1
流程:
RabbitMQ
↓
Message 1
↓
Consumer
↓
处理
↓
ACK
↓
RabbitMQ
↓
Message 2
Consumer 每次只能拥有:
1 条未确认消息
优点:
消息分配比较公平
Consumer 内存压力小
更适合很慢或者很大的消息
但如果 Consumer 本身处理非常快:
处理完
↓
等下一条
↓
再处理
频繁等待 Broker 继续投递,可能无法充分利用消费能力。
因此:
prefetch 太小
可能限制吞吐量。
15. Prefetch 太大又会怎么样?
假设:
prefetch = 10000
RabbitMQ 可以提前给一个 Consumer:
10000 条未 ACK 消息
如果每条消息:
1MB
理论上就可能形成非常大的客户端内存压力。
而且多个 Consumer 时:
Consumer A
可能提前拿走大量消息。
其他 Consumer:
Consumer B
Consumer C
反而没有多少消息可以处理。
RabbitMQ 官方也提醒,高达数千甚至更多的 Prefetch 可能带来和自动确认类似的过载问题,并增加 Broker 端未确认消息相关的内存使用。
所以:
Prefetch 不是越大越快。
16. Spring Boot 怎么设置 Prefetch?
可以:
spring:
rabbitmq:
listener:
simple:
prefetch: 20
表示每个 Consumer:
最多拥有 20 条未 ACK 消息
Spring Boot 当前配置文档对这个参数的描述就是:
Maximum number of unacknowledged messages
that can be outstanding at each consumer
也就是每个 Consumer 最多允许多少条尚未确认的消息。
例如同时配置:
spring:
rabbitmq:
listener:
simple:
concurrency: 5
prefetch: 20
那么理论上:
5 个 Consumer
每个最多:
20 条 Unacked
整体最多可能存在大约:
5 × 20
= 100 条
已经投递给这些 Consumer、但尚未确认的消息。
17. Spring AMQP 默认 Prefetch 是多少?
当前 Spring AMQP 文档中:
默认 prefetch = 250
这是为了让处理速度较快的 Consumer 保持较高利用率。
不过官方同时指出,在下面这些场景应该考虑降低 Prefetch:
消息体很大
消息处理很慢
要求严格顺序
多个 Consumer 下希望分配更加均匀
尤其严格顺序场景,官方建议将 Prefetch 降到:
1
所以博客里不要写成:
prefetch 一定设置为多少最好
因为它并不存在一个适用于所有业务的固定值。
18. Prefetch 和 Concurrency 到底有什么区别?
这是这一篇最需要区分的地方。
Concurrency
控制:
有多少 Consumer 同时干活
例如:
concurrency = 5
就是:
Consumer 1
Consumer 2
Consumer 3
Consumer 4
Consumer 5
Prefetch
控制:
每个 Consumer
最多提前拿多少条未 ACK 消息
例如:
prefetch = 20
就是:
Consumer 1 → 最多 20 条 Unacked
Consumer 2 → 最多 20 条 Unacked
…
所以:
Concurrency
↓
有多少个人干活
Prefetch
↓
每个人手里最多可以拿多少任务
这个比喻非常好记。
19. 用餐厅理解 Concurrency 和 Prefetch
假设餐厅厨房:
订单 Queue
厨师:
Consumer
Concurrency = 3
代表:
有 3 个厨师
厨师 A
厨师 B
厨师 C
Prefetch = 5
代表:
每个厨师桌子上最多可以提前放 5 张订单
于是:
Queue
↓
┌──────────┐
│ 厨师 A │ ← 最多 5 单
├──────────┤
│ 厨师 B │ ← 最多 5 单
├──────────┤
│ 厨师 C │ ← 最多 5 单
└──────────┘
如果:
Prefetch 很大
就像一次给某个厨师桌上堆:
500 张订单
不仅桌子堆满了,也可能导致其他厨师没有任务。
如果:
Prefetch 太小
又像厨师做完一道菜以后:
必须重新跑到前台拿下一张订单
也可能降低效率。
所以需要找到:
合适的值
而不是:
越大越好
20. 第三种方案:增加应用实例
假设现在一台机器运行:
OrderConsumer
即使线程并发开得再多:
concurrency = 50
最终还是共享:
一台机器的 CPU
内存
网络
当单机已经到达瓶颈,就需要:
横向扩容
例如原来:
RabbitMQ
↓
Application A
扩容以后:
┌→ Application A
│
RabbitMQ ────┼→ Application B
│
└→ Application C
三个应用实例都监听:
order.queue
RabbitMQ 会把消息分发给不同 Consumer。
因此总消费能力可以进一步提升。
21. 线程扩容和机器扩容有什么区别?
增加 Concurrency
属于:
单个应用内部扩容
比如:
1 个 JVM
里面 10 个 Consumer
增加应用实例
属于:
横向扩容
例如:
JVM A → 10 Consumers
JVM B → 10 Consumers
JVM C → 10 Consumers
于是:
总共 30 个 Consumer
可以继续提高整体消费能力。
所以通常:
先优化单条处理速度
↓
适当提高 Consumer 并发
↓
单机达到瓶颈
↓
增加应用实例
这是比较自然的扩容路径。
22. 临时出现几十万条积压怎么办?
假设平时系统运行正常。
突然某个 Consumer:
宕机 30 分钟
RabbitMQ 已经堆积:
50 万条消息
现在 Consumer 恢复了。
但正常生产速度:
2000/s
Consumer 正常消费速度也是:
2000/s
那意味着:
新消息刚好全部消费掉
之前那:
50 万条
永远也消化不完。
这时候必须让:
消费速度 > 新消息生产速度
例如临时扩容:
消费能力:
8000/s
新消息:
2000/s
那么每秒可以额外清理:
8000 – 2000
= 6000 条历史消息
原来:
500000 条
理论上大约:
500000 ÷ 6000
≈ 83 秒
能够逐渐消化。
所以处理历史积压的关键是:
必须暂时提供超过当前生产速度的额外消费能力。
23. 加 Consumer 还是不够怎么办?
假设已经:
Consumer × 50
但积压仍然越来越多。
这时候说明问题可能已经不只是 Consumer 数量。
例如:
数据库已经到极限
Redis 到极限
下游服务限流
网络带宽到极限
继续加 Consumer 反而:
把更多压力传递给下游
所以应该重新分析整个链路:
RabbitMQ
↓
Consumer
↓
MySQL
↓
Redis
↓
Third Party API
到底哪一段才是真正瓶颈。
这就是为什么:
RabbitMQ 消息积压其实是系统吞吐能力不足的表现,而不一定是 RabbitMQ 自己的问题。
24. 能不能限制 Queue 最大长度?
可以。
RabbitMQ 支持给 Queue 配置:
最大消息数量
最大消息总字节数
也就是:
x-max-length
x-max-length-bytes
当前官方文档还支持不同的溢出策略,例如:
drop-head
reject-publish
reject-publish-dlx
RabbitMQ 官方更推荐通过 Policy 管理 Queue Length Limit。
但需要注意:
设置 Queue 最大长度并不是“解决消息积压”。
它只是防止:
Queue 无限增长
如果生产速度还是:
5000/s
消费速度还是:
1000/s
根本矛盾依然存在。
甚至达到限制以后,还需要考虑:
消息被拒绝怎么办?
消息被丢弃怎么办?
是否进入 DLX?
所以 Queue Limit 更像:
保护措施
而不是:
性能优化措施
25. 怎么判断到底该增加 Consumer 还是调 Prefetch?
RabbitMQ Management UI 中还有一个很有价值的指标:
Consumer Capacity
它以前叫:
Consumer Utilisation
RabbitMQ 当前文档说明,这个指标可以帮助判断 Queue 是否可能通过:
增加 Consumer
缩短 Consumer 处理时间
提高 Prefetch
获得更高的投递能力。
所以实际排查时不要只盯着:
Queue 有多少条消息
还应该一起看:
Ready
Unacked
Publish Rate
Ack Rate
Consumer Count
Consumer Capacity
这些指标组合起来才能判断真正的问题。
26. 一套比较实用的消息积压排查思路
看到:
Queue 消息越来越多
不要上来就:
prefetch = 10000
或者:
concurrency = 100
更合理的排查顺序是:
发现消息积压
↓
生产速度是多少?
↓
消费速度是多少?
↓
Consumer 是否异常?
↓
单条消息处理为什么慢?
↓
CPU / DB / Redis / 第三方是否已经到瓶颈?
↓
调整消费者并发
↓
合理调整 Prefetch
↓
需要时水平扩容实例
如果最终发现:
一个 Queue 本身已经成为极高吞吐瓶颈
则可能需要进一步考虑:
拆分 Queue
业务分片
RabbitMQ Streams
RabbitMQ 当前官方文档也指出,单个 Queue Replica 的热点执行路径受到单个 CPU Core 的限制,对于逼近单 Queue 吞吐极限的场景,应考虑多个 Queue,或者 Streams / Partitioned Streams。
不过这一部分已经属于比较深入的性能设计了。
入门阶段知道:
一个 Queue 也不是无限吞吐
即可。
27. 把消息积压问题总结成一个公式
其实这篇最核心的东西非常简单:
生产速度 > 消费速度
↓
消息积压
解决问题:
提高消费速度
主要有三种方式:
① 单条消息处理更快
② 更多 Consumer 并发
③ 更多应用实例
而:
Prefetch
主要负责:
控制每个 Consumer 手里
最多有多少条未 ACK 消息
它可以影响吞吐和消息分配,但:
Prefetch 本身不能把一个很慢的业务逻辑突然变快。
总结
RabbitMQ 消息积压的根本原因就是:
Producer
生产太快
>
Consumer
消费太慢
短时间积压:
属于 RabbitMQ 削峰的正常现象
但长期不断增长:
说明消费能力不足
解决问题时,可以按照:
优化业务代码
↓
提高 Consumer 并发
↓
合理调整 Prefetch
↓
增加应用实例
逐步处理。
其中一定要区分:
Concurrency
↓
控制有多少 Consumer 同时工作
和:
Prefetch
↓
控制每个 Consumer
最多持有多少条未 ACK 消息
可以用一句话记:
Concurrency 决定有多少人干活,Prefetch 决定每个人手里最多拿多少任务。
而排查消息积压时,不要只看:
Queue Length
还应该关注:
Ready
Unacked
Publish Rate
Ack Rate
Consumer Count
Consumer Capacity
最后需要记住:
消息积压很多时候只是表象,真正的问题往往是整条消费链路的处理能力跟不上生产速度。
所以 RabbitMQ 性能优化并不是单纯:
把某个参数调大
而是找到:
真正的系统瓶颈
再决定到底应该:
优化代码
增加消费者
调整 Prefetch
扩容机器
还是优化数据库和下游服务

