Files

159 lines
5.4 KiB
Markdown
Raw Permalink Normal View History

2026-05-24 20:51:06 +08:00
---
tags:
- MQ
create time: 2026-05-24 19:52
---
# MQ 背压与流控
## 概述
当生产速率持续超过消费速率时,消息队列会面临内存溢出、磁盘满甚至消息丢失的风险。背压(Backpressure)机制让上游感知下游的处理能力,流控则是系统在过载时保护自身的手段。本文从 Broker 端和 Consumer 端两个维度,剖析主流 MQ 的背压与流控策略。
## 正文
### 问题:生产太快,消费太慢
想象一个场景:大促期间订单量暴增,Producer 疯狂写入消息,而下游 Consumer 处理能力有限。如果不加控制,会发生什么?
1. **内存溢出**:Broker 将消息堆积在内存中,最终 OOM
2. **磁盘满**:持久化消息写满磁盘,Broker 宕机
3. **消息丢失**:触发淘汰策略(TTL / 队列满丢弃),消息悄无声息地消失
> [!question]
> 如果 MQ 天生就是一个缓冲区,消息堆积不是它的基本能力吗?为什么堆积到一定程度反而会出问题?
这涉及到一个关键认知:**缓冲区是有限的**。任何系统都有资源上限——内存、磁盘、CPU。无限制的堆积只是把问题延后,而不是解决。
### 背压(Backpressure)概念
背压的核心思想:**让上游感知下游的处理能力,主动降速**。
```
Producer → Broker → Consumer
↑ ↓
└──── 处理能力反馈 ────┘
```
这不是 MQ 的专利。TCP 的滑动窗口、HTTP/2 的流控、Reactive Streams 的 `request(n)` 都是背压的不同实现形式。
### Broker 端流控
#### RabbitMQ 的信用机制(Credit Flow)
RabbitMQ 使用 **信用机制** 控制消息流速。Producer 发送消息前需要有足够的 credit,Broker 处理完后归还 credit。当 Broker 积压过多,会暂停归还 credit,从而让 Producer 阻塞。
```
Producer --msg1--> Broker (credit: 10→9)
Producer --msg2--> Broker (credit: 9→8)
...积压严重...
Broker 暂停归还 credit
Producer 阻塞,停止发送
```
#### Kafka 的 Producer 端背压
Kafka 通过两个参数实现背压:
- `buffer.memory`:Producer 端缓冲区大小(默认 32MB)
- `max.block.ms`:缓冲区满时,`send()` 方法的阻塞时间
当缓冲区满且超过 `max.block.ms`,Producer 抛出 `TimeoutException`。这是硬性的背压信号。
```go
// Kafka Producer 配置中的背压参数
config := sarama.Config{}
config.Producer.RequiredAcks = sarama.WaitForAll
config.Producer.Return.Successes = true
// 缓冲区满时阻塞 1 秒,超时则报错
// 这就是 Kafka 的背压触发点
```
#### 内存告警与磁盘告警
多数 Broker 有内置的资源监控:
- **RabbitMQ**:内存高水位触发 flow control,阻塞所有连接
- **Kafka**:`log.retention.bytes` 限制分区大小,磁盘满时拒绝写入
- **RocketMQ**:`diskMaxUsedSpaceRatio` 触发磁盘保护
### Consumer 端反压
Consumer 端同样需要流控,核心手段包括:
**拉取速率控制**:Pull 模式天然支持背压——Consumer 按自己的节奏拉取,拉多少处理多少。RabbitMQ 的 `basicQos(prefetchCount)` 就是限制未 ACK 消息数的典型手段。
**处理能力反馈**:Consumer 可以动态上报自己的处理延迟或队列深度,上游据此调整推送速率。
**动态调整消费并发**:根据处理延迟自动扩缩 Consumer 实例数。Kubernetes HPA 基于队列深度的自动伸缩是常见方案。
```mermaid
graph LR
P["Producer"] -->|"生产消息"| B["Broker"]
B -->|"推送/拉取"| C["Consumer"]
C -->|"ACK / 处理反馈"| B
B -->|"Credit / 阻塞信号"| P
style P fill:#4CAF50,color:#fff
style B fill:#2196F3,color:#fff
style C fill:#FF9800,color:#fff
```
### 限流降级策略
当背压来不及响应时,需要更积极的流控手段:
**令牌桶(Token Bucket)**:以恒定速率产生令牌,请求必须持有令牌才能通过。允许一定的突发流量(桶内预存令牌)。
**漏桶(Leaky Bucket)**:请求进入桶中,以恒定速率流出。严格平滑流量,但不允许突发。
**动态调整生产速率**:根据 Broker 的健康指标(队列深度、内存使用率)动态调整 Producer 的发送速率。
```go
// 带背压控制的 Producer 示例
// 利用 channel 的天然阻塞特性实现背压
func producer(ch chan<- string, done <-chan struct{}) {
for {
select {
case <-done:
return
case ch <- "message":
// channel 满时自动阻塞,实现背压
// 上游感知到下游处理不过来,自然降速
}
}
}
func consumer(ch <-chan string) {
for msg := range ch {
// 模拟慢消费
time.Sleep(100 * time.Millisecond)
_ = msg
}
}
func main() {
// 有界 channel 就是一个天然的背压缓冲区
ch := make(chan string, 100)
done := make(chan struct{})
go producer(ch, done)
go consumer(ch)
// 当 consumer 处理不过来时,
// channel 满后 producer 自动阻塞
time.Sleep(5 * time.Second)
close(done)
}
```
这段代码的精髓在于 `make(chan string, 100)`——有界 channel 就是一个天然的背压装置。当缓冲区满时,发送方自动阻塞,不需要额外的信号传递。
> [!question]
> 令牌桶和漏桶看起来很像,它们的核心区别是什么?什么场景下该用哪个?
## 关联笔记
- [[32-MQ-请求-回复模式]]
- [[33-MQ-与流处理]]
- [[34-事件驱动架构-EDA]]