ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

消息中间件存储架构深度解析:Kafka、RocketMQ、JMQ对比

消息中间件存储架构深度解析:Kafka、RocketMQ、JMQ对比 工程师之夜系列分享第三十九篇我们专门聊聊 Kafka、RocketMQ、JMQ 这三款消息中间件的存储架构。存储架构既是消息系统性能的底线也是线上故障排查时最先要怀疑的环节——消息延迟、消费堆积、磁盘涨满、甚至丢消息很多问题最后都能回到“数据到底怎么落盘、怎么建索引”这件事上。这篇内容适合消息中间件运维、数据架构选型的同学也适合正在准备 Kafka 和 RocketMQ 面试的人我会从设计思路、索引机制、刷盘链路、生产案例四个层面展开最后给出一份可以直接拿去用的选型参考。1. 存储架构对比的底层问题消息队列为什么要单独设计一套存储1.1 三个系统各自的出身决定了它们的存储风格Kafka 最初是处理活动日志的产物天生就是“海量写入、高吞吐、允许一定延迟”的定位所以它的存储设计非常依赖操作系统的顺序写能力。RocketMQ 是老牌电商企业在交易场景里反复打磨出来的既要支撑大促峰值又要支持事务消息、延迟消息、消息回溯这些企业级特性因此它在存储上选择了“一份 CommitLog 多份队列索引”的重索引方案。JMQ 则是京东自研的消息中间件从早期版本逐步演进目标是支撑京东体量的订单、物流、营销等场景它吸收了前两者的很多思想但又有自己的取舍比如共享页缓存池、多级存储。同一个“消息存储”问题三家给出了三套差异明显的答案。1.2 存储要回答的三个问题第一写入不能成为瓶颈。消息队列是典型的写多读少系统所有消息都要先落盘才算“被接受”所以如何把随机写变成顺序写是所有存储架构的第一步。第二消费方要能高效定位消息。生产者写入后消费者不可能去全盘扫描必须通过索引从“某个 offset”快速拿到数据。第三系统要能处理消息堆积和回溯。堆积意味着磁盘上的数据是海量的回溯意味着消费者可以按时间或 offset 重新消费这对存储的保留策略、批量删除、索引方式都有很高要求。为了加深理解我习惯把“存储架构”拆成两个部分一个是“日志栈”解决写入一个是“索引栈”解决读取。Kafka、RocketMQ、JMQ 的差异说白了就是这两层如何组合以及每一层各自做到什么程度。2. 三大存储架构逐层拆解从日志文件到索引文件2.1 Kafka 的 Partition Segment 设计目录即分区段文件即日志Kafka 最底层的存储单元是 Partition分区每个 Partition 在磁盘上就是一个目录目录名通常是 topic-partition例如 my-topic-0。Kafka 存储设计的第一步是不做“全局共享日志”而是让每个分区独立维护日志文件。这样做的好处是分区内天然有序消费组可以按分区并行副本同步也只需要按分区处理。坏处也好理解如果 Topic 和分区非常多磁盘上会出现大量小目录对文件系统 inode 和页缓存都不友好。在单个分区目录里数据被切成若干个 Segment段文件。每个段文件是磁盘上的一个 .log 文件默认 1GB 滚动一次log.segment.bytes。Kafka 采用“段内顺序追加段间顺序滚动”的方式一旦当前段写满就新建一个段继续写同时老段可以被清理。与 .log 配套的是 .index 稀疏索引文件和 .timeindex 时间索引文件。稀疏索引的意思是不是每条消息都建索引而是每隔约 4KBlog.index.interval.bytes记录一条“相对 offset 与物理位置”的映射。查找时先二分定位到具体 Segment再在段内根据稀疏索引缩小范围最后从物理位置顺序扫几条。这里有一个很关键的设计Kafka 不主动刷盘。它把数据写进 Page Cache页缓存由操作系统决定何时落盘。生产环境下它的吞吐之所以高正是因为这个“看似偷懒”的设计——Append 到页缓存就好比把货物先堆在收货区理货员OS等空闲了再上架。但这也带来一个副作用如果机器突然断电或进程崩溃页缓存里的数据可能丢失。所以 Kafka 用副本机制来兜底消息在 leader 的日志里“被接受”后还需要被足够多的 follower 同步才会向 producer 返回成功acksall 就是最严格的副本同步模式。2.2 RocketMQ 的 CommitLog ConsumeQueue写一份读多份RocketMQ 没有选择“每 Topic 每分区独立日志”而是采用了一个非常像数据库预写日志WAL的架构所有 Topic 的所有消息都依次写入同一个 CommitLog 文件序列。CommitLog 在磁盘上按固定大小滚动默认 1GB 一个文件路径一般在 store/commitlog 下。生产者发来一条消息Broker 要做的第一件事就是把它追加到当前 CommitLog 的末尾所以这依然是一份极致的顺序写。但“所有消息写同一份日志”会带来一个读取难题消费者不知道自己的消息在 CommitLog 的哪个位置。所以 RocketMQ 在写入 CommitLog 后会异步地把“消息在 CommitLog 中的 offset、消息长度、Tag 哈希”这些信息分发到对应的 ConsumeQueue 索引文件里。ConsumeQueue 是按 Topic 和 Queue 组织的每个目录结构为 store/consumequeue/{topic}/{queueId}里面是固定 20 字节一条的索引记录。消费者消费时先根据 QueueId 和 offset 从 ConsumeQueue 读到索引再到 CommitLog 里用真正的存储偏移去取消息。这套架构的妙处在于无论 Topic 有多少、Queue 有多少CommitLog 始终只有一份写入路径不会因为“Topic 数量爆炸”而退化为随机写。代价是读取时多了一次索引查找消息量大的时候如果 CommitLog 对应的页缓存被大量冷数据挤走消费者拉取消息就可能出现磁盘随机读这就是 RocketMQ 在极高并发读场景下不如 Kafka 那么“省钱”的技术原因。RocketMQ 还额外提供了一个 IndexFile用来按消息 Key 做精确查找主要用于事务消息、消息查询这类场景它不会影响普通消费链路。2.3 JMQ 的存储演进从“借鉴 RocketMQ”到“独立设计”JMQ 是京东的消息中间件我在接触它时印象最深的一点是它经历了从共享型存储到独立存储、再到多级存储的演进。早期版本受基础设施限制多个 Broker 共享一套存储设备虽然节省成本但一个 Topic 的流量很容易影响邻居 Topic。后来走向本地磁盘独立存储才开始真正掌握“隔离”这个词的分量。从公开资料看JMQ 在核心存储模型上很接近 RocketMQ 的“主日志 队列索引”思路所有消息先写到一个主数据日志文件里再为每个 Topic/Queue 建立一份轻量级索引。但它又做了几个重要改良。第一个是共享页缓存池也就是把整个 Broker 的页缓存纳入统一管理避免某个大 Topic 把缓存全部占掉其他 Topic 直接变成冷读这本质上是把操作系统页缓存的管理权从“完全交给 OS”往“应用参与控制”拉了一大步。第二个是多级存储让热数据留在内存或 SSD温数据落到容量型磁盘冷数据进入归档它的存储架构不是“一层日志”而是按数据的冷热度分层放置。第三个是队列模型更扁平在京东的实践里可以支撑百万级队列这是 Kafka 的“每分区目录”和 RocketMQ 的“每 Queue 索引”在超大规模场景下都很难轻松做到的。因此对比这三大架构时不要简单说谁绝对最好。Kafka 把简单和顺序写到极致但牺牲了对海量 Topic 的原生友好RocketMQ 用一份 CommitLog 换来了 Topic 无关性但要为索引读取付出成本JMQ 走的是第三条路保留顺序写的主日志同时通过强力索引和管理策略去吸收前面两者的压力。3. 关键链路对比写入、索引、刷盘与副本3.1 写入路径对比谁更接近“免打扰”我梳理了一个简化写入链路KafkaProducer 发送 - leader Partition 追加到当前 Segment 的页缓存 - follower 同步 - 返回 ACK。这里几乎没有额外的“分发”动作CPU 开销很低。RocketMQProducer 发送 - 写入 CommitLog 对应的 MappedFile内存映射文件- 异步构建 ConsumeQueue 索引 - 返回 ACK。多了一步异步分发但分发本身在内存里批量完成也不是瓶颈。JMQ同样是“主数据日志 队列索引”在写入时也会做索引构建但它更强调批量写、双缓冲以及把索引写入与数据写入解耦从而避免索引频率过高。在实际压测里如果只看“单条 100-200 字节的小消息吞吐”三者都不会差太多Kafka 和 RocketMQ 的量级都很猛。真正拉开差距的是两组极端场景一是消息体积很大比如单条 1MB 甚至更大二是队列数量特别多。第一组极端场景下大消息会把页缓存冲得很难受每次消费都可能触发真实磁盘读谁的索引更省、谁的冷读控制更好谁就有优势。第二组极端场景下Kafka 的每个 Partition 目录和 RocketMQ 的每个 ConsumeQueue 文件都有固定开销队列一旦上百万光是文件数量就够运维难受一阵子JMQ 那种“索引合并、多级共享”的思路就更有参考价值。3.2 索引与查找机制一条消息如何被快速找到我习惯把这套差异用表来记方便大家直接抄走维度KafkaRocketMQJMQ日志组织每 Partition 独立日志按 Segment 滚动所有 Topic 共用 CommitLog按文件大小滚动主数据日志统一写入文件按大小滚动索引文件.index 稀疏索引记录 offset 到物理位置ConsumeQueue 固定 20 字节/条记录 CommitLog offsetsizetag hashQueue Index / 内存索引具体实现随版本演进查找流程二分定位 Segment - 稀疏索引定位近似位置 - 顺序扫ConsumeQueue 读索引 - 拿到 CommitLog offset - 读消息正文类似 RocketMQ但多级存储下先判断数据在哪个层级强项分区内顺序读、零拷贝、吞吐极高Topic 数量与写入路径解耦适合多 Topic海量队列、多租户隔离、支持数据分层弱项海量 Partition 时文件与缓存开销上升读路径多一次索引跳转冷读风险略高复杂度高外部部署和维护成本不容易低估“读路径多一次索引跳转”听起来没什么但在页缓存不命中的情况下它意味着一次额外的随机磁盘 IO。生产上为了缓解这个弱点RocketMQ 通常会让 ConsumeQueue 长期驻留在页缓存里因为它的体积相比 CommitLog 要小很多20 字节一条索引而已这就相当于“用很小的常驻内存换掉了大日志文件的扫描”。JMQ 的思路更彻底它把索引和数据分级管理再配合共享页缓存池让“热 Topic 的索引”尽量留在内存“冷 Topic 的索引”可以接受磁盘 IO体现出明显的确定性管理风格。3.3 刷盘策略与副本谁把数据安全看得更重Kafka 的刷盘策略本质上是依赖操作系统默认情况下数据写入页缓存后由 OS 统一组织落盘Kafka 提供的 log.flush.interval.messages / log.flush.interval.ms 参数默认其实非常保守所以很多人把 Kafka 的可靠性完全寄托在副本上Partition 有多个副本ISR 机制保证至少有多少副本同步后确认。RocketMQ 则把刷盘开关直接交给使用者SYNC_FLUSH 表示每写一批就同步刷磁盘异步刷盘则是攒一批刷一次。同步刷盘的单条延迟上会明显上升但可靠性最强异步刷盘则在大量小消息场景下能把写放大压得很低。JMQ 在刷盘和副本上更强调“可配置 一致性协议”它不满足于仅靠 OS 兜底而是提供多种刷盘模式同时在副本组层面引入类似 Raft 的选主与同步机制让“写入成功”这件事在协议层得到确认。这里我特别想提醒任何消息队列都不能既要求“高性能异步刷盘”又要求“绝对不丢消息”。这三者的本质区别只是把可靠性放在哪个层级去兜底Kafka 主要靠副本RocketMQ 主要靠刷盘配置 副本JMQ 则在存储层引入了更强的协议闭环。4. 生产环境里的真实表现与常见坑4.1 “Kafka 消息延迟高”的排查思路先从存储找原因很多朋友用 Kafka 遇到“消息延迟高”第一反应去查消费者线程数和网络我一般在确认生产端没有堆积后会直接看磁盘层。最常见的原因有这么几类Segment 文件落在冷页缓存之外。长时间运行的大 Topic消息量大到把整个 OS 页缓存都撑破后新写入的段在消费时可能直接命中磁盘读导致拉取延迟飙升。单条消息太大比如 1MB 的消息。Kafka 默认单个消息上限 message.max.bytes1MB不是不能收而是大消息会让页缓存很快被冲掉消费端每次读到不同位置时大概率触发磁盘随机读。尤其在消费端较多、消费并发较高时这种随机读会放大磁盘 IO。分区数量少消费并发上不去。存储的分布没有均匀打散数据全压在两三个 Partition 上消费者再多也只能干等。同步刷盘与 acksall 叠加。如果配置把可靠性拉满写入延迟会被磁盘 fsync 时间和副本同步往返时间拖高这不是 Kafka“慢”而是你选择了可靠。磁盘 IOPS 本身不够。云盘如果被其他虚机抢占 IO也会出现延迟毛刺。此时用 iostat 能看到 %util 很高await 明显上涨。排查的时候可以结合 Kafka 的指标UnderReplicatedPartitions分区副本落后数、RequestQueueTime、LocalTime 等。如果 LocalTime 高说明服务端存储层写入变慢如果 ConsumerLag 高则要先看消费端。我自己的经验是把“磁盘 IO 延迟指标”接到监控大盘上比盯着消息数更早发现问题。4.2 Kafka 集群安装与目录观察从日志落地看存储设计当你真正把 Kafka 装起来才会对存储架构有直观感受。我复盘一下最小安装流程先下载 Kafka现在常用 3.x 版本KRaft 模式可以免 Zookeeper修改 config/server.properties 里的 log.dirs默认 /tmp/kraft-combined-logs。然后第一次启动前执行格式化命令生成集群元数据再启动 broker。如果还是用 Zookeeper 模式的旧版本需要先启动 ZK再启动 broker。起来之后往集群里创建一个 Topic假设叫 engineer-topic分区数设 3。此时到 log.dirs 对应的目录下你会看到类似 engineer-topic-0、engineer-topic-1、engineer-topic-2 这样的目录。每个目录下有一个 00000000000000000000.log 之类的段文件以及对应的 .index 和 .timeindex。这就是 Kafka 存储架构的“可视化结果”一个分区一个目录一堆 Segment 文件。安装时建议重点关注三个参数log.dirs可以写多个磁盘目录Kafka 会把 partition 分布到不同磁盘上。log.segment.bytes默认 1GB不要为了“管理方便”调太小否则段文件多索引和清理压力反而变大。log.retention.hours / log.retention.check.interval.ms决定日志保留多久关系到磁盘容量规划。注意在 KRaft 模式下目录结构和 ZK 模式略有差异但核心的“分区目录 Segment 文件”一直没变。网上很多安装教程还会让你调整 vm.swappiness、文件句柄数 ulimit这些和存储都有关系建议一次性配好。4.3 有没有官方 UI监控工具与存储指标盘点Kafka 官方一直没给 UI所以很多人刚上手时会问“Kafka 有没有 ui 界面”。生态里常用的方案有Kafka UI开源支持 Topic、Consumer Group、Seek、Offset Explorer原名 Kafka Tool图形化桌面工具、Kafka Eagle / EFAK开源社区常用支持告警和监控、CMAK原 Kafka Manager。这些工具的核心功能都在“查看 Topic 的分区、offset、消费组延迟”但真正想观察存储还是要看 Broker 自身指标比如未复制的分区数、日志目录大小、每秒写入字节数等。RocketMQ 官方生态相对完整有 RocketMQ Dashboard现在叫 rocketmq-dashboard可以在页面上看到 Broker 的存储目录、消息轨迹、消费组状态。JMQ 在京东内部则是一套微服务管控平台对外能看到的运维产品较少但从设计上讲它的控制台更偏向“多租户、容量规划、分层存储状态”这类系统级信息。监控存储架构时我最常用的三个指标磁盘写入吞吐直接判断是否接近磁盘上限。Page Cache 命中率能在操作系统的压力统计或监控系统里看到命中率低说明冷读严重。刷盘/fsync 耗时如果刷盘耗时持续走高多半是磁盘本身在过载。5. 高频面试题与选型参考看完可以直接拿去用5.1 消息队列存储面试题还原结合很多人搜“Kafka 面试题及答案”我把存储相关的常问问题整理成一张表常问问题参考答案要点Kafka 为什么写入快顺序追加日志、依赖 Page Cache、批量写、零拷贝 sendfile 读。Kafka 为什么不用主索引用稀疏索引减小索引体积通过二分定位 Segment 近似扫描。RocketMQ 为什么要 CommitLog 和 ConsumeQueue 两份CommitLog 保证统一顺序写、与 Topic 数量解耦ConsumeQueue 提供轻量索引供消费定位。RocketMQ 适合的场景电商交易、事务消息、延迟消息以及需要控制台和消息回溯的场景。JMQ 相比 RocketMQ 做了什么改进共享页缓存池、多级存储、海量队列支持、更强的一致性与多租户管理。消息会不会丢取决于刷盘配置、副本数量和 ACK 机制。堆积上亿条消息怎么办扩容消费者、增加分区需提前规划最终还要靠存储保留策略和消费进度处理。面试时如果能讲出一条“存储架构决定性能边界”的推导链会明显加分。比如Kafka 为了高吞吐做了分区独立日志天然不适合海量 Topic所以使用 Kafka 时要控制分区总数RocketMQ 为了多 Topic把主日志统一化所以它的写入曲线很平但消费路径多一次索引跳转JMQ 则在两者之间找平衡用更复杂的索引和分级存储换取更大的队列规模和隔离性。5.2 我的选型建议与容量规划经验真到了做架构选型这一步我觉得不要把问题变成“谁最强”而是变成“你的压痛点是什么它有没有踩到另一个的雷区”。如果你需要海量日志接入、离线数据管道、以流处理为主的高吞吐场景选 Kafka 很顺。在线核心交易、需要事务消息/延迟消息/消息轨迹、又希望控制台友好选 RocketMQ 更省心。多团队共用一套消息平台、Topic 数量巨大、又要区分不同业务的重要级别那么 JMQ 那种“共享页缓存池 多级存储 多租户”的设计很多理念值得借鉴甚至能帮你优化别的系统配置。容量规划上我的习惯是先算三个数消息平均大小、峰值的每秒写入条数、保留周期。用这三个数估算磁盘总容量再乘一个 1.5 至 2 的余量系数。比如单条平均 4KB峰值每秒写 2 万条保留 3 天那么每天约 4KB × 2 万 × 86400 ≈ 691GB三天大约 2TB预留后建议至少准备 4TB。如果其中夹杂大量 1MB 的大消息这个公式必须按最坏分位数重新算因为大消息对页缓存和随机读的影响不是线性的。个人在实际操作中有个很深的体会消息队列的存储架构不能只看启动时那几条 benchmark 数字。Kafka、RocketMQ、JMQ 三套方案看起来都在做“顺序写 索引”但一个 topic 数量、一个消息大小、一个刷盘模式波动就能把性能曲线拉开几个量级。我会建议团队在实际选型前把真实流量回放到测试集群用存储监控验证“页缓存、刷盘、索引命中”这三个底层指标而不是只看延迟曲线。如果今天这篇内容能帮助你少走一次弯路后续我再把刷盘、副本、消息堆积的详细排障拆成系列分享。
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进