vault backup: 2026-05-24 20:51:06
This commit is contained in:
@@ -0,0 +1,153 @@
|
||||
---
|
||||
tags: [MQ, 消息队列, 设计模式, 负载均衡]
|
||||
create time: 2026-05-24 19:52
|
||||
---
|
||||
|
||||
# MQ Competing Consumers 模式
|
||||
|
||||
## 概述
|
||||
|
||||
Competing Consumers(竞争消费者)模式通过多个消费者同时消费同一个队列来实现负载均衡,是提升消息处理吞吐量最直接的手段。本文深入讲解该模式的工作原理、分区分配算法、Consumer Rebalance 机制及其常见问题与优化方案。
|
||||
|
||||
## 正文
|
||||
|
||||
### 模式定义
|
||||
|
||||
传统单消费者模型中,一个队列只有一个消费者按顺序处理消息。当消息量增大时,单个消费者成为瓶颈。Competing Consumers 模式的核心思想很简单:**让多个消费者"抢"同一个队列的消息**,谁抢到谁处理,从而实现水平扩展。
|
||||
|
||||
```mermaid
|
||||
graph LR
|
||||
P["Producer"] --> Q["Queue"]
|
||||
Q --> C1["Consumer 1"]
|
||||
Q --> C2["Consumer 2"]
|
||||
Q --> C3["Consumer 3"]
|
||||
|
||||
style P fill:#4A90D9,color:#fff
|
||||
style Q fill:#F5A623,color:#fff
|
||||
style C1 fill:#6EC1E0,color:#fff
|
||||
style C2 fill:#6EC1E0,color:#fff
|
||||
style C3 fill:#6EC1E0,color:#fff
|
||||
```
|
||||
|
||||
### 与传统单消费者的对比
|
||||
|
||||
| 维度 | 单消费者 | Competing Consumers |
|
||||
|------|---------|---------------------|
|
||||
| 吞吐量 | 受单机性能限制 | 线性扩展(理论上) |
|
||||
| 消息顺序 | 天然有序 | 需要额外保障(分区有序) |
|
||||
| 可用性 | 消费者故障则停摆 | 单个消费者故障不影响整体 |
|
||||
| 复杂度 | 低 | 高(需要处理 Rebalance) |
|
||||
|
||||
> [!question]
|
||||
> Competing Consumers 提升了吞吐量,但牺牲了全局顺序性。如果你的业务需要"同一用户的操作有序",该如何在分区级别保障顺序?
|
||||
|
||||
### 分区分配算法详解
|
||||
|
||||
在 Kafka 中,Consumer Group 内的消费者通过分区分配算法决定"谁消费哪个 Partition"。不同的分配策略直接影响负载均衡效果和 Rebalance 开销。
|
||||
|
||||
#### Range(范围分配)
|
||||
|
||||
将 Topic 的 Partition 按序号范围分配给消费者。例如 6 个 Partition、3 个消费者:Consumer 0 拿到 [0,1],Consumer 1 拿到 [2,3],Consumer 2 拿到 [4,5]。
|
||||
|
||||
优点是简单直观,但如果多个 Topic 使用 Range 策略,可能导致某个消费者被分配到多个 Topic 的"头部"分区,造成负载不均。
|
||||
|
||||
#### Round-Robin(轮询分配)
|
||||
|
||||
将所有 Partition 轮询分配给消费者。6 个 Partition、3 个消费者:Consumer 0 拿到 [0,3],Consumer 1 拿到 [1,4],Consumer 2 拿到 [2,5]。
|
||||
|
||||
分配更均匀,但要求同一 Consumer Group 内所有消费者订阅的 Topic 完全一致,否则会出现分配异常。
|
||||
|
||||
#### Sticky(粘性分配)
|
||||
|
||||
在 Round-Robin 的基础上增加"粘性":Rebalance 时尽量保留原有的分配关系,只迁移必要的 Partition。这大幅减少了 Rebalance 期间的 Partition 迁移量。
|
||||
|
||||
#### Cooperative Sticky(协作式 Rebalance)
|
||||
|
||||
传统 Rebalance 是"Stop-the-World"——所有消费者暂停消费,等待分配完成。Cooperative Sticky 采用增量式 Rebalance:只暂停被迁移的 Partition,其余 Partition 继续消费。
|
||||
|
||||
```mermaid
|
||||
graph TD
|
||||
subgraph Before["Rebalance 前"]
|
||||
B1["Consumer 1"] --> BP1["Partition 0, 1, 2"]
|
||||
B2["Consumer 2"] --> BP2["Partition 3, 4, 5"]
|
||||
end
|
||||
|
||||
subgraph After["Rebalance 后 (Cooperative Sticky)"]
|
||||
A1["Consumer 1"] --> AP1["Partition 0, 1"]
|
||||
A2["Consumer 2"] --> AP2["Partition 3, 4, 5"]
|
||||
A3["Consumer 3 (新加入)"] --> AP3["Partition 2"]
|
||||
end
|
||||
|
||||
Before --> After
|
||||
|
||||
style Before fill:#F5A623,color:#fff
|
||||
style After fill:#6EC1E0,color:#fff
|
||||
```
|
||||
|
||||
注意图中:Cooperative Sticky 只迁移了 Partition 2,其余 Partition 在 Rebalance 期间保持消费不中断。
|
||||
|
||||
### Consumer Rebalance 过程与问题
|
||||
|
||||
Rebalance 是 Competing Consumers 模式中最关键也最容易出问题的环节。
|
||||
|
||||
**触发条件**:消费者加入/离开 Group、消费者心跳超时、Topic Partition 数量变化。
|
||||
|
||||
**Stop-the-World 问题**:传统 Rebalance 期间,Group 内所有消费者必须停止消费,等待 Coordinator 完成重新分配。在分区数量多或消费者数量大时,这个过程可能持续数秒甚至数十秒。
|
||||
|
||||
**Rebalance 风暴**:当消费者处理消息过慢导致心跳超时,被踢出 Group 触发 Rebalance;Rebalance 期间积压更多消息;消费者重新加入后又因处理不过来被踢出——形成恶性循环。
|
||||
|
||||
### Static Membership 减少不必要的 Rebalance
|
||||
|
||||
Kafka 2.3 引入 Static Membership 机制:每个消费者配置固定的 `group.instance.id`。消费者短暂断开重连时,只要在 `session.timeout.ms` 内恢复,就不会触发 Rebalance。
|
||||
|
||||
这在容器化部署中特别有用——Pod 重启时不会引发整个 Consumer Group 的 Rebalance 风暴。
|
||||
|
||||
> [!question]
|
||||
> Rebalance 期间消费者会暂停消费,在高并发场景下这会造成什么问题?如何缓解?
|
||||
|
||||
### Go 代码:Kafka Consumer 分区分配配置
|
||||
|
||||
```go
|
||||
package main
|
||||
|
||||
import (
|
||||
"github.com/segmentio/kafka-go"
|
||||
)
|
||||
|
||||
func main() {
|
||||
// 创建 Reader 时指定 Consumer Group 和分配策略
|
||||
r := kafka.NewReader(kafka.ReaderConfig{
|
||||
Brokers: []string{"localhost:9092"},
|
||||
Topic: "orders",
|
||||
GroupID: "order-service",
|
||||
// 使用 Cooperative Sticky 分配策略,减少 Rebalance 迁移
|
||||
GroupBalancers: []kafka.GroupBalancer{
|
||||
kafka.CooperativeGroupBalancer{},
|
||||
},
|
||||
// 静态成员:Pod 重启时不触发 Rebalance
|
||||
GroupInstanceID: "order-consumer-pod-0",
|
||||
})
|
||||
defer r.Close()
|
||||
|
||||
for {
|
||||
msg, err := r.ReadMessage(context.Background())
|
||||
if err != nil {
|
||||
// 处理错误,注意 Rebalance 期间会返回特定错误
|
||||
break
|
||||
}
|
||||
process(msg)
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
关键配置说明:
|
||||
- `GroupBalancers` 设置分配策略,`CooperativeGroupBalancer` 对应 Cooperative Sticky。
|
||||
- `GroupInstanceID` 设置静态成员 ID,同一个 Pod 重启后保持相同的 ID,避免触发不必要的 Rebalance。
|
||||
|
||||
## 关联笔记
|
||||
|
||||
- [[02-消息模型/3-MQ-消息模型|MQ 消息模型]]
|
||||
- [[07-主流MQ对比/22-Kafka|Kafka]]
|
||||
- [[05-可靠性保障/15-MQ-顺序性保障|MQ 顺序性保障]]
|
||||
- [[10-监控与运维/37-MQ-消费积压治理|MQ 消费积压治理]]
|
||||
- [[12-架构与实战/49-MQ-客户端连接管理|MQ 客户端连接管理]]
|
||||
@@ -0,0 +1,174 @@
|
||||
---
|
||||
tags: [MQ, 消息队列, CQRS, Event Sourcing, 架构模式]
|
||||
create time: 2026-05-24 19:52
|
||||
---
|
||||
|
||||
# MQ CQRS 与 Event Sourcing
|
||||
|
||||
## 概述
|
||||
|
||||
CQRS(Command Query Responsibility Segregation)将读写模型分离,Event Sourcing 用事件流代替状态快照。两者结合并通过 MQ 广播事件,可以构建高性能、可审计、可追溯的系统架构。本文详解这两个模式的原理、组合方式、与传统 CRUD 的对比,以及事件存储的核心设计。
|
||||
|
||||
## 正文
|
||||
|
||||
### CQRS 概述
|
||||
|
||||
传统 CRUD 模式下,同一个数据模型既用于写入也用于查询。这在简单场景下工作良好,但随着业务复杂化,读写的性能需求和模型结构会逐渐分化。
|
||||
|
||||
CQRS 的核心思想:**Command(写)和 Query(读)使用不同的模型**。
|
||||
|
||||
- **写模型**(Command Side):专注于业务规则校验和状态变更,使用领域模型(Domain Model),可以高度规范化。
|
||||
- **读模型**(Query Side):专注于查询性能,使用反规范化的视图模型(View Model),可以针对不同查询场景定制。
|
||||
|
||||
```mermaid
|
||||
graph LR
|
||||
Client["客户端"] --> CmdAPI["Command API"]
|
||||
Client --> QueryAPI["Query API"]
|
||||
|
||||
CmdAPI --> WriteDB["写模型 (Domain)"]
|
||||
WriteDB -->|"事件通过 MQ 广播"| MQ["Message Queue"]
|
||||
MQ --> ReadModel1["读模型 A (列表视图)"]
|
||||
MQ --> ReadModel2["读模型 B (统计视图)"]
|
||||
ReadModel1 --> QueryAPI
|
||||
ReadModel2 --> QueryAPI
|
||||
|
||||
style Client fill:#4A90D9,color:#fff
|
||||
style CmdAPI fill:#F5A623,color:#fff
|
||||
style QueryAPI fill:#6EC1E0,color:#fff
|
||||
style WriteDB fill:#D0021B,color:#fff
|
||||
style MQ fill:#F5A623,color:#fff
|
||||
style ReadModel1 fill:#6EC1E0,color:#fff
|
||||
style ReadModel2 fill:#6EC1E0,color:#fff
|
||||
```
|
||||
|
||||
### Event Sourcing
|
||||
|
||||
传统方式存储数据的"当前状态"——每次更新都是覆盖写入。Event Sourcing 反其道而行:**不存储当前状态,存储所有导致状态变更的事件**。
|
||||
|
||||
以电商订单为例:
|
||||
|
||||
| 传统 CRUD | Event Sourcing |
|
||||
|-----------|----------------|
|
||||
| `UPDATE orders SET status='paid'` | 追加事件 `OrderPaid{orderId, amount, time}` |
|
||||
| 只保留最终状态 | 保留完整的事件历史 |
|
||||
| 无法追溯"为什么变成这样" | 可以重放任意时间点的状态 |
|
||||
|
||||
Event Sourcing 的优势:
|
||||
- **完整审计追踪**:每个状态变更都有记录,满足金融、医疗等合规要求。
|
||||
- **时间旅行**:通过重放事件,可以重建任意历史时刻的状态。
|
||||
- **调试友好**:出 bug 时可以精确回放复现问题。
|
||||
|
||||
Event Sourcing 的挑战:
|
||||
- **事件版本兼容**:事件 Schema 演进时需要保证新旧版本兼容。
|
||||
- **查询复杂**:获取当前状态需要重放所有事件(可用 Snapshot 优化)。
|
||||
- **最终一致性**:读模型通过 MQ 异步更新,存在短暂延迟。
|
||||
|
||||
> [!question]
|
||||
> Event Sourcing 让历史可追溯,但也带来事件版本兼容的挑战。你会如何处理事件 Schema 的演进?
|
||||
|
||||
### CQRS + Event Sourcing 的组合
|
||||
|
||||
CQRS 和 Event Sourcing 天然互补:
|
||||
|
||||
1. **Command Side** 采用 Event Sourcing,将状态变更以事件形式追加到 Event Store。
|
||||
2. **Event Store** 通过 MQ 将事件广播给所有关心的消费者。
|
||||
3. **Query Side**(Read Model)消费事件,构建针对查询优化的视图。
|
||||
|
||||
这种架构下,MQ 是连接读写两侧的桥梁——写入端产生事件,MQ 负责分发,读取端消费事件并更新视图。
|
||||
|
||||
### 与传统 CRUD 的对比
|
||||
|
||||
| 维度 | 传统 CRUD | CQRS + Event Sourcing |
|
||||
|------|-----------|----------------------|
|
||||
| 数据一致性 | 强一致(同一数据库) | 最终一致(读模型异步更新) |
|
||||
| 查询性能 | 受限于写模型结构 | 读模型可针对查询优化 |
|
||||
| 审计追踪 | 需要额外设计审计表 | 天然支持,事件即审计日志 |
|
||||
| 扩展性 | 读写耦合,难以独立扩展 | 读写独立扩展 |
|
||||
| 复杂度 | 低 | 高(事件设计、版本管理、最终一致性) |
|
||||
| 适用场景 | 简单 CRUD 应用 | 复杂业务、高并发读、审计要求高 |
|
||||
|
||||
> [!question]
|
||||
> CQRS 带来最终一致性,用户下单后立刻查询可能看不到最新状态。在你的业务场景中,这种延迟可以接受吗?如何在 UX 层面缓解?
|
||||
|
||||
### 事件存储设计
|
||||
|
||||
#### 事件版本
|
||||
|
||||
事件一旦写入就不可修改(Immutable),但业务会演进。处理方式:
|
||||
- **Upcast**:读取旧版本事件时,在内存中转换为新版本。
|
||||
- **多版本共存**:事件中携带版本号,消费者根据版本号分别处理。
|
||||
|
||||
#### 快照(Snapshot)
|
||||
|
||||
当事件数量很大时,每次重放全部事件来获取当前状态会很慢。Snapshot 在特定时间点保存状态快照,重放时只需从最近的 Snapshot 开始。
|
||||
|
||||
#### 事件重放
|
||||
|
||||
重放是 Event Sourcing 的核心能力——从头(或从 Snapshot)依次应用事件,重建状态。读模型的重建、bug 修复后的数据修复、新读模型的初始化,都依赖事件重放。
|
||||
|
||||
### Go 代码:简化版 Event Sourcing
|
||||
|
||||
```go
|
||||
package main
|
||||
|
||||
import "time"
|
||||
|
||||
// Event 事件定义:每个事件代表一次状态变更
|
||||
type Event struct {
|
||||
ID string
|
||||
AggregateID string // 聚合根 ID(如订单 ID)
|
||||
Type string // 事件类型
|
||||
Data []byte // 事件数据
|
||||
Version int // 事件版本号
|
||||
Timestamp time.Time
|
||||
}
|
||||
|
||||
// EventStore 事件存储接口
|
||||
type EventStore interface {
|
||||
Save(events []Event) error
|
||||
Load(aggregateID string) ([]Event, error)
|
||||
}
|
||||
|
||||
// OrderAggregate 订单聚合根
|
||||
type OrderAggregate struct {
|
||||
ID string
|
||||
Status string
|
||||
Amount float64
|
||||
}
|
||||
|
||||
// Apply 将单个事件应用到聚合根,更新状态
|
||||
func (o *OrderAggregate) Apply(event Event) {
|
||||
switch event.Type {
|
||||
case "OrderCreated":
|
||||
o.Status = "created"
|
||||
case "OrderPaid":
|
||||
o.Status = "paid"
|
||||
case "OrderShipped":
|
||||
o.Status = "shipped"
|
||||
}
|
||||
}
|
||||
|
||||
// ReplayEvents 从事件流重建聚合根状态
|
||||
func ReplayEvents(aggregateID string, store EventStore) (*OrderAggregate, error) {
|
||||
events, err := store.Load(aggregateID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
agg := &OrderAggregate{ID: aggregateID}
|
||||
for _, event := range events {
|
||||
agg.Apply(event) // 依次应用每个事件
|
||||
}
|
||||
return agg, nil
|
||||
}
|
||||
```
|
||||
|
||||
核心逻辑:`ReplayEvents` 从 Event Store 加载指定聚合根的所有事件,依次调用 `Apply` 重建当前状态。这就是 Event Sourcing 的精髓——状态是事件的函数,`State = f(events)`。
|
||||
|
||||
## 关联笔记
|
||||
|
||||
- [[02-消息模型/3-MQ-消息模型|MQ 消息模型]]
|
||||
- [[09-流处理与事件驱动/34-事件驱动架构-EDA|事件驱动架构 EDA]]
|
||||
- [[12-架构与实战/47-MQ-与微服务|MQ 与微服务]]
|
||||
- [[06-高级特性/20-MQ-Schema-管理与演进|MQ Schema 管理与演进]]
|
||||
- [[12-架构与实战/46-MQ-分布式事务实践|MQ 分布式事务实践]]
|
||||
@@ -0,0 +1,168 @@
|
||||
---
|
||||
tags: [MQ, 消息队列, 设计模式, 性能优化]
|
||||
create time: 2026-05-24 19:52
|
||||
---
|
||||
|
||||
# MQ Claim Check 与消息瘦身
|
||||
|
||||
## 概述
|
||||
|
||||
当消息体包含图片、文件、富文本等大体积数据时,直接放入 MQ 会严重影响 Broker 的吞吐和存储性能。Claim Check(存根/提货单)模式的核心思路是:**消息体外置到对象存储,MQ 中只传递一个轻量的引用(Key)**。消费端按需根据 Key 拉取完整数据,从而实现"消息瘦身"。
|
||||
|
||||
## 正文
|
||||
|
||||
### 问题:大消息的代价
|
||||
|
||||
MQ 的设计哲学是"快速转发小消息"。当消息体积增大时,会引发一系列连锁问题:
|
||||
|
||||
- **网络带宽**:Broker 需要在生产者和消费者之间转发完整消息,大消息占用大量带宽。
|
||||
- **存储压力**:Kafka 的 Partition Log、RocketMQ 的 CommitLog 都是顺序写入,大消息导致磁盘 IO 放大。
|
||||
- **内存占用**:Broker 和 Consumer 的缓冲区需要加载完整消息,GC 压力增大。
|
||||
- **延迟增加**:大消息的序列化、反序列化和网络传输时间更长,端到端延迟上升。
|
||||
|
||||
一个典型的例子:电商系统中,用户上传的商品图片(几 MB)如果直接塞进消息体,MQ 的吞吐量可能下降一个数量级。
|
||||
|
||||
> [!question]
|
||||
> 如果你的系统每天产生 100 万条消息,每条消息包含一张 2MB 的图片,MQ 需要额外存储 2TB 数据。这对 Broker 集群的磁盘和网络意味着什么?
|
||||
|
||||
### Claim Check 模式
|
||||
|
||||
Claim Check 模式的灵感来自衣帽间:你把大衣(消息体)存起来,只拿一张小票(Key)。消费时凭小票取回大衣。
|
||||
|
||||
核心流程:
|
||||
|
||||
1. **Producer 端**:将消息体上传到对象存储(S3/MinIO/OSS),获取 Storage Key。
|
||||
2. **MQ 传递**:Producer 只将 Storage Key 和必要的元数据发送到 MQ。
|
||||
3. **Consumer 端**:消费者收到消息后,根据 Storage Key 从对象存储拉取完整数据。
|
||||
|
||||
```mermaid
|
||||
graph LR
|
||||
P["Producer"] -->|"1. 上传大消息体"| OS["Object Storage (S3/MinIO)"]
|
||||
OS -->|"2. 返回 Storage Key"| P
|
||||
P -->|"3. 发送轻量消息 (Key + 元数据)"| MQ["Message Queue"]
|
||||
MQ -->|"4. 转发消息"| C["Consumer"]
|
||||
C -->|"5. 根据 Key 拉取完整数据"| OS
|
||||
|
||||
style P fill:#4A90D9,color:#fff
|
||||
style OS fill:#F5A623,color:#fff
|
||||
style MQ fill:#6EC1E0,color:#fff
|
||||
style C fill:#D0021B,color:#fff
|
||||
```
|
||||
|
||||
### 权衡分析
|
||||
|
||||
Claim Check 不是免费午餐,它引入了一个权衡:
|
||||
|
||||
| 维度 | 直接发送大消息 | Claim Check 模式 |
|
||||
|------|--------------|-----------------|
|
||||
| MQ 性能 | 差(大消息拖慢 Broker) | 好(消息体极小) |
|
||||
| 网络开销 | 单次大传输 | 多次小传输(MQ + 对象存储) |
|
||||
| 消费延迟 | 低(数据已在消息中) | 略高(需要额外一次 IO) |
|
||||
| 运维复杂度 | 低 | 中(需要管理对象存储) |
|
||||
| 存储成本 | MQ 磁盘成本高 | 对象存储成本低(通常更便宜) |
|
||||
|
||||
大多数场景下,MQ 性能的提升远大于额外一次对象存储 IO 的开销。对象存储(如 S3)本身就是为高吞吐、低延迟的读取设计的。
|
||||
|
||||
### 实现方式
|
||||
|
||||
**消息存 S3/MinIO/OSS**:Producer 先将消息体 PUT 到对象存储,拿到 Key 后封装成轻量消息发送到 MQ。
|
||||
|
||||
**消费端按需拉取**:Consumer 收到消息后,根据 Key 从对象存储 GET 完整数据。如果某些消费者只需要元数据而不需要完整内容(如路由、过滤),就可以跳过拉取步骤。
|
||||
|
||||
**生命周期管理**:对象存储中的数据需要设置过期策略(TTL),避免无限增长。可以与消息的消费确认(ACK)联动——消息被 ACK 后,对象存储中的数据保留 N 天后自动清理。
|
||||
|
||||
### 变体方案
|
||||
|
||||
#### 消息压缩
|
||||
|
||||
如果消息体不算特别大(几百 KB 到几 MB),可以先压缩再发送。Snappy、LZ4、Zstd 等压缩算法可以在几乎不增加延迟的前提下,将消息体积减少 50%-80%。
|
||||
|
||||
#### 消息分片
|
||||
|
||||
超大消息拆成多个小消息发送,消费端组装。这种方式实现复杂,需要处理分片丢失、乱序等问题,一般只在极端场景下使用。
|
||||
|
||||
> [!question]
|
||||
> Claim Check 引入了对象存储这个外部依赖,如果对象存储不可用怎么办?
|
||||
|
||||
### Go 代码:Claim Check 模式实现
|
||||
|
||||
```go
|
||||
package main
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
|
||||
"github.com/minio/minio-go/v7"
|
||||
"github.com/segmentio/kafka-go"
|
||||
)
|
||||
|
||||
// ClaimCheckMessage 轻量消息:只包含引用,不包含实际数据
|
||||
type ClaimCheckMessage struct {
|
||||
StorageKey string `json:"storage_key"` // 对象存储中的 Key
|
||||
Bucket string `json:"bucket"` // 存储桶名
|
||||
Metadata map[string]string `json:"metadata"` // 业务元数据
|
||||
}
|
||||
|
||||
// ========== Producer 端 ==========
|
||||
|
||||
// PublishWithClaimCheck 大消息外置存储,MQ 只传引用
|
||||
func PublishWithClaimCheck(ctx context.Context, minioClient *minio.Client, writer *kafka.Writer, bucket string, data []byte, metadata map[string]string) error {
|
||||
// 1. 生成唯一的 Storage Key
|
||||
key := fmt.Sprintf("msg/%s", generateUUID())
|
||||
|
||||
// 2. 将消息体上传到对象存储
|
||||
_, err := minioClient.PutObject(ctx, bucket, key, bytes.NewReader(data), int64(data.Length()), minio.PutObjectOptions{})
|
||||
if err != nil {
|
||||
return fmt.Errorf("upload to object storage: %w", err)
|
||||
}
|
||||
|
||||
// 3. 构造轻量消息,只包含引用
|
||||
msg := ClaimCheckMessage{
|
||||
StorageKey: key,
|
||||
Bucket: bucket,
|
||||
Metadata: metadata,
|
||||
}
|
||||
payload, _ := json.Marshal(msg)
|
||||
|
||||
// 4. 将轻量消息发送到 MQ
|
||||
return writer.WriteMessages(ctx, kafka.Message{Value: payload})
|
||||
}
|
||||
|
||||
// ========== Consumer 端 ==========
|
||||
|
||||
// ConsumeWithClaimCheck 消费消息时按需拉取完整数据
|
||||
func ConsumeWithClaimCheck(ctx context.Context, minioClient *minio.Client, reader *kafka.Reader) {
|
||||
for {
|
||||
msg, _ := reader.ReadMessage(ctx)
|
||||
|
||||
// 1. 解析轻量消息,获取引用
|
||||
var claim ClaimCheckMessage
|
||||
json.Unmarshal(msg.Value, &claim)
|
||||
|
||||
// 2. 根据 Key 从对象存储拉取完整数据
|
||||
obj, err := minioClient.GetObject(ctx, claim.Bucket, claim.StorageKey, minio.GetObjectOptions{})
|
||||
if err != nil {
|
||||
continue // 拉取失败,可以重试或进死信队列
|
||||
}
|
||||
|
||||
// 3. 处理完整数据
|
||||
buf := new(bytes.Buffer)
|
||||
buf.ReadFrom(obj)
|
||||
processFullData(buf.Bytes(), claim.Metadata)
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
Producer 端的逻辑分两步:先把大消息体上传到 MinIO(步骤 1-2),再把包含 Storage Key 的轻量消息发送到 Kafka(步骤 3-4)。Consumer 端反向操作:先从 Kafka 读取轻量消息,再根据 Key 从 MinIO 拉取完整数据。
|
||||
|
||||
注意 Consumer 端的错误处理——如果对象存储暂时不可用,消息可以重试或进入死信队列,而不是直接丢弃。
|
||||
|
||||
## 关联笔记
|
||||
|
||||
- [[06-高级特性/21-MQ-消息压缩与批处理|MQ 消息压缩与批处理]]
|
||||
- [[05-可靠性保障/12-MQ-消息确认与持久化|MQ 消息确认与持久化]]
|
||||
- [[05-可靠性保障/16-MQ-死信队列与消息回溯|MQ 死信队列与消息回溯]]
|
||||
- [[08-消息设计模式/31-MQ-背压与流控|MQ 背压与流控]]
|
||||
@@ -0,0 +1,158 @@
|
||||
---
|
||||
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]]
|
||||
@@ -0,0 +1,140 @@
|
||||
---
|
||||
tags:
|
||||
- MQ
|
||||
create time: 2026-05-24 19:52
|
||||
---
|
||||
|
||||
# MQ 请求-回复模式
|
||||
|
||||
## 概述
|
||||
|
||||
请求-回复模式(Request-Reply)是通过 MQ 实现同步 RPC 语义的方式:Producer 发送请求消息,Consumer 处理后将响应发送到 Reply-To 队列,Producer 通过 CorrelationID 匹配响应。这种模式在异步通道上模拟了同步调用,适用于跨网络不可靠、需要消息持久化或遗留系统集成的场景。
|
||||
|
||||
## 正文
|
||||
|
||||
### 基本原理
|
||||
|
||||
传统 RPC 是一问一答:客户端发请求,服务端返回响应。请求-回复模式把这个过程搬到了 MQ 上:
|
||||
|
||||
1. Producer 发送请求消息,携带 `ReplyTo`(回复队列名)和 `CorrelationID`(请求标识)
|
||||
2. Consumer 从请求队列消费消息,处理业务逻辑
|
||||
3. Consumer 将响应发送到 `ReplyTo` 队列,携带相同的 `CorrelationID`
|
||||
4. Producer 从 ReplyTo 队列消费响应,用 `CorrelationID` 匹配请求
|
||||
|
||||
```mermaid
|
||||
sequenceDiagram
|
||||
participant P as Producer
|
||||
participant BQ as "Broker Request Queue"
|
||||
participant C as Consumer
|
||||
participant RQ as "Broker Reply Queue"
|
||||
|
||||
P->>BQ: "发送请求 (CorrelationID=abc123, ReplyTo=reply-queue)"
|
||||
BQ->>C: "投递请求"
|
||||
C->>C: "处理业务逻辑"
|
||||
C->>RQ: "发送响应 (CorrelationID=abc123)"
|
||||
RQ->>P: "投递响应"
|
||||
P->>P: "匹配 CorrelationID, 返回结果"
|
||||
```
|
||||
|
||||
关键在于 `CorrelationID`——它是请求和响应之间的"凭证"。没有它,当多个请求并发时,Producer 无法知道哪个响应对应哪个请求。
|
||||
|
||||
### 与直接 RPC 的对比
|
||||
|
||||
| 维度 | gRPC / HTTP | MQ 请求-回复 |
|
||||
|------|------------|-------------|
|
||||
| 耦合度 | 直连,需要知道服务地址 | 通过 Broker 中转,天然解耦 |
|
||||
| 削峰能力 | 无,服务端过载直接拒绝 | Broker 缓冲,削峰填谷 |
|
||||
| 延迟 | 低(一次网络跳转) | 高(至少两次 Broker 中转) |
|
||||
| 持久化 | 无 | 消息可持久化,故障可恢复 |
|
||||
| 复杂度 | 低 | 高(CorrelationID 匹配、超时管理) |
|
||||
|
||||
> [!question]
|
||||
> 请求-回复模式看起来像用 MQ 模拟 RPC,什么情况下这比直接用 gRPC 更好?
|
||||
|
||||
答案藏在"代价"里:你付出的是延迟和复杂度,换来的是解耦、削峰和可靠性。当这些特性比低延迟更重要时,请求-回复模式就是合理的选择。
|
||||
|
||||
### 适用场景
|
||||
|
||||
**跨网络不可靠的远程调用**:网络不稳定时,MQ 的持久化和重试机制比直连更可靠。消息不会因为网络抖动而丢失。
|
||||
|
||||
**需要消息持久化的 RPC**:某些金融场景要求每一次请求都有据可查。MQ 天然支持消息持久化,可以作为审计日志。
|
||||
|
||||
**遗留系统集成**:老系统可能只暴露 MQ 接口(比如 IBM MQ),新系统通过请求-回复模式与其交互,避免大规模改造。
|
||||
|
||||
**异步长任务**:某些请求处理耗时很长(分钟级),用直连 RPC 会超时。通过 MQ,Producer 发完请求就可以去做别的,异步接收结果。
|
||||
|
||||
### 注意事项
|
||||
|
||||
**超时处理**:Producer 不能无限等待回复。需要设置超时,超时后清理本地的 pending 请求映射。超时的响应如果后续到达,应该被丢弃或记录日志。
|
||||
|
||||
**Reply-To 队列的生命周期管理**:
|
||||
|
||||
- **持久队列**:一个 Producer 固定使用一个回复队列,适合长期运行的服务
|
||||
- **临时队列**:每次请求创建一个临时队列,用完即删。适合短生命周期的客户端,但创建和销毁队列有额外开销
|
||||
|
||||
**消息确认**:响应消息也需要 ACK 机制,否则 Broker 会反复投递,导致重复处理。
|
||||
|
||||
### Go 实现
|
||||
|
||||
下面是通过 RabbitMQ 实现请求-回复模式的 Producer 和 Consumer 示例。
|
||||
|
||||
**Producer 端**:发送请求并等待回复
|
||||
|
||||
```go
|
||||
func rpcClient(ch *amqp.Channel, request string) (string, error) {
|
||||
// 声明一个独占的临时回复队列
|
||||
replyQ, err := ch.QueueDeclare("", false, false, true, false, nil)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
|
||||
corrID := uuid.New().String()
|
||||
|
||||
// 发送请求,附带 ReplyTo 和 CorrelationID
|
||||
ch.PublishWithContext(ctx, "", "rpc_queue", false, false,
|
||||
amqp.Publishing{
|
||||
ContentType: "text/plain",
|
||||
CorrelationId: corrID,
|
||||
ReplyTo: replyQ.Name,
|
||||
Body: []byte(request),
|
||||
})
|
||||
|
||||
// 等待匹配的回复
|
||||
msgs, _ := ch.Consume(replyQ.Name, "", true, false, false, false, nil)
|
||||
for msg := range msgs {
|
||||
if msg.CorrelationId == corrID {
|
||||
return string(msg.Body), nil
|
||||
}
|
||||
}
|
||||
return "", fmt.Errorf("timeout")
|
||||
}
|
||||
```
|
||||
|
||||
**Consumer 端**:处理请求并发送回复
|
||||
|
||||
```go
|
||||
func rpcServer(ch *amqp.Channel) {
|
||||
msgs, _ := ch.Consume("rpc_queue", "", false, false, false, false, nil)
|
||||
for msg := range msgs {
|
||||
// 处理请求
|
||||
result := processRequest(msg.Body)
|
||||
|
||||
// 将回复发送到 ReplyTo 队列,携带相同的 CorrelationID
|
||||
ch.PublishWithContext(ctx, "", msg.ReplyTo, false, false,
|
||||
amqp.Publishing{
|
||||
ContentType: "text/plain",
|
||||
CorrelationId: msg.CorrelationId,
|
||||
Body: []byte(result),
|
||||
})
|
||||
msg.Ack(false)
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
Producer 用一个临时队列接收回复,通过 `CorrelationID` 精确匹配。Consumer 只需要把处理结果发回 `ReplyTo` 队列即可——这就是请求-回复模式的全部核心逻辑。
|
||||
|
||||
## 关联笔记
|
||||
|
||||
- [[31-MQ-背压与流控]]
|
||||
- [[33-MQ-与流处理]]
|
||||
- [[34-事件驱动架构-EDA]]
|
||||
Reference in New Issue
Block a user