163 lines
6.8 KiB
Markdown
163 lines
6.8 KiB
Markdown
---
|
||
tags: [mq/rabbitmq, pull-model, qos, prefetch, backpressure, consumer-balance]
|
||
create time: 2026-08-08 18:53
|
||
update time: 2026-08-08 18:53
|
||
---
|
||
|
||
# 推拉结合消费模式
|
||
|
||
## 概述
|
||
|
||
RabbitMQ 的 Consumer API 名为 `basic.consume`(拉取),但实际运行模型是 Server Push(服务端推送)。这个"名不副实"的设计源于 AMQP 协议的早期约定,也是理解 RabbitMQ 消费模型的关键起点。深入掌握预取策略(QoS)、背压处理和多消费者负载均衡机制,才能在吞吐量和可靠性之间找到最佳平衡点。
|
||
|
||
## 核心原理
|
||
|
||
### 消费模型的本质
|
||
|
||
RabbitMQ 的消费入口确实是拉取操作(`basic.consume` 或 `basic.get`),两者行为截然不同:
|
||
|
||
| 模式 | 方法 | 行为 | 特点 |
|
||
|------|------|------|-----|
|
||
| 推模式(Push) | `basic.consume` | 建立订阅后,Broker 持续主动投递 | 默认模式,适用于大多数场景 |
|
||
| 拉模式(Pull) | `basic.get` | 每次调用阻塞获取单条消息 | 适用于管理界面、健康检查等非实时场景 |
|
||
|
||
`basic.consume` 的工作流程:
|
||
|
||
```mermaid
|
||
sequenceDiagram
|
||
participant C as 消费者
|
||
participant RMQ as RabbitMQ Broker
|
||
participant Q as Queue
|
||
|
||
C->>RMQ: basic.consume(queue="orders", no_ack=false, prefetch=10)
|
||
RMQ-->>C: consume-ok(consumer_tag="ct1")
|
||
Note over C,RMQ: 流控通道已建立
|
||
loop 消息到达队列
|
||
RMQ->>C: basic.deliver(tag=1)<br/>queue="orders"
|
||
C->>C: 处理消息...
|
||
C->>RMQ: basic.ack(tag=1)
|
||
C->>RMQ: window refill (available=10)
|
||
end
|
||
```
|
||
|
||
消费者调用 `basic.consume` 后,RabbitMQ 会在该 Channel 上建立一个流控窗口。每当收到 `basic.deliver` 消息时,窗口减少;消费者发送 `basic.ack` 后窗口恢复。Window 耗尽时,Broker 暂停投递直到窗口回收——这就是**内建背压机制**。
|
||
|
||
### Prefetch Count(QoS)设置原理
|
||
|
||
`prefetch_count` 定义了单个 Channel 上允许的最大未确认消息数。它是流量控制的调节阀。
|
||
|
||
设 prefetch = N 时:
|
||
|
||
1. Broker 可以在没有收到 ack 的情况下,最多向该 Channel 发送 N 条消息
|
||
2. 这 N 条消息累积在未确认(unacked)状态
|
||
3. 当 ack 一条消息后,unacked 数减一,Broker 再补发一条
|
||
|
||
```go
|
||
// Go 代码示例:Qos 配置
|
||
// prefetch=1 启用公平分发(Fair Dispatch)
|
||
ch.Qos(1, // prefetch count
|
||
0, // prefetch size(字节限制,0表示不限制)
|
||
false) // global=false 作用于单个 channel
|
||
```
|
||
|
||
为什么 prefetch=1 能实现公平分发?假设消费者 A 处理快、B 处理慢:
|
||
|
||
- 如果 prefetch=N(大数值),Broker 可能在短时间内把 N 条消息全部发给 A,而 B 空闲等待
|
||
- 如果 prefetch=1,A 每处理完一条、发一个 ack,才会收到下一条,自然形成了速度匹配
|
||
|
||
> [!WARNING]
|
||
> `global=true` 是遗留参数,在所有现代客户端中建议设为 `false`。它会对整个连接(而非单个 Channel)生效,在多 Channel 场景下会导致意外行为。
|
||
|
||
### 预取策略对吞吐量和内存的影响
|
||
|
||
| 预取值 | 吞吐量表现 | 内存占用 | 适用场景 |
|
||
|--------|-----------|---------|---------|
|
||
| 1 | 较低(串行节奏) | 最低 | 处理耗时差异大的消费者组 |
|
||
| 10~50 | 较高 | 中等 | 通用生产环境默认值 |
|
||
| 100~500 | 最高 | 较高 | 短处理、低延迟需求 |
|
||
| 0(不设置) | 取决于默认值(通常为 0 即无限) | 可能溢出 | **严禁在生产中使用** |
|
||
|
||
prefetch 越大意味着 Broker 可以积压更多消息到消费者内存中。这降低了每次投递的往返开销(减少网络 RTT),但同时增加了消费者的内存压力和崩溃恢复时的损失量。
|
||
|
||
### 背压处理机制
|
||
|
||
RabbitMQ 的背压在两个层面上工作:
|
||
|
||
**Layer 1:Channel 级别的 Window**(prefetch 控制)
|
||
- 消费者窗口耗尽时,Broker 停止向该 Channel 投递
|
||
- 消费者恢复正常后自动解封
|
||
|
||
**Layer 2:Queue 级别的持久化**
|
||
- 当消费者整体处理能力低于生产者速度时,消息堆积在 Queue 中
|
||
- Queue 受 `x-max-length` 约束时可拒绝新消息
|
||
- 不受限时占用磁盘空间,可能拖垮 Broker
|
||
|
||
```
|
||
生产者速度 > 消费者速度 = 消息在 Queue 中排队等待
|
||
↓
|
||
Queue 满 + 有限长 = 新消息被拒
|
||
↓
|
||
走 DLX 或触发告警
|
||
```
|
||
|
||
### 多消费者并发与负载均衡
|
||
|
||
同一个 Queue 可以有多个消费者(构成 Consumer Group),RabbitMQ 以 Round-Robin 方式将消息均匀分配给各消费者:
|
||
|
||
```mermaid
|
||
graph LR
|
||
Q["Order Queue"] -->|msg 1| C1["Consumer A"]
|
||
Q -->|msg 2| C2["Consumer B"]
|
||
Q -->|msg 3| C1
|
||
Q -->|msg 4| C2
|
||
Q -->|msg 5| C1
|
||
Q -->|msg 6| C2
|
||
style Q fill:#fff3e0
|
||
style C1 fill:#e8f5e9
|
||
style C2 fill:#e8f5e9
|
||
```
|
||
|
||
前提条件是 `prefetch_count` 足够大(否则 prefetch=1 时只有一个消费者活跃,另一个空等)。Round-Robin 的分发方式是平均的但不一定是公平的——处理快的消费者实际上承担了更多工作量,这正是 prefetch=1 要解决的问题。
|
||
|
||
### 预取策略调优经验
|
||
|
||
高吞吐场景(如日志采集、指标上报):
|
||
|
||
- `prefetch=100~500`,消费者端做好消息批处理
|
||
- 处理逻辑尽量无状态,避免加锁
|
||
- Monitor unacked 数量和 Queue depth,发现趋势性增长立即扩容
|
||
|
||
低延迟场景(如实时通知、在线游戏):
|
||
|
||
- `prefetch=1~10`,让消费者尽快处理完再拿下一条
|
||
- 消费者数量尽可能多,分散到不同机器
|
||
- 监控 p99 延迟而非平均值,个别慢消费者可能拉高尾部
|
||
|
||
> [!TIP]
|
||
> 调优黄金法则:先设 prefetch=1 观察各消费者的处理耗时分布,再根据 p99 耗时最高的那个消费者来估算合适的 batch size,初始设为预估值的 2~3 倍即可。
|
||
|
||
## 实践场景
|
||
|
||
**高并发商品详情页缓存预热**:每日定时从数据库全量读取商品信息,通过 MQ 推送到缓存集群预热。日终跑批期间数据量大但每条消息处理简单,采用 prefetch=200 配合批量更新 Redis pipeline,在保证不淹没缓存节点的前提下最大化吞吐。
|
||
|
||
**即时订单状态推送**:用户下单后需要实时更新前端页面。这里 latency 优先级高于 throughput,设置 prefetch=3,消费者部署 10 个副本分布在不同的 Pod 中,配合 Manual ACK 确保状态变更的可靠性。
|
||
|
||
```go
|
||
// Go 代码示例:根据环境动态调整 Qos
|
||
func configureQos(env string) int {
|
||
switch env {
|
||
case "high-throughput":
|
||
return 200 // 大批量批处理
|
||
case "low-latency":
|
||
return 3 // 追求快速响应
|
||
default:
|
||
return 10 // 通用场景
|
||
}
|
||
}
|
||
```
|
||
|
||
## 扩展阅读
|
||
- [[Exchange 路由机制]]
|
||
- [[消息持久化与可靠性投递]]
|
||
- [[ACK 确认与死信队列]]
|