欢迎光临
我们一直在努力

RabbitMQ 消息积压怎么办?Prefetch、消费者并发与扩容

前面我们已经把 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

扩容机器

还是优化数据库和下游服务

赞(0)
未经允许不得转载:171主机测评 » RabbitMQ 消息积压怎么办?Prefetch、消费者并发与扩容
分享到: 更多 (0)

评论 抢沙发

  • 昵称 (必填)
  • 邮箱 (必填)
  • 网址