RocketMQ源码级实现原理02-存储系统概要设计
2026/9/11 10:25:28 网站建设 项目流程

最新的NVME协议的SSD硬盘,磁盘存储:

文件写入速度,顺序读写 3g/s左右,随机读写2g/s左右(内存条的读写速度是10个g左右)

存储系统概要设计

主要就是讲消息从producer到达broker后,怎么落地存储

CommitLog的写是随机写,但是读可能会有随机读,当然也可能直接读pagecache

大方向是对的——这是两者存储模型根本差异决定的经典对比。但"RocketMQ 完全不会下降"要打个折扣。我按代码和架构说清楚。

Kafka 为什么会随 topic(更准确说是分区)数下降

Kafka 的存储模型是:每个 topic-partition 是一个独立的日志目录 + 独立的一组文件(.log / .index / .timeindex)。

- 写入时,broker 同时在给"多个分区的文件"追加 → 分区数一多,写入就从"顺序写一个大文件"变成"在几十上百个文件之间来回写",即随机 I/O;
- 每个分区都要占用文件句柄和 page cache,分区多了缓存被切碎、命中率下降,且各分区的尾部数据互相挤占缓存;
- 每个分区还有独立的索引、元数据,分区数增长还会带来更长的 rebalance、leader 选举、controller 元数据压力等(这些和写入无直接关系,但拖累整体)。

所以 Kafka 在分区数(不是单纯的 topic 数)到几百上千时,写入性能会明显下滑——社区普遍建议单 broker 分区数别太多。到 100 这个量级是否"明显",取决于磁盘(HDD
明显、SSD 弱化)、缓存大小和配置。

RocketMQ 为什么写入路径几乎不受 topic 数影响

RocketMQ 用的是**"单 CommitLog + 消费队列索引"**模型:

1. 所有 topic 的所有消息,全部追加到同一个 CommitLog 文件(DefaultMessageStore.commitLog = new CommitLog(this),putMessage
就是往这一个文件顺序追加)——写入永远是单文件顺序写,跟 topic 数量无关;
2. 追加完成后,由一个后台线程ReputMessageService(DefaultMessageStore.java:1864)从 commitlog 顺序读取,再把"消息位置"分发写进各个 topic-queue 对应的
ConsumeQueue;
3. ConsumeQueue 每条记录是定长 20 字节(ConsumeQueue.CQ_STORE_UNIT_SIZE = 20,只是"物理偏移+大小+tag hash"),很小。

关键点:真正贵的那一步(消息数据落盘)始终是单文件顺序写,不管你有 1 个还是 1000 个 topic。多 topic 只影响第二步"分发索引",而那是小的定长记录、且是异步的。这就是
RocketMQ 敢宣称"支持海量 topic"的底气。

一个好的文件系统,它的性能肯定是要优于分布式KV和newSQL;

Kafka和RocketMQ存储模型的差别

Kafka,会有多个文件存储,不像rocketmq总是写一个文件,能实现一直往文件末尾追加的随机写

RocketMQ存储架构图

1. ConsumeQueue里面放的是一个个索引条目,索引条目有三个字段:CommitLogOffSet ,MessageSize、TagHashCode

2. doDispatch 异步线程 构建ConsumeQueue

CommitLog底层结构

这么多MappedFile,如何管理起来对外暴露,就用一个MappedFileQueue;

Message格式

存储模块实现

总体代码层级设计架构

存储这块的设计很好,层次很清晰

存储逻辑层,比如CommitLog类中出了对外暴露putMessage()的接口外,还需要提供一些内部类,比如FlushRealTimeService这类的刷盘线程逻辑

存储IO层,则直接和PageCache、和磁盘打交道

各业务拥有独立线程池

也就是,不同的逻辑处理,RocketMQ都会为它分配不同的线程池,比如这里,为了处理发送端发过来的消息写入磁盘,就通过sendMssageExecutor线程池来执行sendMssageProcessor中的代码,完成一条消息的写入磁盘

CommitLog写入模块

未关闭自动创建topic开关

生产中,需要关闭自动创建topic

DefaultMessageStore#putMessage()

可以看到这里指定了每个mappedFile大小是10M,实际中这个大小是1G

commitLog作为共享资源写入要加锁

commitLog.putMessage()会有加锁释放锁的逻辑,因为要保证同一时间,只有一个线程去往MappedFile中写入消息数据

byteBuffer.position(),代表的就是当前在00000000000000000174080这个commitlog文件中的相对偏移,表示在当前要写入的消息之前,上一条消息已经写到了174080 + byteBuffer.position()的位置了,所以当前要写入的消息是从174080 + byteBuffer.position()往后开始写入

因为一个消息不能跨两个mappedFile,所以当最后一个mappedFile的剩余空间小于 (当前文件大小 + 8个字节的结尾魔数),则就要开始创建下一个mappedFile,把当前消息写入下一个mappedFile了

ConsumeQueue写入模块

物理结构

Broker端,每个topic下的各个队列queue持久化后的消费进度数据

主要代码组件

有独立的线程FlushConsumeQueueService线程,把构建好的consumeQueue对应的mappedFile刷入磁盘

代码实现

已经刷过盘的消息offset - 当前已经reput过的消息offset,就是剩余的还需要reput的消息

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询