vault backup: 2026-06-08 23:08:57
This commit is contained in:
@@ -51,24 +51,91 @@ ConsumeQueue 是 CommitLog 和 Consumer 之间的桥梁。每个 Topic 的每个
|
||||
| 字段 | 大小 | 说明 |
|
||||
|------|------|------|
|
||||
| CommitLog Offset | 8 字节 | 消息在 CommitLog 中的物理偏移 |
|
||||
| Size | 4 字节 | 消息长度 |
|
||||
| Size | 4 字节 | 消息在 CommitLog 中的存储大小(含消息头 + 消息体) |
|
||||
| Tag HashCode | 8 字节 | 消息 Tag 的哈希值,用于过滤 |
|
||||
|
||||
Consumer 拉取消息时,先读 ConsumeQueue 获取 offset 列表,再根据 offset 去 CommitLog 读取完整消息。由于 ConsumeQueue 的条目是定长的,可以直接通过下标计算偏移量进行随机读取,效率很高。
|
||||
|
||||
Tag 过滤也是在 ConsumeQueue 层完成的:Consumer 携带订阅的 Tag 哈希,Broker 遍历 ConsumeQueue 条目时先比对 Tag HashCode,不匹配的直接跳过,避免了读取 CommitLog 的开销。
|
||||
|
||||
> [!question] ConsumeQueue 只存了 Tag 的哈希值,如果两个不同的 Tag 碰撞到了同一个 HashCode,会发生什么?
|
||||
>
|
||||
> 这确实是一个**哈希碰撞**的场景。ConsumeQueue 的 Tag 过滤只是一个**粗过滤(Bloom Filter 思想)**:不匹配的一定跳过,匹配的还需要到 CommitLog 中读取完整消息再做精确的 Tag 字符串比对。因此碰撞不会导致消息丢失,只会少量增加无效的 CommitLog 读取,实际上概率极低。
|
||||
|
||||
**ConsumeQueue 文件结构**:每个 ConsumeQueue 文件由一个 **32 字节的固定文件头** 和后续的定长条目组成。文件头记录了该队列的元信息:
|
||||
|
||||
| 字段 | 大小 | 说明 |
|
||||
|------|------|------|
|
||||
| MinPhysicOffset | 8 字节 | 当前 CQ 引用的最小 CommitLog 物理偏移 |
|
||||
| MinLogicOffset | 8 字节 | 当前 CQ 的最小逻辑偏移(用于计算条目下标) |
|
||||
| IndexCount | 4 字节 | 当前已写入的条目总数 |
|
||||
|
||||
每个 ConsumeQueue 文件默认存储约 30 万个条目(约 5.72 MB),写满后自动滚动到下一个文件。
|
||||
|
||||
**消费进度(ConsumerOffset)**:Consumer 的消费进度并不是存在 ConsumeQueue 里,而是由 Broker 端的 `ConsumerOffsetManager` 单独管理,持久化到 `$HOME/store/config/consumerOffset.json`。进度的 key 是 `Topic@ConsumerGroup@QueueId`,value 是已消费到的 ConsumeQueue 逻辑偏移量。这与 Kafka 将 offset 存到 `__consumer_offsets` Topic 的设计形成了有趣的对比。
|
||||
|
||||
### 4. IndexFile:按 Key 查消息
|
||||
|
||||
IndexFile 是可选的哈希索引文件,结构类似 HashMap:通过 MessageKey 的哈希值定位到槽位(Slot),每个槽位指向一个链表,链表节点存储消息的 CommitLog Offset 和时间戳。
|
||||
IndexFile 是可选的哈希索引文件,结构类似 HashMap,每个文件默认大小为 400MB,由三部分组成:
|
||||
|
||||
| 区域 | 大小 | 说明 |
|
||||
|------|------|------|
|
||||
| IndexHeader | 40 字节 | 文件元信息:beginTimestamp、endTimestamp、beginPhyOffset、endPhyOffset、hashSlotCount、indexCount |
|
||||
| SlotTable | 500万 × 4 字节 | 每个槽位存储该哈希桶中**最新一条**索引条目的编号(逻辑下标) |
|
||||
| IndexArea | 2000万 × 20 字节 | 索引条目区,每条 20 字节:keyHash(4B) + phyOffset(8B) + timeDiff(4B) + slotValue(4B) |
|
||||
|
||||
查找流程:对 MessageKey 取哈希 → 在 SlotTable 中定位槽位 → 拿到该槽位最新条目的编号 → 在 IndexArea 中沿链表(`slotValue` 指向前一条同哈希条目)遍历比对 → 找到匹配的 CommitLog 偏移量。
|
||||
|
||||
```mermaid
|
||||
graph LR
|
||||
Key["MessageKey"] -->|"hash % 5000000"| Slot["SlotTable 槽位"]
|
||||
Slot -->|"最新条目编号"| Entry3["IndexArea 条目 3"]
|
||||
Entry3 -->|"slotValue"| Entry1["IndexArea 条目 1"]
|
||||
Entry1 -->|"slotValue = 0"| End["链表结束"]
|
||||
|
||||
style Key fill:#4A90D9,color:#fff
|
||||
style Slot fill:#F5A623,color:#fff
|
||||
style Entry3 fill:#6EC1E0,color:#fff
|
||||
style Entry1 fill:#6EC1E0,color:#fff
|
||||
```
|
||||
|
||||
> [!question] SlotTable 每个槽位只存一个编号,如果同一个哈希桶有多条消息,怎么找到它们?
|
||||
>
|
||||
> 这就是链表的设计:每个 IndexArea 条目的 `slotValue` 字段存的是**前一条同哈希消息的条目编号**。新消息写入时,先读取当前槽位的编号作为前驱,写入新条目后把新条目的编号更新到槽位。这是一个**头插法链表**,查找时从最新往最旧遍历。
|
||||
|
||||
这使得 RocketMQ 支持按 MessageKey 精确查询消息,也支持按时间范围回溯消息——在排查问题、重放历史消息时非常有用。
|
||||
|
||||
> [!tip] 实用建议
|
||||
> IndexFile 是可选的。如果业务不需要按 Key 查询消息(大多数纯消费场景不需要),可以关闭索引构建以减少 I/O 开销。
|
||||
|
||||
### 5. MappedFile 与 mmap
|
||||
|
||||
RocketMQ 使用 `MappedFile` 抽象封装了 mmap(内存映射文件)操作。每个 CommitLog / ConsumeQueue / IndexFile 文件都对应一个 MappedFile 实例。
|
||||
RocketMQ 使用 `MappedFile` 抽象封装了 mmap(内存映射文件)操作。每个 CommitLog / ConsumeQueue / IndexFile 文件都对应一个 MappedFile 实例。多个 MappedFile 由 `MappedFileQueue` 统一管理,形成一个逻辑上连续的大文件。
|
||||
|
||||
mmap 的核心优势:将磁盘文件映射到进程虚拟内存空间,写入时直接操作内存,由操作系统内核负责异步刷盘。相比 `write()` 系统调用,mmap 减少了一次内核态到用户态的数据拷贝,写入性能更高。
|
||||
**mmap 的核心优势**:将磁盘文件映射到进程虚拟内存空间后,写入操作直接操作 Page Cache 中的内存页,由操作系统内核负责异步刷盘。相比 `write()` 系统调用需要将用户空间缓冲区的数据 **CPU 拷贝**到内核空间,mmap 的写入直接落到了 Page Cache,省去了这次拷贝——写入性能更高。
|
||||
|
||||
**MappedFileQueue 管理机制**:
|
||||
|
||||
```mermaid
|
||||
graph LR
|
||||
MFQ["MappedFileQueue"] --> MF1["MappedFile 0<br/>0000000000"]
|
||||
MFQ --> MF2["MappedFile 1<br/>0000000001"]
|
||||
MFQ --> MF3["MappedFile 2 (Active)<br/>0000000002"]
|
||||
|
||||
MF1 -->|"写满,只读"| RO["不可写入"]
|
||||
MF2 -->|"写满,只读"| RO
|
||||
MF3 -->|"当前写入"| WR["可追加写入"]
|
||||
|
||||
style MFQ fill:#4A90D9,color:#fff
|
||||
style MF3 fill:#2ECC71,color:#fff
|
||||
style MF1 fill:#999,color:#fff
|
||||
style MF2 fill:#999,color:#fff
|
||||
```
|
||||
|
||||
MappedFileQueue 维护所有 MappedFile 的有序列表,`GetLastMappedFile()` 始终返回当前活跃的(正在写入的)文件。当活跃文件写满时,自动创建新的 MappedFile 并追加到队列尾部。
|
||||
|
||||
> [!tip] 内存锁定(mlock)
|
||||
> 在生产环境中,建议通过 `mlockall()` 将进程的内存页锁定,防止操作系统在内存紧张时将 Page Cache 换出到 swap。一旦 mmap 映射的内存页被 swap 到磁盘,读写操作会触发严重的缺页中断,导致性能断崖式下降。RocketMQ 提供了 `lockInMmap` 配置项来启用此特性。
|
||||
|
||||
```go
|
||||
// 简化版 CommitLog 写入逻辑
|
||||
@@ -98,12 +165,18 @@ func (cl *CommitLog) PutMessage(msg *Message) error {
|
||||
|
||||
### 6. 同步/异步刷盘与主从同步
|
||||
|
||||
**刷盘策略**决定了消息写入 MappedFile(内存)后何时真正持久化到磁盘:
|
||||
**刷盘策略**决定了消息写入 MappedFile(内存)后何时真正持久化到磁盘。RocketMQ 内部有两个核心刷盘实现类:
|
||||
|
||||
- **同步刷盘(SYNC_FLUSH)**:写入后立即调用 `fsync`,数据安全性高,但吞吐较低。适合对可靠性要求极高的场景(如金融交易)。
|
||||
- **异步刷盘(ASYNC_FLUSH)**:由后台 FlushCommitLogService 线程定时刷盘,吞吐高但宕机时可能丢失少量消息。适合大多数业务场景。
|
||||
| 刷盘策略 | 实现类 | 行为 | 适用场景 |
|
||||
|----------|--------|------|----------|
|
||||
| 同步刷盘(SYNC_FLUSH) | `GroupCommitService` | 写入后阻塞等待 `fsync` 完成才返回 ACK | 金融交易等对可靠性要求极高的场景 |
|
||||
| 异步刷盘(ASYNC_FLUSH) | `FlushCommitLogService` | 后台线程定时(默认 500ms)刷盘,写入 Page Cache 即返回 | 大多数业务场景,吞吐优先 |
|
||||
|
||||
**主从同步**方面,RocketMQ 4.x 引入了基于 Raft 协议的 **DLedger** 模式替代传统的 Master-Slave 异步复制。DLedger 要求消息写入多数节点后才算成功,从根本上解决了主从异步复制时 Master 宕机丢数据的问题。
|
||||
> [!question] 同步刷盘时,`GroupCommitService` 是如何做到"批量高效"的?
|
||||
>
|
||||
> `GroupCommitService` 内部使用了**批处理**机制:当多个 Producer 线程几乎同时写入时,它们会被放入一个请求批次。刷盘线程对整个批次执行一次 `fsync`,然后唤醒所有等待的 Producer 线程。这样 N 个并发写入只需要一次磁盘 I/O,显著降低了同步刷盘的性能损耗。
|
||||
|
||||
**主从同步**方面,RocketMQ 4.4+ 引入了基于 Raft 协议的 **DLedger** 模式替代传统的 Master-Slave 异步复制。DLedger 要求消息写入多数节点后才算成功,从根本上解决了主从异步复制时 Master 宕机丢数据的问题。RocketMQ 5.0 进一步演进为 **Controller 模式**,将 DLedger 的 Raft 能力集成到 Broker 内部,不再需要独立部署 DLedger 组件,运维更加简洁。
|
||||
|
||||
```mermaid
|
||||
graph LR
|
||||
@@ -119,6 +192,19 @@ graph LR
|
||||
style Slave2 fill:#6EC1E0,color:#fff
|
||||
```
|
||||
|
||||
### 7. 文件过期清理机制
|
||||
|
||||
RocketMQ 的存储文件不会无限增长,有独立的清理线程 `CleanCommitLogService` 定期扫描并删除过期文件。清理触发条件有两个(满足任一即清理):
|
||||
|
||||
- **按时间**:文件最后修改时间距今超过 `fileReservedTime`(默认 72 小时)
|
||||
- **按磁盘空间**:磁盘使用率超过 `diskMaxUsedSpaceRatio`(默认 75%),从最旧文件开始强制删除
|
||||
|
||||
清理时的一个关键约束:**不能删除被 ConsumeQueue 引用的 CommitLog 文件**。CleanCommitLogService 会先检查所有 ConsumeQueue 的最小引用偏移量(`MinPhysicOffset`),只有 CommitLog 文件的尾部偏移量小于该值时才能安全删除。ConsumeQueue 和 IndexFile 的清理则由 `CleanConsumeQueueService` 和 `CleanIndexFileService` 以类似逻辑跟进。
|
||||
|
||||
> [!question] 如果一个消费者长期不消费,导致 ConsumeQueue 的 MinPhysicOffset 一直很小,旧的 CommitLog 文件无法删除,磁盘满了怎么办?
|
||||
>
|
||||
> 这就是为什么 RocketMQ 会监控消费积压情况并发出告警。极端情况下,`diskMaxUsedSpaceRatio` 触发的强制清理机制会丢弃消费者尚未消费的消息来保全系统——宁可丢消息也不能让 Broker 挂掉。
|
||||
|
||||
## 关联笔记
|
||||
|
||||
- [[04-存储引擎/8-MQ-存储引擎设计|MQ 存储引擎设计]]
|
||||
|
||||
@@ -7,7 +7,7 @@ create time: 2026-05-24 19:52
|
||||
|
||||
## 概述
|
||||
|
||||
RabbitMQ 基于 Erlang/OTP 平台构建,消息存储依赖 Erlang 进程模型和 Mnesia 分布式数据库。本文深入分析 RabbitMQ 的消息持久化机制、队列类型演进(经典队列、仲裁队列、流式队列)、内存管理策略,以及消息从写入到删除的完整生命周期。
|
||||
RabbitMQ 基于 Erlang/OTP 平台构建,其存储体系分为两层:**元数据**存储在 Mnesia 分布式数据库中(Exchange 定义、Binding 关系等),**消息体**则由各队列进程通过 ETS 表和磁盘文件自行管理。本文深入分析 RabbitMQ 的消息持久化机制、队列类型演进(经典队列、仲裁队列、流式队列)、内存管理策略,以及消息从写入到删除的完整生命周期。
|
||||
|
||||
## 正文
|
||||
|
||||
@@ -15,7 +15,9 @@ RabbitMQ 基于 Erlang/OTP 平台构建,消息存储依赖 Erlang 进程模型
|
||||
|
||||
RabbitMQ 的每个 Queue 本质上是一个 Erlang 进程(gen_server),拥有独立的 mailbox 和状态。这种设计天然隔离了不同队列的故障,但也意味着每个队列的资源消耗(内存、调度时间)受到 Erlang 虚拟机(BEAM)的约束。
|
||||
|
||||
持久化元数据(Exchange 定义、Binding 关系、用户权限等)存储在 **Mnesia** 分布式数据库中。Mnesia 是 Erlang 原生的数据库,支持事务和分布式复制,但它的设计目标是存储小量元数据,而非大规模消息数据。实际的消息体存储由各队列进程自行管理。
|
||||
持久化元数据(Exchange 定义、Binding 关系、用户权限等)存储在 **Mnesia** 分布式数据库中。Mnesia 是 Erlang 原生的数据库,支持事务和分布式复制,但它的设计目标是存储小量元数据,而非大规模消息数据。
|
||||
|
||||
实际的消息体存储由各队列进程自行管理,核心依赖 Erlang 的 **ETS(Erlang Term Storage)** 表实现内存层缓存。ETS 是 BEAM 虚拟机内置的高性能键值存储,读写无需经过消息传递,单进程即可高效访问。经典队列在内存中使用 ETS 表暂存待投递的消息,磁盘层则通过自定义的文件存储引擎持久化。仲裁队列和流式队列则各自有不同的底层实现。
|
||||
|
||||
```mermaid
|
||||
graph TD
|
||||
@@ -36,30 +38,35 @@ graph TD
|
||||
|
||||
### 2. 消息持久化机制
|
||||
|
||||
RabbitMQ 中消息要真正持久化,需要同时满足两个条件:
|
||||
RabbitMQ 中消息要真正持久化,需要同时满足三个条件:
|
||||
|
||||
- **Queue 声明为 durable**:队列的元数据在 Broker 重启后保留。
|
||||
- **Message 的 delivery mode 设为 2(persistent)**:消息体写入磁盘。
|
||||
- **启用 Publisher Confirms**:生产者端获得 Broker 的写入确认,确保消息不会在传输窗口中丢失(详见 §7)。
|
||||
|
||||
> [!question] durable 和 persistent 有什么区别?
|
||||
>
|
||||
> `durable` 是队列的属性,控制的是"队列本身在重启后是否还存在";`persistent` 是消息的属性,控制的是"这条消息是否需要写入磁盘"。两者必须同时设置才能实现真正的持久化。如果只有 durable 没有 persistent,队列重启后还在,但里面的消息会丢失;反过来则队列重启后直接消失,消息也就无从谈起了。
|
||||
> `durable` 是队列的属性,控制的是"队列本身在重启后是否还存在";`persistent` 是消息的属性,控制的是"这条消息是否需要写入磁盘"。两者必须同时设置才能实现真正的持久化。如果只有 durable 没有 persistent,队列重启后还在,但里面的消息会丢失;反过来则队列重启后直接消失,消息也就无从谈起了。不过即使两者都设置了,还需要配合 **Publisher Confirms** 才能让生产者端确认消息已被 Broker 安全接收(见 §7)。
|
||||
|
||||
消息写入磁盘的时机并非立即的。RabbitMQ 使用一种"惰性"写入策略:消息先写入内存,当满足一定条件(如内存压力达到阈值、或显式 flush)时才批量刷入磁盘。这意味着在极端情况下(如 Broker 突然崩溃),少量消息可能丢失。
|
||||
消息写入磁盘的时机因消息类型而异。对于**持久化消息**(`delivery mode = 2`),经典队列会将其立即提交给独立的 **persister 进程**,由 persister 异步地批量刷入磁盘。消息在被 persister 确认写入磁盘之前,会一直保存在内存中,所以即使 Broker 崩溃也不会丢。但对于**非持久化消息**,它们只存在于内存中,Broker 重启后即丢失。
|
||||
|
||||
> [!tip] persister 进程
|
||||
>
|
||||
> 经典队列的磁盘写入并不由队列进程自己完成,而是委托给一个专门的 persister 进程。这种异步设计避免了 I/O 阻塞队列调度,但也意味着"消息写入磁盘"和"消息被队列接受"之间存在一个短暂的窗口。这就是为什么需要 **Publisher Confirms**(见下文)来给生产者端到端的持久化保证。
|
||||
|
||||
### 3. Queue 类型演进
|
||||
|
||||
RabbitMQ 3.8+ 引入了三种队列类型,各自面向不同的可靠性与性能需求:
|
||||
|
||||
**经典队列(Classic Queue)**是 RabbitMQ 最早的队列实现。消息存储在 Erlang Mnesia 表中(实际上从 3.x 开始使用自定义的文件存储引擎),支持可选持久化。单 Master 架构,不支持复制。在消息堆积时性能急剧下降,因为队列进程需要维护一个按 offset 索引的磁盘文件结构,大量消息会导致 GC 压力和 I/O 放大。
|
||||
**经典队列(Classic Queue)**是 RabbitMQ 最早的队列实现,经历了底层存储的重大演进:早期版本使用 Mnesia 表存储消息,从 3.x 起改用自定义的**文件存储引擎**(类似日志结构的段文件)。经典队列支持可选持久化,采用单 Master 架构,不支持原生复制(镜像队列插件已废弃)。在消息堆积时性能急剧下降,因为队列进程需要维护按消息 ID 索引的磁盘文件结构,大量消息会导致 Erlang 进程 GC 压力增大和 I/O 放大。
|
||||
|
||||
**仲裁队列(Quorum Queue)**基于 Raft 共识协议,每个队列在多个节点上维护副本(奇数个,默认 3 个)。消息写入需要多数节点确认,天然保证了数据安全。底层使用 WAL(Write-Ahead Log)顺序写盘,性能在大量消息堆积时保持稳定。
|
||||
**仲裁队列(Quorum Queue)**基于 Raft 共识协议,每个队列在多个节点上维护副本(奇数个,默认 3 个)。消息写入需要多数派(majority)确认后才返回成功,天然保证了数据安全。底层使用 **WAL(Write-Ahead Log)** 顺序写盘——所有队列共享同一个 WAL 文件,避免了经典队列的随机 I/O 问题。消息被消费并 ACK 后,通过 Raft 快照(snapshot)机制定期清理已确认的记录,性能在大量消息堆积时保持稳定。
|
||||
|
||||
**流式队列(Stream)**是 3.9 引入的新类型,借鉴了 Kafka 的设计思想。消息以 append-only log 形式存储,支持多消费者独立读取(非竞争消费),消息不会被消费后删除,支持回溯重放。适合事件驱动、审计日志等场景。
|
||||
**流式队列(Stream)**是 3.9 引入的新类型,借鉴了 Kafka 的设计思想。消息以 **append-only log**(段文件,segment file)形式存储在磁盘上,每个 segment 有固定大小,写满后滚动到下一个。支持多消费者独立读取(非竞争消费),各消费者维护自己的 offset,消息不会被消费后删除,支持按时间或 offset 回溯重放。适合事件驱动、审计日志等场景。
|
||||
|
||||
| 特性 | Classic Queue | Quorum Queue | Stream |
|
||||
|------|:---:|:---:|:---:|
|
||||
| 复制 | 不支持 | Raft 多副本 | 可选 |
|
||||
| 复制 | 不支持 | Raft 多副本 | Leader-Follower |
|
||||
| 持久化 | 可选 | 必须 | 必须 |
|
||||
| 消费后删除 | 是 | 是 | 否 |
|
||||
| 消息回溯 | 不支持 | 不支持 | 支持 |
|
||||
@@ -83,7 +90,7 @@ RabbitMQ 的内存管理是一个三层防护体系:
|
||||
|
||||
**流控(Flow Control)**:当某个 Erlang 进程(如队列进程)的消息积压超过一定阈值,会主动向 TCP 连接进程发送 `pause` 信号,让生产者的 TCP 接收窗口降为 0,从而在协议层面实现背压。
|
||||
|
||||
**信用机制(Credit Flow)**:这是 Erlang 进程间的流控机制。每个进程维护一个 credit 计数器,消息传递消耗 credit,消费端处理完毕后归还 credit。当 credit 归零时,发送端阻塞,防止过快地往下游进程灌消息。
|
||||
**信用机制(Credit Flow)**:这是 Erlang 进程间的流控机制。每个进程维护一个 credit 计数器(初始值默认 200,称为 `credit_flow_default_credit`),每传递一条消息消耗 1 个 credit,消费端处理完毕后归还 credit(默认每处理一条归还 1 个)。当 credit 归零时,发送端进程被阻塞(blocked),不再向下游发送,直到收到下游归还的 credit 重新恢复。这个机制保证了 Erlang 进程间的"管道"不会溢出。
|
||||
|
||||
```mermaid
|
||||
graph TD
|
||||
@@ -103,44 +110,84 @@ graph TD
|
||||
|
||||
消息从进入 RabbitMQ 到最终消失,经历以下阶段:
|
||||
|
||||
1. **写入内存**:消息到达队列进程后,先存入进程内存(Erlang ETS 表或进程状态)。
|
||||
2. **写入磁盘**(如果 persistent):当内存压力达到一定阈值,或消息被标记为 persistent 时,消息体被写入磁盘文件。
|
||||
3. **投递消费者**:消息从内存中读取并发送给消费者。
|
||||
4. **等待 ACK**:消息在未收到 ACK 前不会被删除。
|
||||
1. **写入内存**:消息到达队列进程后,先存入进程内存(ETS 表)。
|
||||
2. **写入磁盘**(如果 persistent):持久化消息被立即提交给 persister 进程异步写入磁盘文件;非持久化消息仅保留在内存中。
|
||||
3. **投递消费者**:消息从内存中读取并发送给消费者。此时消息被标记为 "unacked",不会重复投递给其他消费者。
|
||||
4. **等待 ACK**:消息在未收到 ACK 前不会被删除。如果消费者断开连接,unacked 消息会被重新入队。
|
||||
5. **收到 ACK 后删除**:消息从内存和磁盘中移除。对于持久化消息,磁盘文件中的标记被清除,空间在文件 compaction 时回收。
|
||||
|
||||
### 7. Go 代码示例
|
||||
### 7. Publisher Confirms:生产端的持久化保证
|
||||
|
||||
前面提到,消息写入磁盘是异步的。`ch.Publish` 返回 `nil` 只表示消息进入了 Broker 进程的内存,并不保证已经落盘。要获得端到端的持久化保证,必须使用 **Publisher Confirms** 机制。
|
||||
|
||||
工作原理:生产者在通道上开启 confirm 模式(`channel.confirmSelect()`),之后每条消息都会被 Broker 分配一个序列号。当消息被成功处理(写入所有镜像/副本)后,Broker 回调 confirm 并携带该序列号;如果处理失败,回调会携带 nack 标记。
|
||||
|
||||
> [!question] Publisher Confirms 和 AMQP 事务(tx)有什么区别?
|
||||
>
|
||||
> AMQP 的 `tx.select/commit/rollback` 是同步阻塞的事务模式,性能很差(每秒只能处理几百条消息)。Publisher Confirms 是异步的,吞吐量比事务高 1-2 个数量级,是生产环境的唯一推荐方案。
|
||||
|
||||
### 8. 消息策略:TTL 与 Max-Length
|
||||
|
||||
RabbitMQ 提供了两个队列级参数来控制消息的生命周期和队列大小,防止无限堆积:
|
||||
|
||||
- **`x-message-ttl`**(毫秒):消息在队列中的最大存活时间,超时后自动丢弃(或进入死信队列)。适合对时效性敏感的业务场景。
|
||||
- **`x-max-length`**:队列中允许的最大消息数量(包括 unacked)。达到上限后新消息的处理策略由 `x-overflow` 决定:默认 `drop-head`(丢弃队头最老的消息),也可以设为 `reject-publish`(拒绝新消息)。
|
||||
|
||||
这两个参数在声明队列时通过 `x-arguments` 传入,也可以通过 Policy 在运行时动态设置。
|
||||
|
||||
### 9. Go 代码示例
|
||||
|
||||
```go
|
||||
package main
|
||||
|
||||
import (
|
||||
"log"
|
||||
|
||||
amqp "github.com/rabbitmq/amqp091-go"
|
||||
)
|
||||
|
||||
func main() {
|
||||
conn, _ := amqp.Dial("amqp://guest:guest@localhost:5672/")
|
||||
conn, err := amqp.Dial("amqp://guest:guest@localhost:5672/")
|
||||
if err != nil {
|
||||
log.Fatalf("连接失败: %v", err)
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
ch, _ := conn.Channel()
|
||||
ch, err := conn.Channel()
|
||||
if err != nil {
|
||||
log.Fatalf("打开通道失败: %v", err)
|
||||
}
|
||||
defer ch.Close()
|
||||
|
||||
// 声明持久化队列:durable = true, autoDelete = false, exclusive = false
|
||||
ch.QueueDeclare("order_events", true, false, false, false, amqp.Table{
|
||||
"x-queue-type": "quorum", // 使用仲裁队列,替代默认的经典队列
|
||||
// 使用仲裁队列(quorum)替代默认的经典队列,获得 Raft 多副本保障
|
||||
_, err = ch.QueueDeclare("order_events", true, false, false, false, amqp.Table{
|
||||
"x-queue-type": "quorum",
|
||||
})
|
||||
if err != nil {
|
||||
log.Fatalf("声明队列失败: %v", err)
|
||||
}
|
||||
|
||||
// 发布持久化消息:DeliveryMode = Persistent
|
||||
ch.Publish("", "order_events", false, false, amqp.Publishing{
|
||||
err = ch.Publish("", "order_events", false, false, amqp.Publishing{
|
||||
ContentType: "application/json",
|
||||
DeliveryMode: amqp.Persistent, // 消息持久化
|
||||
DeliveryMode: amqp.Persistent, // 消息持久化,配合 durable 队列实现真正持久化
|
||||
Body: []byte(`{"order_id":"12345","status":"created"}`),
|
||||
})
|
||||
if err != nil {
|
||||
log.Fatalf("发布消息失败: %v", err)
|
||||
}
|
||||
|
||||
log.Println("消息发布成功")
|
||||
}
|
||||
```
|
||||
|
||||
上面的代码展示了两个关键点:一是 `QueueDeclare` 设置 `durable: true` 并指定 `x-queue-type: quorum` 使用仲裁队列;二是 `Publishing` 设置 `DeliveryMode: amqp.Persistent` 确保消息持久化。两者缺一不可。
|
||||
|
||||
> [!tip] 生产环境建议启用 Publisher Confirms
|
||||
>
|
||||
> 代码中 `ch.Publish` 返回 `nil` 只代表消息进入了 TCP 缓冲区,并不代表 Broker 已持久化。要获得端到端的写入确认,需要调用 `ch.Confirm(false)` 启用确认模式,然后通过 `ch.NotifyPublish` 监听 Broker 的 ack/nack。这在高可靠场景下是必须的。
|
||||
|
||||
## 关联笔记
|
||||
|
||||
- [[07-主流MQ对比/23-RabbitMQ|RabbitMQ]]
|
||||
|
||||
@@ -18,6 +18,16 @@ MQ 的本质是一个"写多读多"的中转站:生产者不断写入消息,
|
||||
> [!question] 思考
|
||||
> 如果让你从零设计一个 MQ,你会选择把消息存在哪里?关系数据库?Redis?还是直接写文件?
|
||||
|
||||
> [!tip] 回答:直接写文件(Append-Only Log)
|
||||
>
|
||||
> **关系数据库**(如 MySQL):B+ 树索引面向随机读写优化,每条消息写入都要走索引更新,高并发下页分裂和随机 I/O 是严重瓶颈;消息过期清理需要 VACUUM,产生碎片。写入吞吐通常在万级 TPS,远不及文件顺序写。
|
||||
>
|
||||
> **Redis**:全量消息放内存不现实(MQ 消息量通常 TB 级),持久化方案(RDB 丢数据、AOF 本质又是写文件且重写有性能抖动),容量受限,消息堆积易 OOM。适合做热缓存,不适合做核心存储。
|
||||
>
|
||||
> **直接写文件**(Kafka、RocketMQ 的选择):顺序写磁头几乎不寻道,HDD 可达 600MB/s;配合零拷贝(`sendfile`)和 Page Cache,热数据读取几乎零磁盘 I/O;稀疏索引 + 偏移量定位,索引开销极低;过期 Segment 直接截断,无碎片。
|
||||
>
|
||||
> 一句话:MQ 的读写模式是"追加写 + 偏移量读",和 Append-Only Log 文件天然同构,任何额外的抽象层(B+ 树、内存数据结构)都是在为不需要的能力买单。
|
||||
|
||||
### 磁盘顺序写 vs 随机写
|
||||
|
||||
传统认知里,磁盘(尤其是机械硬盘 HDD)是性能瓶颈。但这里有个关键区分:**顺序写**和**随机写**的性能差距是数量级的。
|
||||
@@ -39,27 +49,29 @@ graph LR
|
||||
|
||||
### 零拷贝(Zero-Copy)技术
|
||||
|
||||
传统网络传输一条消息,数据需要经历 4 次拷贝和 4 次上下文切换:
|
||||
传统网络传输一条消息(`read` + `write`),数据需要经历 **4 次拷贝**和 **2 次上下文切换**:
|
||||
|
||||
1. 磁盘 → 内核缓冲区(DMA 拷贝)
|
||||
1. 磁盘 → 内核缓冲区(DMA 拷贝)—— `read` 系统调用触发
|
||||
2. 内核缓冲区 → 用户空间缓冲区(CPU 拷贝)
|
||||
3. 用户空间缓冲区 → Socket 缓冲区(CPU 拷贝)
|
||||
3. 用户空间缓冲区 → Socket 缓冲区(CPU 拷贝)—— `write` 系统调用触发
|
||||
4. Socket 缓冲区 → 网卡(DMA 拷贝)
|
||||
|
||||
`sendfile` 系统调用可以让数据直接从内核缓冲区传到网卡,跳过用户空间的两次拷贝,只需要 2 次 DMA 拷贝 + 1 次上下文切换。**Kafka 正是利用 sendfile 实现了消费者拉取数据时的零拷贝**,这让消费者读取消息几乎不消耗 CPU 资源。
|
||||
其中 2 次 CPU 拷贝是纯浪费——应用层根本没有修改数据,只是"搬运"。
|
||||
|
||||
`sendfile` 系统调用可以让数据直接从内核缓冲区传到网卡,跳过用户空间的两次 CPU 拷贝,只需要 **2 次 DMA 拷贝 + 1 次上下文切换**。**Kafka 正是利用 sendfile 实现了消费者拉取数据时的零拷贝**——消费者读取消息时,数据从磁盘经 Page Cache 直接到网卡,几乎不消耗 CPU 资源。
|
||||
|
||||
```mermaid
|
||||
graph LR
|
||||
subgraph "传统拷贝"
|
||||
D1["磁盘"] -->|"DMA"| K1["内核缓冲区"]
|
||||
K1 -->|"CPU"| U1["用户空间"]
|
||||
U1 -->|"CPU"| S1["Socket 缓冲区"]
|
||||
S1 -->|"DMA"| N1["网卡"]
|
||||
subgraph "传统拷贝 read+write"
|
||||
D1["磁盘"] -->|"DMA 1"| K1["内核缓冲区"]
|
||||
K1 -->|"CPU 拷贝"| U1["用户空间"]
|
||||
U1 -->|"CPU 拷贝"| S1["Socket 缓冲区"]
|
||||
S1 -->|"DMA 2"| N1["网卡"]
|
||||
end
|
||||
|
||||
subgraph "零拷贝 sendfile"
|
||||
D2["磁盘"] -->|"DMA"| K2["内核缓冲区"]
|
||||
K2 -->|"DMA"| N2["网卡"]
|
||||
D2["磁盘"] -->|"DMA 1"| K2["内核缓冲区"]
|
||||
K2 -->|"DMA 2"| N2["网卡"]
|
||||
end
|
||||
|
||||
style D1 fill:#F5A623,color:#fff
|
||||
@@ -97,9 +109,31 @@ MQ 通常不会每条消息都触发一次磁盘 I/O,而是将多条消息在
|
||||
为什么 MQ 不用数据库(如 MySQL)存储消息?核心原因有三:
|
||||
|
||||
1. **写入模式不匹配**:数据库面向"随机读写"优化,B+ 树索引在高并发写入下会成为瓶颈;而 MQ 是纯粹的"顺序追加 + 顺序读取",Append-Only 日志天然适配。
|
||||
2. **无索引开销**:消息不需要按内容检索,只需要按偏移量(offset)顺序读取,省去了索引维护的代价。
|
||||
2. **索引开销极低**:消息不需要按内容检索,只需要按偏移量(offset)顺序读取。配合**稀疏索引**(仅记录部分消息的 offset→物理文件位置映射),查找时先定位到稀疏索引区间,再在区间内顺序扫描,维护成本几乎可以忽略。
|
||||
3. **删除成本低**:过期日志直接截断文件头,不需要像数据库那样做 VACUUM 或标记删除。
|
||||
|
||||
### 文件分段(Segment)策略
|
||||
|
||||
Append-Only 日志文件不可能无限增长——否则文件越大,索引查找越慢,过期清理也无法"截断"。实际做法是将日志切分为多个固定大小的 **Segment 文件**:
|
||||
|
||||
```mermaid
|
||||
graph LR
|
||||
A["Segment 000000.log"] -->|"写满"| B["Segment 000001.log"]
|
||||
B -->|"写满"| C["Segment 000002.log"]
|
||||
C -->|"写满"| D["Segment 000003.log (Active)"]
|
||||
|
||||
style D fill:#2ECC71,color:#fff
|
||||
```
|
||||
|
||||
每个 Segment 对应一个**稀疏索引文件**(如 `000000.index`),记录部分 offset 到物理位置的映射。查找某条消息时:
|
||||
|
||||
1. 二分查找定位到目标 Segment
|
||||
2. 在该 Segment 的稀索引中找到最近的索引项
|
||||
3. 从该位置开始顺序扫描
|
||||
|
||||
> [!question] 思考
|
||||
> Segment 文件的大小应该如何设定?太小会导致文件数量爆炸,太大会让过期清理的粒度变粗。Kafka 默认 1GB,这个数字是基于什么考虑的?
|
||||
|
||||
### 存储写入流程总览
|
||||
|
||||
```mermaid
|
||||
@@ -108,8 +142,11 @@ graph TD
|
||||
B --> M["消息序列化写入内存缓冲区"]
|
||||
M --> BQ{"是否达到批量阈值?"}
|
||||
BQ -->|"否"| M
|
||||
BQ -->|"是"| FS["追加写入磁盘日志文件"]
|
||||
FS --> F{"刷盘策略?"}
|
||||
BQ -->|"是"| SEG["追加写入当前 Active Segment"]
|
||||
SEG --> FULL{"Segment 写满?"}
|
||||
FULL -->|"是"| NEW["创建新 Segment + 稀疏索引"]
|
||||
FULL -->|"否"| F{"刷盘策略?"}
|
||||
NEW --> F
|
||||
F -->|"同步"| FSYNC["fsync 确保落盘"]
|
||||
F -->|"异步"| PC["写入 Page Cache 后返回"]
|
||||
PC --> BG["后台线程定时 fsync"]
|
||||
@@ -126,6 +163,12 @@ graph TD
|
||||
下面的代码展示了一个最简的 Append-Only Log 的核心逻辑——顺序追加写入和按偏移量读取。
|
||||
|
||||
```go
|
||||
import (
|
||||
"encoding/binary"
|
||||
"os"
|
||||
"sync"
|
||||
)
|
||||
|
||||
// AppendLog 是一个简化版的顺序写日志存储
|
||||
type AppendLog struct {
|
||||
mu sync.Mutex
|
||||
@@ -181,7 +224,7 @@ func (a *AppendLog) ReadAt(offset int64) ([]byte, error) {
|
||||
}
|
||||
```
|
||||
|
||||
**核心要点**:写入只做追加(`O_APPEND`),不修改已有数据;读取通过偏移量直接定位,无需索引。这就是 Kafka、RocketMQ 存储引擎的最简原型。
|
||||
**核心要点**:写入只做追加(`O_APPEND`),不修改已有数据;读取通过偏移量直接定位,配合稀疏索引即可快速查找。这就是 Kafka、RocketMQ 存储引擎的最简原型。
|
||||
|
||||
## 关联笔记
|
||||
|
||||
|
||||
+163
-30
@@ -30,15 +30,68 @@ partition-0/
|
||||
|
||||
文件名是该 Segment 的 **base offset**(起始偏移量)。当 Segment 达到配置的大小或时间阈值时,Kafka 会"滚动"出一个新的 Segment。
|
||||
|
||||
### 消息物理格式:RecordBatch 与 Record
|
||||
|
||||
`.log` 文件中存储的不是裸消息,而是经过精心设计的二进制格式。Kafka 0.11+ 引入了 **RecordBatch** 结构——一个 Batch 包含多条 Record,整体写入 `.log` 文件:
|
||||
|
||||
```mermaid
|
||||
graph LR
|
||||
subgraph "RecordBatch"
|
||||
BH["Batch Header<br/>BaseOffset + Length + ..."]
|
||||
R1["Record 0"]
|
||||
R2["Record 1"]
|
||||
R3["Record 2"]
|
||||
end
|
||||
|
||||
BH --> R1 --> R2 --> R3
|
||||
|
||||
style BH fill:#4A90D9,color:#fff
|
||||
```
|
||||
|
||||
**RecordBatch Header** 的关键字段:
|
||||
|
||||
| 字段 | 大小 | 说明 |
|
||||
|------|------|------|
|
||||
| `BaseOffset` | 8B | 该 Batch 的起始 offset |
|
||||
| `Length` | 4B | 整个 Batch 的字节长度 |
|
||||
| `PartitionLeaderEpoch` | 4B | 用于检测 Leader 切换,防止过期数据写入 |
|
||||
| `Magic` | 1B | 格式版本号(v2 = 0x02) |
|
||||
| `CRC` | 4B | 整个 Batch 的校验和(CRC-32C),检测数据损坏 |
|
||||
| `MaxTimestamp` | 8B | Batch 中最大的消息时间戳,用于基于时间的索引和清理 |
|
||||
| `LastOffsetDelta` | 4B | 最后一条 Record 相对 BaseOffset 的偏移 |
|
||||
| `RecordCount` | 2B | Batch 中 Record 的数量 |
|
||||
|
||||
每条 **Record** 的结构更为紧凑:
|
||||
|
||||
```
|
||||
┌──────────────────────────────────────────┐
|
||||
│ Length (varint) │
|
||||
│ Attributes (1B) │
|
||||
│ Timestamp Delta (varint) │ ← 相对于 Batch 的 MaxTimestamp
|
||||
│ Offset Delta (varint) │ ← 相对于 Batch 的 BaseOffset
|
||||
│ Key Length (varint) + Key │
|
||||
│ Value Length (varint) + Value │
|
||||
│ Headers Count + Headers[] │
|
||||
└──────────────────────────────────────────┘
|
||||
```
|
||||
|
||||
> [!question] 思考
|
||||
> 为什么 Kafka 要把多条消息打包成 RecordBatch,而不是每条消息独立存储?提示:想想网络传输、磁盘 I/O 和 CRC 校验三个维度。
|
||||
|
||||
两个关键设计思想:
|
||||
|
||||
1. **增量编码**:每条 Record 的 offset 和 timestamp 只存**相对于 Batch 头的增量**(varint 编码),而非绝对值。消息越密集,增量越小,varint 编码字节数越少——典型场景下一条 Record 的元数据开销只有 5-6 字节。
|
||||
2. **整体验校验**:CRC-32C 覆盖整个 Batch 而非单条消息,减少了校验和的存储和计算开销,同时仍然能检测出任何位置的数据损坏。
|
||||
|
||||
```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"]
|
||||
P0 --> S1["Segment 0<br/>base offset = 0"]
|
||||
P0 --> S2["Segment 1<br/>base offset = 523840"]
|
||||
P0 --> S3["Segment 2<br/>base offset = 1048576"]
|
||||
|
||||
S1 --> L1["00000.log"]
|
||||
S1 --> I1["00000.index"]
|
||||
@@ -81,6 +134,10 @@ Kafka 不是每条消息都建索引——那会带来巨大的存储和维护
|
||||
|
||||
因为索引本身是有序的,二分查找效率为 O(log N);而顺序扫描的范围通常只有几条消息(4KB 内),开销极小。这种设计用极少的索引空间换来了接近 O(1) 的查找性能。
|
||||
|
||||
**索引加载机制**:`.index` 文件通过 **mmap** 映射到内存——Broker 启动或打开 Segment 时,并不读取整个索引文件,而是将其映射到虚拟地址空间。首次访问某页时触发缺页中断加载数据,之后的查找操作完全在内存中完成。索引文件默认最大 10MB(`log.index.size.max.bytes`),对于 1GB 的 Segment 只占约 1% 的额外空间。
|
||||
|
||||
**时间戳索引(`.timeindex`)**:与 `.index` 按 offset 建索引不同,`.timeindex` 记录的是 `timestamp → offset` 的映射。当消费者使用 `offsetsForTimes()` API(即"从某个时间点开始消费")时,Kafka 先在 `.timeindex` 中二分查找目标时间戳对应的 offset,再通过 `.index` 定位物理位置。`MaxTimestamp` 字段(RecordBatch Header 中)使得整个 Batch 只需一个时间戳条目,索引开销极低。
|
||||
|
||||
### 日志压缩(Log Compaction)
|
||||
|
||||
普通的消息保留策略是按时间或大小删除过期 Segment。但 Kafka 还提供了另一种策略——**Log Compaction**:保留每个 Key 的最后一条消息,删除之前的旧版本。
|
||||
@@ -94,9 +151,45 @@ Kafka 不是每条消息都建索引——那会带来巨大的存储和维护
|
||||
|
||||
**典型场景**:变更数据捕获(CDC)、状态快照同步。比如数据库的一张用户表,每次更新都发一条消息到 Kafka(Key = userId),下游系统通过 Log Compaction 拿到的就是每个用户的最新状态。
|
||||
|
||||
**Tombstone 记录**:如果想在 Compacted Topic 中**删除**某个 Key 的所有记录,Producer 需要发送一条 `value = null` 的消息(称为 Tombstone)。Cleaner 会先保留这条 Tombstone 一段时间(由 `log.cleaner.delete.retention.ms` 控制,默认 24 小时),之后才彻底清除该 Key 的所有痕迹。如果直接不发 Tombstone 而只是停止写入,Compaction 会永远保留该 Key 的最后一条有效消息。
|
||||
|
||||
**关键配置**:`log.cleaner.min.compaction.lag.ms`(默认 0)控制一条消息写入后至少保留多久才能被 Compaction 清理。这在"先写入后立即回读"的场景中非常重要——防止刚写入的消息还没被消费者读到就被清理了。
|
||||
|
||||
> [!question] 思考
|
||||
> Log Compaction 保证的是"每个 Key 至少保留最新一条",但如果有两个 Key 相同的消息几乎同时到达,Compaction 会保留哪条?
|
||||
|
||||
### 数据保留策略(Retention Policy)
|
||||
|
||||
除了 Log Compaction,Kafka 最常用的过期数据清理方式是**基于时间和大小的删除策略**。每个 Topic 可以独立配置:
|
||||
|
||||
- **`log.retention.hours`**(默认 168 小时 = 7 天):消息保留的最大时间。超过该时间的 Segment 会被标记删除。
|
||||
- **`log.retention.bytes`**(默认 -1,即不限制):每个 Partition 的最大保留大小。超出时从最旧的 Segment 开始删除。
|
||||
|
||||
两者是**或**的关系——任一条件触发就会执行清理。
|
||||
|
||||
Topic 的清理行为由 **`cleanup.policy`** 控制,有三个取值:
|
||||
|
||||
| 值 | 行为 | 适用场景 |
|
||||
|---|---|---|
|
||||
| `delete`(默认) | 按时间和大小删除过期 Segment | 大多数业务消息 |
|
||||
| `compact` | 保留每个 Key 的最新消息 | CDC、状态快照 |
|
||||
| `compact,delete` | 先 Compaction 再按时间删除 | 需要 Compaction 但也想限制历史深度 |
|
||||
|
||||
### `__consumer_offsets`:Consumer Offset 的存储
|
||||
|
||||
消费者提交的 offset 存在哪里?Kafka 把它存在一个**内置 Topic** `__consumer_offsets` 中(默认 50 个 Partition)。这个 Topic 有两个有趣的特性:
|
||||
|
||||
1. **使用 Compaction 策略**:每个 Key(格式为 `groupId + topic + partition`)只保留最新的 offset 值,历史提交自动被清理。
|
||||
2. **内部 Topic,不对外暴露**:消费者通过 `OffsetCommit` / `OffsetFetch` 请求间接读写它,不需要直接 produce/consume。
|
||||
|
||||
**Partition 路由**:一个消费者组的 offset 存储在哪个 Partition 中?Kafka 对 `groupId` 做哈希取模:`hash(groupId) % 50`。这保证了同一消费者组的所有 Topic-Partition 的 offset 都集中在同一个 `__consumer_offsets` Partition 上,便于事务性地批量提交和拉取。
|
||||
|
||||
> [!question] 思考
|
||||
> 如果 `__consumer_offsets` 使用 `delete` 策略而不是 `compact`,会导致什么问题?
|
||||
|
||||
> [!tip] 答案
|
||||
> 如果使用 `delete` 策略,当消费者组长时间不活跃(超过 `offsets.retention.minutes`,默认 10080 分钟 = 7 天),其提交的 offset 记录会被删除。当该消费者组重新上线时,找不到之前的消费位置,只能根据 `auto.offset.reset` 策略从头消费(`earliest`)或跳到最新(`latest`)——前者导致大量重复消费,后者导致消息丢失。`compact` 策略保证每个 Key 的最新 offset 永远不会被删除,除非该消费者组被显式废弃。
|
||||
|
||||
### 时间轮(Timing Wheel)
|
||||
|
||||
Kafka 内部有大量的定时任务:延迟消息、会话过期、日志清理等。如果每个定时任务都开一个 goroutine(或 Java 的 Timer),任务数多了之后,调度开销会非常大。
|
||||
@@ -106,19 +199,19 @@ Kafka 使用**层级时间轮(Hierarchical Timing Wheel)**来高效管理定
|
||||
- 时间轮是一个环形数组,每个槽位(slot)代表一个时间区间。
|
||||
- 新任务根据到期时间插入对应槽位。
|
||||
- 时钟每推进一个 tick,处理当前槽位的所有任务。
|
||||
- 当低层时间轮溢出时,任务会被"滴答"到上层时间轮,等待合适时机再降下来。
|
||||
- 当低层时间轮溢出时(即任务的到期时间超出了当前层的总跨度),任务会被"提升"到上层时间轮的某个槽位;随着时钟推进,上层的任务到期时会被重新分配到下层的精确槽位中。
|
||||
|
||||
时间轮的插入和删除都是 O(1) 操作,远优于优先队列的 O(log N)。
|
||||
时间轮的插入和删除都是 O(1) 操作——新任务只需计算目标槽位(`currentTime / tickDuration % wheelSize`),直接链表插入;远优于优先队列的 O(log N)。
|
||||
|
||||
```mermaid
|
||||
graph TD
|
||||
TW["时间轮"] --> L1["第一层: 1ms 精度\n20 个槽位"]
|
||||
TW --> L2["第二层: 20ms 精度\n20 个槽位"]
|
||||
TW --> L3["第三层: 400ms 精度\n20 个槽位"]
|
||||
TW["时间轮"] --> L1["第一层: 每槽位 1ms, 共 20 槽位"]
|
||||
TW --> L2["第二层: 每槽位 20ms, 共 20 槽位"]
|
||||
TW --> L3["第三层: 每槽位 400ms, 共 20 槽位"]
|
||||
|
||||
L1 --> SLOT1["slot 0: 任务 A"]
|
||||
L1 --> SLOT2["slot 5: 任务 B"]
|
||||
L2 --> SLOT3["slot 3: 任务 C\n溢出后降级到第一层"]
|
||||
L2 --> SLOT3["slot 3: 任务 C"]
|
||||
|
||||
style TW fill:#4A90D9,color:#fff
|
||||
style L1 fill:#6EC1E0,color:#fff
|
||||
@@ -143,6 +236,25 @@ Kafka 的读写路径深度依赖 Linux 内核的两个能力:
|
||||
> [!question] 思考
|
||||
> 如果 Kafka Broker 的内存足够大,Page Cache 能缓存大量数据。但如果消费者需要回溯到很早的消息(不在 Page Cache 中),会发生什么?性能会下降多少?
|
||||
|
||||
### Clean Shutdown 与 Recovery
|
||||
|
||||
Kafka 默认采用**异步刷盘**(`log.flush.interval.messages` 和 `log.flush.interval.ms` 控制),这意味着 Page Cache 中可能还有未落盘的数据。Broker 关闭或崩溃时的恢复机制至关重要:
|
||||
|
||||
**Clean Shutdown**(正常关闭):Broker 收到 SIGTERM 时,会:
|
||||
1. 停止接受新的 Produce/Fetch 请求
|
||||
2. 将所有 Active Segment 的 Page Cache 数据 **flush 到磁盘**
|
||||
3. 写入 `recovery-point-offset-checkpoint` 文件(记录每个 Partition 已成功刷盘的最新 offset)
|
||||
4. 关闭所有文件句柄
|
||||
|
||||
**Unclean Recovery**(崩溃恢复):如果 Broker 异常终止(SIGKILL、断电),恢复时需要:
|
||||
1. 读取上次保存的 `recovery-point-offset-checkpoint`
|
||||
2. 对每个 Partition,从 checkpoint 记录的 offset 开始,逐条验证 RecordBatch 的 CRC-32C 校验和
|
||||
3. 截断最后一个 CRC 校验失败的 RecordBatch 之后的所有数据
|
||||
4. 重建 `.index` 和 `.timeindex`(从 `.log` 文件扫描生成)
|
||||
|
||||
> [!tip] 实用建议
|
||||
> 在生产环境中,如果 `unclean.leader.election.enable=false`(默认),那么当一个 Partition 的所有同步副本(ISR)都宕机时,该 Partition 不会自动恢复服务。虽然这保证了数据一致性,但会导致分区不可用。需要在**数据安全**和**可用性**之间做出权衡。
|
||||
|
||||
### Go 代码:简化版 Segment 文件的读写
|
||||
|
||||
下面展示一个简化版的 Segment 文件实现,包含日志文件和稀疏索引的写入与查找。
|
||||
@@ -153,6 +265,8 @@ type Segment struct {
|
||||
logFile *os.File
|
||||
indexFile *os.File
|
||||
index []indexEntry // 内存中的稀疏索引
|
||||
writePos int64 // 当前写入位置(避免每次调用 stat.Size())
|
||||
nextOffset int64 // 下一条消息的 offset
|
||||
}
|
||||
|
||||
type indexEntry struct {
|
||||
@@ -161,62 +275,66 @@ type indexEntry struct {
|
||||
}
|
||||
|
||||
// Write 追加一条消息到 Segment
|
||||
func (s *Segment) Write(offset int64, data []byte) error {
|
||||
// 记录当前写入位置作为物理偏移
|
||||
stat, _ := s.logFile.Stat()
|
||||
pos := stat.Size()
|
||||
func (s *Segment) Write(data []byte) error {
|
||||
pos := s.writePos
|
||||
|
||||
// Length-Prefix 编码写入日志文件
|
||||
// Length-Prefix 编码写入日志文件: [4字节长度][数据]
|
||||
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
|
||||
}
|
||||
s.writePos += int64(4 + len(data))
|
||||
|
||||
// 每隔一定消息数写入一条索引(稀疏索引,不是每条都写)
|
||||
if len(s.index) == 0 || offset-s.index[len(s.index)-1].offset >= 4 {
|
||||
entry := indexEntry{offset: offset, position: int32(pos)}
|
||||
// 稀疏索引:每隔一定字节写入一条索引项
|
||||
// 简化演示:每隔 4 条消息写一条索引(实际 Kafka 按字节间隔)
|
||||
if len(s.index) == 0 || s.nextOffset-s.index[len(s.index)-1].offset >= 4 {
|
||||
entry := indexEntry{offset: s.nextOffset, position: int32(pos)}
|
||||
s.index = append(s.index, entry)
|
||||
// 同步写入 .index 文件(简化:每条 12 字节 = 8 offset + 4 position)
|
||||
// 同步写入 .index 文件(每条 12 字节 = 8 offset + 4 position)
|
||||
buf := make([]byte, 12)
|
||||
binary.BigEndian.PutUint64(buf[:8], uint64(offset))
|
||||
binary.BigEndian.PutUint64(buf[:8], uint64(s.nextOffset))
|
||||
binary.BigEndian.PutUint32(buf[8:], uint32(pos))
|
||||
s.indexFile.Write(buf)
|
||||
}
|
||||
s.nextOffset++
|
||||
return nil
|
||||
}
|
||||
|
||||
// Read 根据 offset 查找并读取消息
|
||||
func (s *Segment) Read(offset int64) ([]byte, error) {
|
||||
// 第一步:二分查找稀疏索引,找到最近的索引项
|
||||
// 第一步:二分查找稀疏索引,找到不大于目标 offset 的最近索引项
|
||||
idx := sort.Search(len(s.index), func(i int) bool {
|
||||
return s.index[i].offset > offset
|
||||
}) - 1
|
||||
|
||||
var startPos int32
|
||||
var scanOffset int64
|
||||
if idx >= 0 {
|
||||
startPos = s.index[idx].position // 从索引指向的位置开始
|
||||
startPos = s.index[idx].position
|
||||
scanOffset = s.index[idx].offset // 从索引记录的 offset 开始扫描
|
||||
} else {
|
||||
scanOffset = s.baseOffset // 没找到索引则从头开始
|
||||
}
|
||||
|
||||
// 第二步:从 startPos 开始顺序扫描 .log 文件
|
||||
// 第二步:从索引位置开始顺序扫描 .log 文件,逐条对比 offset
|
||||
pos := int64(startPos)
|
||||
for {
|
||||
for scanOffset <= offset {
|
||||
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 {
|
||||
if scanOffset == offset {
|
||||
data := make([]byte, length)
|
||||
s.logFile.ReadAt(data, pos+4)
|
||||
return data, nil
|
||||
}
|
||||
pos += int64(4 + length)
|
||||
pos += int64(4 + length) // 跳到下一条消息
|
||||
scanOffset++
|
||||
}
|
||||
return nil, fmt.Errorf("offset %d not found", offset)
|
||||
}
|
||||
```
|
||||
|
||||
@@ -227,6 +345,21 @@ func (s *Segment) Read(offset int64) ([]byte, error) {
|
||||
|
||||
答案是:**理论上行,但实践中代价太大**。如果只用 .log,消费者查找某个 offset 的消息时,必须从 Segment 头部开始顺序扫描,时间复杂度为 O(N)。在消息量巨大的场景下(单 Segment 可达 1GB、数百万条消息),这会严重拖慢消费速度。稀疏索引将查找降为 O(log N) 的二分 + 少量顺序扫描,代价只是每个 Segment 多几 KB 的索引文件,性价比极高。此外,`.timeindex` 文件支持按时间戳查找消息(用于"从某时刻开始消费"的场景),仅靠 `.log` 文件无法高效实现。
|
||||
|
||||
### 存储相关配置参数速查
|
||||
|
||||
| 参数 | 默认值 | 说明 |
|
||||
|------|--------|------|
|
||||
| `log.segment.bytes` | 1 GB | 单个 Segment 的最大大小,达到后滚动新 Segment |
|
||||
| `log.roll.ms` / `log.roll.hours` | 7 天 | 即使 Segment 未满,超过该时间也会滚动 |
|
||||
| `log.index.size.max.bytes` | 10 MB | 单个 `.index` 文件的最大大小 |
|
||||
| `log.index.interval.bytes` | 4 KB | 稀疏索引的写入间隔(每隔多少字节写一条索引) |
|
||||
| `log.retention.hours` | 168 (7天) | 消息最大保留时间 |
|
||||
| `log.retention.bytes` | -1 (不限) | 每个 Partition 的最大保留大小 |
|
||||
| `cleanup.policy` | delete | Topic 级清理策略:delete / compact / compact,delete |
|
||||
| `log.cleaner.min.compaction.lag.ms` | 0 | 消息写入后至少保留多久才能被 Compaction |
|
||||
| `log.cleaner.delete.retention.ms` | 24 小时 | Tombstone 记录的保留时间 |
|
||||
| `offsets.retention.minutes` | 10080 (7天) | Consumer Group 不活跃后 offset 的保留时间 |
|
||||
|
||||
## 关联笔记
|
||||
|
||||
- [[04-存储引擎/8-MQ-存储引擎设计|MQ 存储引擎设计]]
|
||||
|
||||
Reference in New Issue
Block a user