消息队列选型:RabbitMQ、Kafka、RocketMQ、Pulsar 的真实吞吐对比和我踩过的坑

🔑 关键词:消息队列选型对比,Kafka,RocketMQ,RabbitMQ,Pulsar

📖 摘要:四个主流消息队列的吞吐量级、延时消息、顺序性、积压排查和运维成本对比,附 Kafka 与 RabbitMQ 积压排查的具体命令步骤,以及我自己的选型判断顺序。

先说我被问得最多的那个问题

图片

上周三晚上,一个做跨境电商的哥们给我发微信:我们 RabbitMQ 快扛不住了,是不是该换 Kafka。

我问了他三个问题。你们峰值多少 QPS?他说大概两万。消费者是什么写的?他说 Java 和 Python 各一半。你们团队有几个人能动这个集群?他停了一下,说,就他自己。

然后我说,你先别换。

他那点量,RabbitMQ 单机完全吃得下。他真正的问题出在消费者:Python 写的,一条消息要调三个外部 HTTP 接口,平均 200ms 才 ACK 一次。这意味着单个消费者实例一秒最多消费 5 条。两万 QPS 你得上四千个消费者实例才追得平。

这事跟中间件选型一毛钱关系都没有。但它特别典型——大多数人选 MQ,是在解决一个自己还没搞清楚的瓶颈。

四个主流 MQ 的关键数字堆一块

下面这个表是我们几个环境(8C16G、SSD、三节点)压出来和翻文档凑的,硬件和参数不一样差别很大,你当量级看就行。

RabbitMQ 3.13 Kafka 3.7 RocketMQ 5.x Pulsar 3.x
持久化吞吐(3 副本) 3~8 万/秒(classic queue),quorum queue 大概对半 30~80 万/秒(开批量) 8~15 万/秒 与 Kafka 同量级,瓶颈通常在 bookie 磁盘
单条延迟(不攒批) 亚毫秒到几毫秒 几毫秒到几十毫秒 毫秒级 毫秒级
顺序性 单队列 单分区 单队列 单分区
延时消息 要装 rabbitmq_delayed_message_exchange 插件,时间不准 原生不支持 原生支持 原生支持
事务消息 没有 有,但很慢,别用 有,这个是真能打 有
消息回溯 没有 有,改 offset 就行 有,按时间戳 有
依赖的外部组件 无,quorum queue 是内置 Raft ZK 或 KRaft 自研 NameServer ZooKeeper 存元数据 + BookKeeper 存数据

图片

我先说一句可能不太讨喜的:这张表里最没用的那一列是吞吐。因为你真到需要几十万 QPS 那天,瓶颈大概率不在 MQ 上,而在你的数据库和下游接口。

RabbitMQ:最该担心的不是性能,是镜像队列已经没了

还有不少人跑在 classic mirrored queue 上。这个功能在 3.12 就被标记废弃,3.13 直接移除(3.13.0 是 2024 年 2 月发的)。也就是说你现在做版本升级,镜像队列那套配置全部作废,得迁到 quorum queue。

我去年干过一次这个迁移,一个四十多个队列的订单系统。麻烦的地方在于 quorum queue 是 Raft,官方文档上写的就是至少 3 个节点才有意义,两个节点等于没有冗余。另外它不支持部分 classic queue 的特性,优先级队列、exclusive 消费这些,之前用了的就得改设计。

性能上的坑:quorum queue 每个队列的吞吐有上限,我们压出来大概是 classic queue 的六成左右(3 副本)。队列多、每个队列消息量还不小的时候,先扛不住的是 CPU。

再就是 prefetch。basic.qos 默认值是 0,意思是「有多少推多少」。低延迟小消息场景看着很美,但消费者一旦处理变慢,消息全堆在客户端内存里,JVM 分分钟 OOM。我们生产环境一般设 20~50,配手动 ack,改完之后消费慢的队列内存曲线肉眼可见地平了。

Kafka:分区和顺序性这两个坑,我见得太多了

图片

一个是分区数。

很多团队建 topic 的时候随手给 6 个分区,觉得六六大顺。某天流量涨了要扩到 12,命令就一行:

kafka-topics.sh --bootstrap-server broker:9092 --alter --topic order-event --partitions 12

分区数只能加不能减。更要命的是,如果你之前拿 orderId 当 key 发消息,扩完分区之后同一个 orderId 的哈希会落到不同分区,原来保证的顺序就没了。这个错误不报错、不告警,只会在某天对账的时候突然冒出来。

另一个是 max.poll.interval.ms,默认 300000,也就是 5 分钟。你要是一批 500 条(max.poll.records 默认就是 500)处理超过 5 分钟,消费者会被踢出消费组,触发 rebalance,这批消息重新投递,你以为处理了一次,其实处理了两次。日志里会刷一堆 heartbeat failed 或者 Member xxx is leaving。

调这个参数的时候,记得跟 max.poll.records 一起调,把单批处理时间压到 interval 的三分之一以内。别直接把它怼到 30 分钟,那只是把 rebalance 往后拖。

还有一个我觉得挺反直觉的事:大部分业务根本不需要全局顺序。为了顺序把 topic 压成一个分区,吞吐直接掉回 RabbitMQ 水平,然后转头骂 Kafka 慢,这个我见过不止两次。

RocketMQ:延时消息是它真能打的地方

图片

延时消息这个需求国内业务到处都是,订单 15 分钟未支付关单、确认收货 7 天自动确认。用 Kafka 做,你得自己写时间轮加定时扫表加重新投递,还得琢磨重启之后时间轮怎么恢复。

RocketMQ 是原生的。4.x 支持 18 个固定级别(1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h),你要的 15 分钟它会就近落到 20 分钟那一档,这个坑不少人上线之后才发现。5.0 之后支持任意时间精度的定时消息,好很多。

5.0 还有个 POP 消费模式,解决的是老 Push 模式的经典问题:一个消费组挂了几百上千个实例,每次有人上下线都 rebalance,期间整个组停摆几十秒。POP 改成 broker 侧按需拉取投递,不再做全组的队列重分配。

坑也有。RocketMQ 默认一个 Topic 是 4 个写队列、8 个读队列。压测的时候经常有人只开一个生产者线程,然后说性能不行——是你自己并发没上去。建 Topic 的时候按消费并发把队列数定好:

mqadmin updateTopic -n nameserver:9876 -c DefaultCluster -t order_topic -w 8 -r 16

w 是写队列数,r 是读队列数,写队列数直接决定了生产端的并发上限。

Pulsar:架构最漂亮,也最费人

图片

存算分离、broker 无状态、多租户、分层存储,听起来全是优点。broker 挂了随便加一台,扩容不用像 Kafka 那样挪分区数据。

代价是你要同时运维 BookKeeper 和 ZooKeeper。bookie 那边光是把 journal 和 ledger 分到不同盘、调 RocksDB 的写缓存读缓存、盯各种水位线参数,就够啃两天文档。而且 bookie 扩完要做 ledger 的 rebalance,不 rebalance 新 bookie 就是空的,老节点还是热点。

我的看法比较直接:Pulsar 适合的是有专职中间件团队、且确实要做大规模多租户的公司。中小团队上 Pulsar,本质上就是给自己新增了一个分布式存储的运维岗位。另外国内招 Pulsar 的人,比招 Kafka 的人难一个量级,这个成本选型的时候很少有人算进去。

真积压了,怎么查

这块我直接给步骤,不讲道理。

Kafka:

  1. kafka-consumer-groups.sh --bootstrap-server broker:9092 --describe --group order-consumer,看 LAG 和 CONSUMER-ID 那一列。如果某个分区的 CONSUMER-ID 是 -,说明这个分区压根没人在消费,先去看消费者是不是崩了。
  2. LAG 在涨但消费者都活着,看 --describe --members --verbose 里每个成员分到的分区,是不是倾斜了。
  3. 确认是消费慢,先量单条处理耗时。加消费者实例,但实例数超过分区数就没用了;要再加分区数,加之前先想清楚 key 哈希那个问题。

RabbitMQ:

图片

  1. rabbitmqctl list_queues name messages consumers memory,看哪个队列在涨、有没有消费者。
  2. 有消费者还积压,多半是 prefetch 太大或者消费逻辑本身慢,先把 prefetch 降下来。
  3. 内存水位一超(默认是物理内存的 40%),RabbitMQ 会直接 block 生产者。现象是连接不报错,但就是发不进去,很多人第一次遇到会懵很久。这时候别重启——重启之后消息从磁盘往内存里恢复,直接二次雪崩。

如果队列里堆了几千万条而且大部分已经过期,别在后端管理界面点删除,界面会卡死。用 rabbitmqctl purge_queue,或者写脚本一批两千条往外捞。我们上次处理两千万积压,是写脚本捞了两个多小时才捞干净的。

我的选择顺序

按顺序问自己四个问题:

  1. 峰值到底多少?低于 5 万 QPS,RabbitMQ 或 RocketMQ 单集群基本都能扛,别为了性能上 Kafka。
  2. 业务里有没有延时消息、事务消息、消息回溯?有就直接 RocketMQ,能省你至少两个月开发量。
  3. 有没有日志、埋点、流式场景?量到百万级 QPS,Kafka。
  4. 团队里有几个人真的懂这些中间件的内部机制?一个都没有的话,选生态最大、出问题搜得到答案的那个,通常就是 Kafka 和 RocketMQ。

最后还是那句:中间件选型,本质上是在选你未来两年半夜要爬起来处理的问题类型。性能可以从容优化,运维成本是每天都在发生的。别被几张 benchmark 图带偏。

以上都是我自己踩坑踩出来的,不一定对,有不同意见欢迎来喷。

🏷️ 标签: