那个周五晚上
2023 年 11 月的一个周五,我在公司楼下刚端上一碗麻辣烫,手机开始震。生产集群的订单消息 lag 冲到 470 万,Kafka 的消费曲线像心电图。等我打车回到工位,发现根因是某个消费者实例卡在了一次 poll 之后的批量入库上——下游 MySQL 那会儿有个慢查询,一批 500 条消息处理了 4 分钟,加上网络抖动,刚好越过 max.poll.interval.ms 默认的 300000 毫秒。然后就是经典的 rebalance 风暴,整个消费组停摆 3 分 12 秒。
那次之后我做了一件事:把公司 12 个业务线的消息量拉了个表。结果挺意外,日均消息量过 5000 万的只有 3 个,剩下 9 个业务线平均在 200 万到 800 万之间,峰值 QPS 没超过 3000。
也就是说,我们养着 6 台 16C32G 的 Kafka broker,大部分时间 CPU 利用率不到 8%。
压测到底怎么压的
我拉了测试环境的 3 台 8C16G(阿里云 ecs.g7.2xlarge 那档,ESSD PL1 云盘,千兆内网),分别装了 Kafka 3.6(KRaft 模式,3 台混布 controller + broker)、RocketMQ 5.1(2 台 NameServer 塞在同一批机器上,broker 用 2 主 2 从,实际上 3 台机器跑不满副本数,最后只跑了 2 主 1 从)、Pulsar 3.0(3 broker + 3 bookie + 1 个 ZK 挤在一起,说实话这个配置对 Pulsar 不太公平,但预算就这么多)。
消息体统一 1KB,生产者 128 线程,acks=all,同步刷盘(Kafka 调小 flush.messages,RocketMQ 用 flushDiskType=SYNC_FLUSH,Pulsar 开 journalSyncData=true)。压了整整两周,每天跑 4 小时,中间因为云盘 IOPS 抖动重跑过三次。
数字放这儿
写入吞吐(单机,1KB 消息,acks=all):
- Kafka 3.6:18.4 万 msg/s,P99 45ms
- RocketMQ 5.1:11.2 万 msg/s,P99 28ms
- Pulsar 3.0:7.1 万 msg/s,P99 62ms
注意 RocketMQ 这里是同步刷盘。改成异步刷盘(ASYNC_FLUSH)能到 26 万左右,但断电丢消息的风险得自己评估。另外 RocketMQ 的 P99 最低,我猜和它的 CommitLog 顺序写 + 零拷贝(mmap 到 sendfile)那条链路相对短有关。
Kafka 的 P99 反而最高,这和 KRaft 模式下元数据同步有关,另外我们只给了 3GB page cache 缓冲,数据量上去之后刷盘抖动明显。Pulsar 慢在 BookKeeper 的三副本写入,每条消息要等 2 个 bookie 确认(ackQuorum=2),网络往返多了。
延迟消息,这才是分水岭
很多人选 MQ 只看吞吐,我觉得延迟消息的支持程度更能说明问题。
RocketMQ 5.0 之后支持任意秒级的定时消息,精度 1 秒。4.x 时代是 18 个固定级别(1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h),要发一个 15 分钟的延迟消息得手动映射到 20m,很别扭。我们测下来 5.1 的定时消息在 10 万条并发下投递误差在 800ms 以内。
RabbitMQ 想延迟只能靠 TTL + 死信队列。这个方案有个经典坑:队列是先进先出的,如果队头那条消息 TTL 是 1 小时,后面 TTL 是 5 秒的消息全得堵着。我们之前用这个方案跑 10 万条延迟消息,Erlang 进程内存涨了 800 多 MB,最后是给每个延迟级别建一个独立队列绕过去的。
Kafka 没有原生延迟消息。TimingWheel 那套得自己写在应用层,还得考虑重启之后怎么恢复,我见过一个团队用 Redis ZSet 存延迟索引,结果 Redis 挂了消息全丢。
所以如果你的业务里有「下单 30 分钟未支付自动取消」这类需求,RocketMQ 能省掉一整个定时任务服务。
事务消息,别被名词骗了
Kafka 有事务,RocketMQ 也有事务消息,但这俩说的完全不是一回事。
Kafka 的事务(transactional.id + beginTransaction / commitTransaction)解决的是流处理里的 exactly-once 语义——消费-处理-生产这个链路不重复。它不管本地数据库事务,也不管消息发出去之后业务有没有执行成功。
RocketMQ 的事务消息是 half message + 回查机制。发 half 消息,执行本地事务,成功了 commit,失败了 rollback;如果 broker 长时间没收到二次确认,会反查生产者「这笔本地事务到底成没成」。这才是用来解决「数据库扣款成功但消息没发出去」的场景。默认回查次数 15 次,间隔从 10 秒开始退避。
我们有个支付回调的业务,之前用 Kafka + 本地消息表 + 定时补偿,代码大概 600 多行。用 RocketMQ 的事务消息重写之后,200 行不到。这是实打实的收益。
顺序消息的坑比你想的深
Kafka 保证分区内有序,RocketMQ 保证队列内有序,听起来一样,但扩容的时候差别很大。
Kafka 的 key 到 partition 的映射是 hash(key) % numPartitions。你从 8 个分区扩到 16 个,同一个订单号的 key 可能就落到不同分区了,顺序直接断掉。所以 Kafka 想扩容分区,基本得停写。
RocketMQ 的 MessageQueueSelector 可以自己实现。扩容的时候把新队列标成只读、等旧队列消费完再切,比 Kafka 好操作一些。
还有一个很多人不知道的点:Kafka 的 rebalance 期间,分区会在消费者之间迁移,迁移过程中同一分区的消费是断的,顺序自然也就没了。用 CooperativeStickyAssignor 能减少 rebalance 范围,但解决不了根本问题。
运维成本才是真成本
Kafka 3.6 上了 KRaft 之后不用 ZooKeeper 了,确实是好事。但 controller 选举出问题的时候,日志里那一堆 KRaft 前缀的堆栈,看了头疼。有一次我们三个 controller 有两个 OOM,集群直接只读,恢复花了 40 分钟。
RocketMQ 的 mqadmin 命令行工具很好用,mqadmin consumerProgress 一眼能看到每个消费组的 lag。但 Dashboard 那个界面确实丑,而且 5.x 之后 Proxy 模式和 NameServer 模式混着用,文档有点乱。
Pulsar 的存算分离是个好东西,扩容 broker 几分钟搞定。但 BookKeeper 得单独运维,bookie 的 journal 盘写满会造成 ledger 写入失败,而且是静默失败——我们那次磁盘到 95% 才发现,已经丢了 3 个 ledger 的副本。另外 Pulsar 的 broker.conf 和 bookkeeper.conf 两套配置,调优门槛比前两个高一个量级。
我们最后没换
评估了一圈,结论有点扫兴:我们最后留在 Kafka 上,只是把 max.poll.records 从 500 降到 200,加了 CooperativeStickyAssignor,给消费逻辑加了 30 秒的超时保护,然后把 9 个业务的 broker 拆到不同集群。
为什么不换 RocketMQ?因为 12 个业务线一共 47 个生产者,分布在 9 个团队,最近的一个版本迭代还在下个月。迁移的隐性成本——协调、回归测试、双跑对账——加起来比省下的 3 台机器贵得多。
这就是我最大的一个体会:MQ 选型从来不是技术问题。先数清楚谁能半夜起来看 broker 日志,谁能在 Grafana 上配一套 lag 告警,这比吞吐量重要得多。
如果你的团队少于 5 个人,日均消息量不到 1000 万,直接用云上的托管 MQ,把运维那部分钱花在业务上更值。