8.2 KiB
tags, create time
| tags | create time | |
|---|---|---|
|
2026-05-24 19:52 |
MQ 存储引擎设计
概述
存储引擎是 MQ 性能的基石。消息从生产者发出到被消费者拉取,中间经历的"写入—持久化—读取"全链路,每一步都与存储设计息息相关。本文从磁盘 I/O 模型出发,依次讲解顺序写、零拷贝、Page Cache、批量刷盘、Append-Only 日志等核心原理,帮你建立对高性能消息存储的系统认知。
正文
为什么存储引擎是 MQ 的性能基石
MQ 的本质是一个"写多读多"的中转站:生产者不断写入消息,消费者不断拉取消息。如果底层存储扛不住写入吞吐,上层的协议优化、集群扩展全是空谈。可以说,存储引擎的设计直接决定了 MQ 的性能上限。
[!question] 思考 如果让你从零设计一个 MQ,你会选择把消息存在哪里?关系数据库?Redis?还是直接写文件?
磁盘顺序写 vs 随机写
传统认知里,磁盘(尤其是机械硬盘 HDD)是性能瓶颈。但这里有个关键区分:顺序写和随机写的性能差距是数量级的。
- 随机写:磁头需要反复寻道,HDD 的 IOPS 通常只有 100-200 次/秒。
- 顺序写:磁头几乎不需要移动,HDD 顺序写吞吐可达 600MB/s 以上,甚至媲美 SATA SSD 的顺序写性能。
这就是为什么 Kafka、RocketMQ 等高性能 MQ 都选择了顺序写盘的策略——即使是廉价的 HDD 也能获得极高的写入吞吐。
graph LR
A["Producer"] -->|"顺序追加写入"| B["磁盘文件"]
B -->|"Consumer 按偏移量顺序读取"| C["Consumer"]
style A fill:#4A90D9,color:#fff
style B fill:#F5A623,color:#fff
style C fill:#6EC1E0,color:#fff
零拷贝(Zero-Copy)技术
传统网络传输一条消息,数据需要经历 4 次拷贝和 4 次上下文切换:
- 磁盘 → 内核缓冲区(DMA 拷贝)
- 内核缓冲区 → 用户空间缓冲区(CPU 拷贝)
- 用户空间缓冲区 → Socket 缓冲区(CPU 拷贝)
- Socket 缓冲区 → 网卡(DMA 拷贝)
sendfile 系统调用可以让数据直接从内核缓冲区传到网卡,跳过用户空间的两次拷贝,只需要 2 次 DMA 拷贝 + 1 次上下文切换。Kafka 正是利用 sendfile 实现了消费者拉取数据时的零拷贝,这让消费者读取消息几乎不消耗 CPU 资源。
graph LR
subgraph "传统拷贝"
D1["磁盘"] -->|"DMA"| K1["内核缓冲区"]
K1 -->|"CPU"| U1["用户空间"]
U1 -->|"CPU"| S1["Socket 缓冲区"]
S1 -->|"DMA"| N1["网卡"]
end
subgraph "零拷贝 sendfile"
D2["磁盘"] -->|"DMA"| K2["内核缓冲区"]
K2 -->|"DMA"| N2["网卡"]
end
style D1 fill:#F5A623,color:#fff
style D2 fill:#F5A623,color:#fff
Page Cache 与内存映射(mmap)
操作系统内核会自动将最近访问的磁盘数据缓存在 Page Cache 中。MQ 的写入和读取天然契合 Page Cache 的工作模式:
- 写入时:消息先写入 Page Cache,由操作系统异步刷盘,应用层返回极快。
- 读取时:如果消费者读取的是刚写入的数据(大多数场景),直接从 Page Cache 返回,无需访问磁盘。
mmap(内存映射) 则更进一步:将磁盘文件直接映射到进程的虚拟地址空间,读写文件就像读写内存一样,省去了 read/write 系统调用的开销。RocketMQ 的 MappedFile 就是基于 mmap 实现的。
[!question] 思考 Page Cache 加速了读写,但如果 Broker 进程崩溃,Page Cache 中还没刷盘的数据会丢失吗?这对"异步刷盘"策略意味着什么?
批量写入与刷盘策略
MQ 通常不会每条消息都触发一次磁盘 I/O,而是将多条消息在内存中攒批后再统一写入,这就是批量写入。刷盘策略则决定了数据何时真正落盘:
| 策略 | 机制 | 可靠性 | 性能 |
|---|---|---|---|
| 同步刷盘 | 每批消息写入后,调用 fsync 等待数据落盘才返回 |
高:数据不丢 | 较低:受磁盘 I/O 限制 |
| 异步刷盘 | 消息写入 Page Cache 即返回,后台线程定时刷盘 | 较低:宕机可能丢最后几条 | 高:写入延迟极低 |
[!question] 思考 异步刷盘宕机可能丢数据,为什么很多 MQ 还是默认异步刷盘?
答案在于实际场景的权衡:大多数业务场景可以容忍极少量消息丢失(配合重试机制),但无法容忍高延迟。而且异步刷盘配合主从同步(消息复制到其他节点后再返回),可以在性能和可靠性之间找到平衡。只有金融级场景才需要同步刷盘。
日志追加(Append-Only)设计
为什么 MQ 不用数据库(如 MySQL)存储消息?核心原因有三:
- 写入模式不匹配:数据库面向"随机读写"优化,B+ 树索引在高并发写入下会成为瓶颈;而 MQ 是纯粹的"顺序追加 + 顺序读取",Append-Only 日志天然适配。
- 无索引开销:消息不需要按内容检索,只需要按偏移量(offset)顺序读取,省去了索引维护的代价。
- 删除成本低:过期日志直接截断文件头,不需要像数据库那样做 VACUUM 或标记删除。
存储写入流程总览
graph TD
P["Producer 发送消息"] --> B["Broker 接收"]
B --> M["消息序列化写入内存缓冲区"]
M --> BQ{"是否达到批量阈值?"}
BQ -->|"否"| M
BQ -->|"是"| FS["追加写入磁盘日志文件"]
FS --> F{"刷盘策略?"}
F -->|"同步"| FSYNC["fsync 确保落盘"]
F -->|"异步"| PC["写入 Page Cache 后返回"]
PC --> BG["后台线程定时 fsync"]
FSYNC --> ACK["返回 ACK 给 Producer"]
BG --> ACK
style P fill:#4A90D9,color:#fff
style B fill:#6EC1E0,color:#fff
style ACK fill:#2ECC71,color:#fff
Go 伪代码:一个简单的 Append-Only Log 存储
下面的代码展示了一个最简的 Append-Only Log 的核心逻辑——顺序追加写入和按偏移量读取。
// AppendLog 是一个简化版的顺序写日志存储
type AppendLog struct {
mu sync.Mutex
file *os.File
offset int64 // 当前写入位置
}
// NewAppendLog 打开或创建日志文件
func NewAppendLog(path string) (*AppendLog, error) {
f, err := os.OpenFile(path, os.O_CREATE|os.O_APPEND|os.O_WRONLY, 0644)
if err != nil {
return nil, err
}
// 获取当前文件大小作为起始偏移量
stat, _ := f.Stat()
return &AppendLog{file: f, offset: stat.Size()}, nil
}
// Append 顺序追加一条消息,返回该消息在文件中的偏移量
func (a *AppendLog) Append(data []byte) (int64, error) {
a.mu.Lock()
defer a.mu.Unlock()
// 写入长度前缀 + 数据(Length-Prefix 编码,方便读取时知道边界)
buf := make([]byte, 4+len(data))
binary.BigEndian.PutUint32(buf[:4], uint32(len(data)))
copy(buf[4:], data)
n, err := a.file.Write(buf) // 顺序追加,OS 自动利用 Page Cache
if err != nil {
return 0, err
}
pos := a.offset
a.offset += int64(n)
return pos, nil // 返回偏移量,消费者后续按此读取
}
// ReadAt 从指定偏移量读取消息(消费者按 offset 拉取)
func (a *AppendLog) ReadAt(offset int64) ([]byte, error) {
// 读取 4 字节长度头
header := make([]byte, 4)
if _, err := a.file.ReadAt(header, offset); err != nil {
return nil, err
}
length := binary.BigEndian.Uint32(header)
// 读取消息体
data := make([]byte, length)
if _, err := a.file.ReadAt(data, offset+4); err != nil {
return nil, err
}
return data, nil
}
核心要点:写入只做追加(O_APPEND),不修改已有数据;读取通过偏移量直接定位,无需索引。这就是 Kafka、RocketMQ 存储引擎的最简原型。