169 lines
6.9 KiB
Markdown
169 lines
6.9 KiB
Markdown
|
|
---
|
|||
|
|
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 背压与流控]]
|