237 lines
10 KiB
Markdown
237 lines
10 KiB
Markdown
---
|
||
tags: [MQ]
|
||
create time: 2026-05-24 19:52
|
||
---
|
||
|
||
# Kafka 存储设计
|
||
|
||
## 概述
|
||
|
||
Kafka 的高吞吐能力不仅来自"顺序写盘"和"零拷贝"等通用优化,更源于其精心设计的存储结构:Partition Log 的分段存储、稀疏索引、日志压缩,以及内部定时任务调度所用的时间轮。本文深入 Kafka 存储层的每一个关键设计,帮你理解"为什么 Kafka 这么快"。
|
||
|
||
## 正文
|
||
|
||
### Partition Log 的整体结构
|
||
|
||
Kafka 的存储层次可以概括为:**Topic → Partition → Segment**。每个 Partition 是一个有序的、不可变的消息序列,底层由多个 Segment 文件组成。
|
||
|
||
一个 Partition 目录下的典型文件结构:
|
||
|
||
```
|
||
partition-0/
|
||
├── 00000000000000000000.log # 第一个 Segment 的消息数据
|
||
├── 00000000000000000000.index # 稀疏偏移索引
|
||
├── 00000000000000000000.timeindex # 时间戳索引
|
||
├── 00000000000000523840.log # 第二个 Segment(起始 offset = 523840)
|
||
├── 00000000000000523840.index
|
||
├── 00000000000000523840.timeindex
|
||
└── ...
|
||
```
|
||
|
||
文件名是该 Segment 的 **base offset**(起始偏移量)。当 Segment 达到配置的大小或时间阈值时,Kafka 会"滚动"出一个新的 Segment。
|
||
|
||
```mermaid
|
||
graph TD
|
||
T["Topic: orders"] --> P0["Partition 0"]
|
||
T --> P1["Partition 1"]
|
||
T --> P2["Partition 2"]
|
||
|
||
P0 --> S1["Segment 0\nbase offset = 0"]
|
||
P0 --> S2["Segment 1\nbase offset = 523840"]
|
||
P0 --> S3["Segment 2\nbase offset = 1048576"]
|
||
|
||
S1 --> L1["00000.log"]
|
||
S1 --> I1["00000.index"]
|
||
S1 --> TI1["00000.timeindex"]
|
||
|
||
S2 --> L2["523840.log"]
|
||
S2 --> I2["523840.index"]
|
||
S2 --> TI2["523840.timeindex"]
|
||
|
||
style T fill:#4A90D9,color:#fff
|
||
style P0 fill:#6EC1E0,color:#fff
|
||
style P1 fill:#6EC1E0,color:#fff
|
||
style P2 fill:#6EC1E0,color:#fff
|
||
```
|
||
|
||
> [!question] 思考
|
||
> 为什么不把整个 Partition 写成一个巨大的日志文件,而要拆分成多个 Segment?
|
||
|
||
### 分段存储策略
|
||
|
||
Segment 滚动的触发条件有两个(满足任一即滚动):
|
||
|
||
- **按大小**:默认 1GB(`log.segment.bytes`)
|
||
- **按时间**:默认 7 天(`log.roll.ms` / `log.roll.hours`)
|
||
|
||
分段的好处是多方面的:
|
||
|
||
1. **快速定位**:消费者只需根据 offset 找到对应 Segment,再在小文件内查找,而不是在一个 TB 级大文件中扫描。
|
||
2. **高效清理**:过期数据直接删除整个 Segment 文件,无需逐条标记删除或做 compaction。
|
||
3. **并行 I/O**:不同 Segment 可以被不同的消费者线程并发读取。
|
||
|
||
### 稀疏索引(Sparse Index)
|
||
|
||
Kafka 不是每条消息都建索引——那会带来巨大的存储和维护开销。它采用**稀疏索引**:每隔一定字节(默认 4KB,由 `log.index.interval.bytes` 控制)才记录一条索引项,格式为 `offset → 文件物理位置`。
|
||
|
||
当消费者需要查找某个 offset 的消息时,流程如下:
|
||
|
||
1. 在 `.index` 文件中用**二分查找**找到不大于目标 offset 的最近索引项。
|
||
2. 从该索引项指向的物理位置开始,在 `.log` 文件中**顺序扫描**直到找到目标 offset。
|
||
|
||
因为索引本身是有序的,二分查找效率为 O(log N);而顺序扫描的范围通常只有几条消息(4KB 内),开销极小。这种设计用极少的索引空间换来了接近 O(1) 的查找性能。
|
||
|
||
### 日志压缩(Log Compaction)
|
||
|
||
普通的消息保留策略是按时间或大小删除过期 Segment。但 Kafka 还提供了另一种策略——**Log Compaction**:保留每个 Key 的最后一条消息,删除之前的旧版本。
|
||
|
||
工作原理:
|
||
|
||
- Log Compaction 由 Cleaner 线程在后台执行。
|
||
- Cleaner 为每个 Key 维护一个"最新 offset"映射,只保留最新消息。
|
||
- 清理后的 Segment 中,每个 Key 只有一条记录。
|
||
- Consumer 从头消费时,仍然能拿到每个 Key 的最新状态。
|
||
|
||
**典型场景**:变更数据捕获(CDC)、状态快照同步。比如数据库的一张用户表,每次更新都发一条消息到 Kafka(Key = userId),下游系统通过 Log Compaction 拿到的就是每个用户的最新状态。
|
||
|
||
> [!question] 思考
|
||
> Log Compaction 保证的是"每个 Key 至少保留最新一条",但如果有两个 Key 相同的消息几乎同时到达,Compaction 会保留哪条?
|
||
|
||
### 时间轮(Timing Wheel)
|
||
|
||
Kafka 内部有大量的定时任务:延迟消息、会话过期、日志清理等。如果每个定时任务都开一个 goroutine(或 Java 的 Timer),任务数多了之后,调度开销会非常大。
|
||
|
||
Kafka 使用**层级时间轮(Hierarchical Timing Wheel)**来高效管理定时任务:
|
||
|
||
- 时间轮是一个环形数组,每个槽位(slot)代表一个时间区间。
|
||
- 新任务根据到期时间插入对应槽位。
|
||
- 时钟每推进一个 tick,处理当前槽位的所有任务。
|
||
- 当低层时间轮溢出时,任务会被"滴答"到上层时间轮,等待合适时机再降下来。
|
||
|
||
时间轮的插入和删除都是 O(1) 操作,远优于优先队列的 O(log N)。
|
||
|
||
```mermaid
|
||
graph TD
|
||
TW["时间轮"] --> L1["第一层: 1ms 精度\n20 个槽位"]
|
||
TW --> L2["第二层: 20ms 精度\n20 个槽位"]
|
||
TW --> L3["第三层: 400ms 精度\n20 个槽位"]
|
||
|
||
L1 --> SLOT1["slot 0: 任务 A"]
|
||
L1 --> SLOT2["slot 5: 任务 B"]
|
||
L2 --> SLOT3["slot 3: 任务 C\n溢出后降级到第一层"]
|
||
|
||
style TW fill:#4A90D9,color:#fff
|
||
style L1 fill:#6EC1E0,color:#fff
|
||
style L2 fill:#F5A623,color:#fff
|
||
style L3 fill:#D0021B,color:#fff
|
||
```
|
||
|
||
### 页缓存与 Sendfile:读写路径中的内核级优化
|
||
|
||
Kafka 的读写路径深度依赖 Linux 内核的两个能力:
|
||
|
||
**写入路径**:Producer 发来的消息写入 Page Cache(内存),由内核异步刷盘。应用层不调用 `fsync`,写入延迟极低。多个 Partition 的写入共享 Page Cache,操作系统会自动管理缓存淘汰。
|
||
|
||
**读取路径**:Consumer 拉取消息时,Kafka 通过 `sendfile` 系统调用直接将 Page Cache 中的数据传输到网卡,**数据完全不经过用户空间**。这意味着:
|
||
|
||
- 没有用户态/内核态的上下文切换
|
||
- 没有内存拷贝
|
||
- 大量 Consumer 并发拉取时,CPU 开销几乎不增长
|
||
|
||
这就是为什么 Kafka 在普通硬件上就能达到百万级 TPS 的核心秘密——它把操作系统的缓存和网络能力用到了极致。
|
||
|
||
> [!question] 思考
|
||
> 如果 Kafka Broker 的内存足够大,Page Cache 能缓存大量数据。但如果消费者需要回溯到很早的消息(不在 Page Cache 中),会发生什么?性能会下降多少?
|
||
|
||
### Go 代码:简化版 Segment 文件的读写
|
||
|
||
下面展示一个简化版的 Segment 文件实现,包含日志文件和稀疏索引的写入与查找。
|
||
|
||
```go
|
||
type Segment struct {
|
||
baseOffset int64
|
||
logFile *os.File
|
||
indexFile *os.File
|
||
index []indexEntry // 内存中的稀疏索引
|
||
}
|
||
|
||
type indexEntry struct {
|
||
offset int64 // 消息的逻辑偏移量
|
||
position int32 // 消息在 .log 文件中的物理位置
|
||
}
|
||
|
||
// Write 追加一条消息到 Segment
|
||
func (s *Segment) Write(offset int64, data []byte) error {
|
||
// 记录当前写入位置作为物理偏移
|
||
stat, _ := s.logFile.Stat()
|
||
pos := stat.Size()
|
||
|
||
// Length-Prefix 编码写入日志文件
|
||
header := make([]byte, 4+len(data))
|
||
binary.BigEndian.PutUint32(header[:4], uint32(len(data)))
|
||
copy(header[4:], data)
|
||
if _, err := s.logFile.Write(header); err != nil {
|
||
return err
|
||
}
|
||
|
||
// 每隔一定消息数写入一条索引(稀疏索引,不是每条都写)
|
||
if len(s.index) == 0 || offset-s.index[len(s.index)-1].offset >= 4 {
|
||
entry := indexEntry{offset: offset, position: int32(pos)}
|
||
s.index = append(s.index, entry)
|
||
// 同步写入 .index 文件(简化:每条 12 字节 = 8 offset + 4 position)
|
||
buf := make([]byte, 12)
|
||
binary.BigEndian.PutUint64(buf[:8], uint64(offset))
|
||
binary.BigEndian.PutUint32(buf[8:], uint32(pos))
|
||
s.indexFile.Write(buf)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// Read 根据 offset 查找并读取消息
|
||
func (s *Segment) Read(offset int64) ([]byte, error) {
|
||
// 第一步:二分查找稀疏索引,找到最近的索引项
|
||
idx := sort.Search(len(s.index), func(i int) bool {
|
||
return s.index[i].offset > offset
|
||
}) - 1
|
||
|
||
var startPos int32
|
||
if idx >= 0 {
|
||
startPos = s.index[idx].position // 从索引指向的位置开始
|
||
}
|
||
|
||
// 第二步:从 startPos 开始顺序扫描 .log 文件
|
||
pos := int64(startPos)
|
||
for {
|
||
header := make([]byte, 4)
|
||
if _, err := s.logFile.ReadAt(header, pos); err != nil {
|
||
return nil, err
|
||
}
|
||
length := binary.BigEndian.Uint32(header)
|
||
data := make([]byte, length)
|
||
s.logFile.ReadAt(data, pos+4)
|
||
|
||
// 简化:假设 offset 是连续递增的
|
||
currentOffset := s.baseOffset + (pos / int64(4+length))
|
||
if currentOffset == offset {
|
||
return data, nil
|
||
}
|
||
pos += int64(4 + length)
|
||
}
|
||
}
|
||
```
|
||
|
||
**核心要点**:写入时不是每条消息都建索引(稀疏索引),读取时先二分查找索引再顺序扫描——用极小的索引空间和极少的扫描范围实现高效定位。
|
||
|
||
> [!question] 思考
|
||
> Kafka 的 Segment 为什么要同时维护 .log 和 .index 两个文件?只用 .log 行不行?
|
||
|
||
答案是:**理论上行,但实践中代价太大**。如果只用 .log,消费者查找某个 offset 的消息时,必须从 Segment 头部开始顺序扫描,时间复杂度为 O(N)。在消息量巨大的场景下(单 Segment 可达 1GB、数百万条消息),这会严重拖慢消费速度。稀疏索引将查找降为 O(log N) 的二分 + 少量顺序扫描,代价只是每个 Segment 多几 KB 的索引文件,性价比极高。此外,`.timeindex` 文件支持按时间戳查找消息(用于"从某时刻开始消费"的场景),仅靠 `.log` 文件无法高效实现。
|
||
|
||
## 关联笔记
|
||
|
||
- [[04-存储引擎/8-MQ-存储引擎设计|MQ 存储引擎设计]]
|
||
- [[04-存储引擎/10-RocketMQ-CommitLog|RocketMQ CommitLog]]
|
||
- [[07-主流MQ对比/22-Kafka|Kafka]]
|
||
- [[10-监控与运维/40-MQ-性能调优|MQ 性能调优]]
|
||
- [[12-架构与实战/50-MQ-设计与实现|MQ 设计与实现]]
|