Kafka、RocketMQ、Pulsar 怎么选?我把三个生产集群的存储目录翻了一遍

🔑 关键词:消息中间件,Kafka,RocketMQ,Pulsar,存储结构选型

📖 摘要:不从吞吐量对比表出发,而是从三个消息中间件底层的文件布局、索引结构和扩容方式倒推选型逻辑,含真实参数与踩坑记录。

上周三凌晨两点,我盯着 Grafana 上一个 200 分区的 topic 发愣。Kafka 3.6 的集群,32 个 broker,某个消费组的 lag 从 4 万跳到 180 万只用了 7 分钟。查到最后不是消费逻辑慢,是 max.poll.interval.ms(默认 300000,也就是 5 分钟)没动过,消费者在 rebalance 里来回横跳,日志里刷满 Attempt to heartbeat failed since group is rebalancing。把参数调到 900000 重启,lag 在 40 分钟内清完。

图片

这件事之后我把三个中间件的 store 目录挨个翻了一遍,因为那些「Kafka 吞吐高、RocketMQ 适合电商、Pulsar 架构先进」的对比表,没有一个能解释我凌晨两点遇到的那个问题。

RocketMQ 的 CommitLog 和那 20 字节

打开 RocketMQ 的 store 目录,最显眼的是几个 1GB 的 CommitLog 文件和一堆 ConsumeQueue。ConsumeQueue 每条固定 20 字节:8 字节 CommitLog offset、4 字节消息长度、8 字节 tag hashcode。注意 tag 存的是 hash,不是原文,所以两个 tag 哈希撞车时,broker 会把消息捞出来再让客户端过滤一遍,这也是为什么 tag 写得太随意会影响性能。

图片

这套设计是顺序写 CommitLog、随机读 ConsumeQueue,但 ConsumeQueue 本身很小且连续,读取时又被 OS 的 page cache 拉回顺序 IO。2012 年阿里内部就是这么干的,到 5.x 基本没动过大结构,说明它扛住了双十一那种量级。代价是 CommitLog 是全局共享的,一个 broker 上所有 topic 抢同一块写入带宽,单个 topic 的写入性能天花板不明显,但故障时影响面是整机级别的。

Kafka 的稀疏索引和分区的物理含义

Kafka 的 partition 就是一个目录,topic-0/ 下面是 00000000000000000000.log.index.timeindex 三件套。.index 默认每写 4096 字节(index.interval.bytes=4096)才加一个索引项,是稀疏索引,所以按 offset 查一条消息要先二分索引再顺着 log 扫。segment 默认滚到 1GB(segment.bytes=1073741824)切新文件,log.retention.hours 默认 168 小时,也就是 7 天。

图片

分区的物理含义决定了它的运维天花板:一个分区只能被消费组内一个消费者线程消费,所以你的并行度上限就是分区数。有人为了吞吐把分区开到 200,然后又来问为什么全局有序没了——分区内有序、分区间无序,这是写死在结构里的事,不是配置能调的。

Pulsar 把存储彻底推出去了

Pulsar 的 broker 自己不存消息,消息落在 BookKeeper 的 bookie 上。你打开 bookie 目录看到的是 journal 和 ledger 两套东西,journal 做写前日志保证恢复,ledger 才是最终存储。单个 ledger 的滚动有两个触发条件:写满阈值或者超过 4 小时(managedLedgerMaxLedgerRolloverTimeMinutes=240)。

图片

这是它和前面两个最根本的区别。broker 无状态意味着扩容就是加机器,不用搬数据;bookie 无状态意味着存储层可以单独扩。代价是组件从一套变成三套——ZooKeeper、BookKeeper、Broker 都得维护,排查链路也长。我们有一次 topic 卡住,最后定位到某个 bookie 的 journal 盘 IO 打满,根因是部署时 journalDirectoryledgerDirectories 没分到不同物理盘,这条规范当时没严格执行。

扩容这件事,比吞吐量重要得多

图片

Kafka 加 broker 之后要手动跑 kafka-reassign-partitions.sh 搬分区,PB 级集群搬一次是几天起步,搬的时候 IO 和网络一起飙。我带过一个 6 节点、1200 分区的集群重平衡,跑了两天半。RocketMQ 相对轻松,队列按 broker 分片,新 broker 上线后扩 topic 队列数即可,老数据不动——但前提是你一开始 readQueueNums 就规划够了,队列数加不了太多,顺序消息的语义会断。Pulsar 在这点上最省心,加机器即插即用。

所以我的实际判断顺序是:先看你的运维团队能扛几套 ZK,再看故障时你希望影响一个分区、一个 broker 还是整个存储层,最后才看吞吐。

队列模型和流模型,以及 POP 模式

图片

RocketMQ 5.0 推了 POP 消费模式,本质是打破「一个队列同一时刻只能被组内一个消费者持有」这个约束。老的 PushConsumer 下,20 个消费者对着 8 个队列,剩下 12 个就在空转。POP 改成 broker 侧按消息维度投递,popInvisibleTime 默认 60000ms 做超时重投。Kafka 那边的对应概念是 max.poll.records(默认 500)加上分区数,本质还是分区独占。

说句可能不中听的:大部分团队的问题不是选错了中间件,是把中间件当数据库用了。我见过在 RocketMQ 里存 30 天消息、然后按 msgId 查历史订单的,也见过拿 Kafka 当配置中心推配置的。消息中间件的 offset/ack 语义是为「处理完就往前走」设计的,你一旦要「回头翻」,成本就以数量级往上跳。真要能翻,得上 Pulsar 的分层存储 offload 到对象存储,或者干脆别用消息队列。

我的选型口径不一定对:已经在 Hadoop/Spark 生态里、要做流批一体的用 Kafka;订单、事务消息、定时消息这类业务语义重的用 RocketMQ(5.0 之后定时消息支持秒级任意精度,之前的版本只有 18 个固定延迟级别);多租户、跨地域复制、不想管分区的用 Pulsar。其他情况,先算算你有几个人能值 ZK 的班。

🏷️ 标签: