Chapter 04

对标 Kafka:何时选谁

前三章建立了 RabbitMQ 的模型、可靠性与内部机制——这一章把"智能 broker、ack 后删除"这条主线放到 Kafka 的"不可变日志、offset 可重放"旁边,逐条对比并给出选型判据。读者已经熟悉 Kafka,所以每一条对比都落到机制层:差异从哪来,而不是只贴一个标签。

本章你将建立的 schema

  • 架构分叉:RabbitMQ 是被路由的临时信箱,Kafka 是可重放的持久账本,一切差异都从这里长出来
  • 投递模型:RabbitMQ 用 prefetch 限流的 push,Kafka 用按 offset long-poll 的 pull——谁定速率,谁会被压垮
  • 顺序与重放:Kafka 分区内全序且可任意倒带;RabbitMQ 的顺序在 competing consumers 与 requeue 下天然易碎
  • 吞吐与延迟的真实数字(Confluent 2024 基准):Kafka 赢的是吞吐不是单条延迟,低负载下 RabbitMQ 反而更快
  • 选型判据:RabbitMQ、Kafka、RabbitMQ Streams 三者各自的归属场景,以及它们正从两端向中间收敛

4.1架构对照:智能 broker vs 不可变日志

RabbitMQ 把状态、路由、ack 记账都压在 broker 里,消息 ack 即删;Kafka 的 broker 只追加不可变字节,不路由、不删除,消费进度由消费者自己用 offset 持有。

这条分叉在 01 章 §1.4 已被立为全教程的轴心。这里把它摆到 Kafka 旁边逐项展开:同一个"消息从生产者到消费者"的过程,两套系统把哪些职责放在了哪一侧。

表 4.1 · 职责落点对照
职责RabbitMQ(智能 broker)Kafka(不可变日志)
路由决策broker 侧:exchange + binding 把消息分发到 0~N 个 queue无 broker 侧路由:producer 按 key 哈希选 partition,仅此而已
消息存储broker 持有消息状态(Ready / Unacked),随消费推进而变broker 存不可变字节序列;写入后只读、不改
消费进度broker 替消费者记账:谁 ack 了、谁还欠着消费者自己持有 offset;broker 不知道谁读到哪
消费后ack 即删除,不留记录消息留存,按 time/size 策略到期才删,可被重读
重放能力classic/quorum 队列:ack 后即消失,无法倒带任意消费者可把 offset 重置到任意位置重读

底层机制(比文档深一层):差异的根在"谁持有消费游标"。RabbitMQ 的 broker 为每个 queue 维护一份消息集合,每次投递分配一个 delivery tag 并把消息从 Ready 移到 Unacked;收到 ack 就从存储里抹掉这条。游标在 broker 内部,且是单调销毁的——读过的位置不再存在。Kafka 反过来:broker 把每个 partition 写成一个仅追加的 segment 文件序列,写入即定型,broker 对"谁消费到哪"一无所知;offset 是消费者自己提交到 __consumer_offsets 这个内部 topic 里的一个数字。因为消息不随消费而变,多个消费者读同一份字节互不干扰,任何一个都能把自己的 offset 倒回去重读。一句话:RabbitMQ 的状态机活在 broker 里,Kafka 的状态机活在消费者里。

RabbitMQ · 被路由的临时信箱 Kafka · 可重放的持久账本 Producer exchange 路由 queue broker 持状态 Consumer binding push ack→删除 Producer partition(仅追加 · 不可变) 0 1 2 3 4 … append Consumer offset=3 可倒带重读
图 4.1左侧 RabbitMQ:消息经 exchange 路由进 queue,push 给 consumer,ack 后从 broker 删除。右侧 Kafka:producer append 到不可变 partition,consumer 拿自己的 offset 指针前移。注意:左侧消息消费后不复存在,右侧消息留存、offset 可任意倒回——这是后续每一条差异的源头。
洞察 · 反直觉点

对一个 Kafka 老手,最容易低估的不是"RabbitMQ 删消息",而是 Kafka 的 broker 完全不做路由与过滤。Kafka broker 收到的是一串字节,它只负责按 partition 追加、复制、按策略过期;任何"这条消息该给谁、要不要过滤掉"的逻辑,要么落在 producer 的分区键里,要么落在消费端 / Kafka Streams 里。RabbitMQ 的 exchange 把这套逻辑放进了 broker——这正是它路由能力的来源,也是它吞吐上限的来源。两端的取舍,4.4 用数字说话。

想一想

一个消费者把一批消息处理完并 ack 了,事后发现处理逻辑有 bug,需要拿原始消息重跑一遍。RabbitMQ(classic/quorum 队列)和 Kafka 各能怎么办?

展开答案(先停 10 秒再点)

Kafka:把消费者的 offset 重置到那批消息的起始位置,重新消费一遍即可——消息还在日志里。RabbitMQ 的 classic/quorum 队列:消息在 ack 那一刻已被 broker 删除,没有任何办法从队列里再取出来,只能寄望于上游能重新生产,或事先把消息落到别处。这就是"ack 后删除"在运维上最痛的一面,也是 RabbitMQ Streams(4.5 详谈)想补的洞。

4.2投递与消费:push/prefetch vs pull/offset

RabbitMQ 把消息 push 给消费者,用每消费者的 prefetch(最多多少条未 ack 在途)限流;Kafka 让消费者 pull——按 offset 长轮询拉取,速率由消费者自己定。

表 4.2 · 投递模型对照
维度RabbitMQ(push + prefetch)Kafka(pull + long-poll)
谁发起投递broker 主动把消息推给消费者消费者主动按 offset 发 fetch 请求拉
速率由谁定broker;消费者只能用 prefetch 设上限消费者自己;想多快拉多快
慢消费者后果未 ack 堆到 prefetch 上限即停推,易被压垮自己少拉即可,不阻塞别人,落后了之后追上
单条延迟低——消息一就绪就推出去略高——要等下一次 fetch,靠批量摊薄
吞吐驱动逐条投递 + 逐条 ack 记账批量 fetch,一次拉一大段

底层机制(比文档深一层):RabbitMQ 的 push 是 broker 通过 basic.deliver 主动把消息帧推到消费者的 channel 上。流控的唯一旋钮是 prefetch(basic.qos 的 prefetch count):它限定一个消费者最多有多少条"已推送但未 ack"的消息在途。prefetch 设 1,broker 推一条、等 ack、再推下一条,延迟最低但吞吐受单条往返限制;prefetch 设大,在途消息多、吞吐高,但若消费者处理慢,这些消息全卡在它手里(Unacked),其它空闲消费者也分不到。Kafka 的 pull 是消费者发 fetch 请求带上 offset,broker 用 long-poll——有数据立即返回,没数据则挂起到 fetch.max.wait.ms 再返回空,避免空轮询打爆 CPU。消费者一次能拉 max.poll.records 条,处理完再拉下一批,速率完全自控;一个慢消费者只是自己 offset 前进得慢,partition 里的数据不动,它随时能追上。push 把"会不会压垮消费者"的风险交给 broker 的 prefetch 旋钮去兜;pull 把这个风险从系统里消除了——消费者永远不会被喂太多。

洞察 · 把 backpressure 想清楚

backpressure(背压,下游处理不过来时让上游慢下来)在两套系统里是两种东西。Kafka 的 pull 让背压自然成立:消费者不发 fetch,数据就在 partition 里待着,天生不会过载。RabbitMQ 的 push 没有这种天然背压,它用 prefetch 来人为制造一个上限——这是个必须调对的旋钮:太小则吞吐被单条往返拖死,太大则消息全堆在一个慢消费者手里、其它消费者饿着。prefetch 的取舍在 03 章 §3.2 的单队列处理模型里有更细的展开。

想一想

一个 RabbitMQ 队列接了 3 个消费者,prefetch 设成 1000。其中一个消费者卡住了(不崩、但处理极慢,迟迟不 ack)。队列里堆积的消息会怎么分配?

展开答案(先停 10 秒再点)

broker 会先把多达 1000 条消息推给那个卡住的消费者(填满它的 prefetch 窗口),这 1000 条就锁在它手里成了 Unacked,其它两个健康消费者干瞪眼。结果是吞吐塌方却没有任何报错。把 prefetch 调小(如 10~50),broker 就不会一次喂太多给慢消费者,消息更均匀地流向健康消费者。换成 Kafka 不存在这个问题——每个消费者按自己的速率 pull,慢的那个只是自己落后。这正是 push 模型必须调 prefetch 的原因。

4.3顺序与重放:RabbitMQ 的顺序为何易碎

Kafka 在一个 partition 内提供全序,同 key 进同 partition 即有序;RabbitMQ 按队列到达顺序投递,但 competing consumers 让完成顺序错位,requeue 还会把消息重新插回队列(常在头部),"FIFO 队列"不等于"FIFO 处理"。

消费并行与重放也在 01 章 §1.4 埋下了伏笔。这里把它和 Kafka 的 consumer group 摆在一起看。

表 4.3 · 并行、顺序、重放对照
维度RabbitMQKafka
并行模型competing consumers:一个 queue 接多个消费者,broker 把消息一人一条派出去consumer group:每个 partition 最多 1 个消费者
加并行直接加消费者,零协调,立刻分担受 partition 数封顶;加消费者超过 partition 数则空转
重平衡无——加减消费者不触发协调消费者加入/离开触发 rebalance,期间短暂停顿
顺序保证队列到达有序,但 competing consumers + requeue 下处理顺序易碎partition 内严格全序(跨 partition 无序)
重放classic/quorum:ack 即删,无法重放任意 offset 自由重读

底层机制(比文档深一层):Kafka 的顺序来自"同 key 哈希到同一 partition + 一个 partition 在组内只归一个消费者"这两条的合力——同一实体的消息走同一条单线日志,被同一个消费者按 offset 顺序读,所以有序。RabbitMQ 的队列投递是按到达顺序的,但顺序在两处断裂:其一,competing consumers——broker 把消息一人一条地派给多个消费者,谁先处理完不取决于谁先拿到,msg-1 给了消费者 A、msg-2 给了消费者 B,B 更快,于是 msg-2 先完成,完成顺序和到达顺序脱钩。其二,requeue——一条被 nack 或因消费者断连而退回的消息,会被重新插入队列,常常插在队首(取决于版本与配置),于是它越过了本来排在它后面的消息,顺序当场被打乱(这个 requeue 回插的细节,01 章 §1.4 的状态机图里画过那条回环)。结论很反直觉:RabbitMQ 的 queue 是 FIFO 的数据结构,但一旦有并行消费或重投,得到的不是 FIFO 的处理。

要在 RabbitMQ 上复刻 Kafka 那种"按 key 有序又能并行",得动用 consistent-hashing exchange(一致性哈希交换机,按消息 key 哈希到固定的某个 queue),把相关消息全路由到同一个 queue,再给该 queue 配单一活跃消费者(single active consumer,同一时刻只有一个消费者在消费这个队列)。等于用"一 key 一队列、一队列一消费者"手工搭出 Kafka 的 partition 语义——能做到,但要自己拼装,而 Kafka 里这是默认行为。

洞察 · 容易高估的对称性

Kafka 的并行上限被 partition 数死死封顶:partition 建了 12 个,最多 12 个消费者并行,第 13 个只能空转等接替;想加并行得先增 partition(且增了不能减)。RabbitMQ 没有这个上限——队列接多少消费者都行,加一个立刻分担、零重平衡。但这份"自由"的代价正是上面的顺序易碎。两者在这里不是对称的优劣,而是把"顺序"和"弹性扩容"放在了天平的两端:Kafka 锁顺序、限弹性;RabbitMQ 给弹性、弃顺序。

想一想

某业务要求"同一个用户的事件必须按发生顺序处理,不同用户之间可以并行"。在 RabbitMQ 上要怎么搭?为什么单纯"一个队列 + 多个消费者"不行?

展开答案(先停 10 秒再点)

单纯"一队列多消费者"不行:competing consumers 会把同一用户的相邻事件派给不同消费者并行处理,完成顺序错乱。正解是 consistent-hashing exchange 按 userId 哈希,把同一用户的事件恒定路由到同一个 queue,再给每个这样的 queue 配 single active consumer。不同用户落在不同 queue、由不同消费者并行处理,同一用户落在同一 queue、被一个消费者顺序处理。这恰好是 Kafka 用"key→partition + 一 partition 一消费者"免费给你的东西。

4.4吞吐与延迟:真实数字与原因

Kafka 在吞吐上压倒性领先,但赢的是吞吐不是单条延迟——在低负载下 RabbitMQ 的端到端延迟反而更低。

Confluent 2024 的基准测试(i3en.2xlarge 实例、3 个 broker、1 KB 消息、3 副本、4 producer/4 consumer)给出一组常被引用的数字:

表 4.4 · Confluent 2024 基准(i3en.2xlarge · 3 broker · 1KB · 3× 副本)
系统峰值吞吐延迟特征
Kafka~605 MB/s200K msg/s 时 p99 约 5 ms;吞吐压倒性领先
Pulsar~305 MB/s200K msg/s 时 p99 约 25 ms
RabbitMQ~38 MB/s~30K msg/s 以上 CPU 瓶颈、p99 急剧抬升;但 ~30 MB/s 低负载下 p99 约 1 ms,低于 Kafka

底层机制(比文档深一层):Kafka 赢吞吐有四个叠加的原因。① 顺序磁盘写——partition 是仅追加文件,每个字节只在一条优化了近十年的代码路径上落盘一次,顺序写对机械盘和 SSD 都远快于随机写。② OS page cache——Kafka 不自建缓存层,直接复用 Linux 页缓存,读多写少时数据热在内核内存里。③ zero-copy(零拷贝)——用 sendfile 把文件数据从页缓存直接送到网卡,不经过用户态的来回拷贝。④ producer 批量——基准里 producer 攒批至多 1 MB、最多等 10 ms(linger),把多条消息合成一次写、一次 fsync,把固定开销摊薄。RabbitMQ 的上限来自相反的设计:它要为每一条消息做路由匹配、维护 Ready/Unacked 状态、记 ack 账,逐条记账的固定开销无法摊薄;更关键的是每个队列由单个 Erlang 进程处理(single-process-per-queue bottleneck,03 章 §3.2 拆过),一个队列的吞吐被一个进程的单核能力封顶,30K msg/s 以上就撞到 CPU 墙。

峰值吞吐(MB/s · 越长越快) Kafka ~605 MB/s RabbitMQ ~38 MB/s 同硬件、1KB、3 副本(Confluent 2024)
图 4.2同一套硬件上 Kafka 峰值吞吐约为 RabbitMQ 的 16 倍。注意:这张图只说吞吐——在 ~30 MB/s 的低负载下 RabbitMQ 的 p99 延迟(约 1 ms)反而低于 Kafka,吞吐与单条延迟是两件事,别用前者的差距去推断后者。
洞察 · 诚实地表述这个差距

"Kafka 比 RabbitMQ 快 16 倍"这句话只在吞吐这一维成立。换到单条消息的端到端延迟,结论会翻转:在远低于 Kafka 的吞吐区间(约 10 MB/s 以下),RabbitMQ 能做到亚毫秒级 p99,比 Kafka 更低——因为 Kafka 的批量机制本身要等 linger、攒批,给单条消息加了固有延迟。所以一个低流量、要求每条消息尽快送达的请求-响应场景,RabbitMQ 反而是更快的那个。把"Kafka 快"理解成"Kafka 吞吐高",而不是"Kafka 每条都更快"。

想一想

一个内部 RPC 系统,QPS 几千、每条请求都要尽快拿到响应,延迟敏感但吞吐谈不上高。光看"Kafka 605 vs RabbitMQ 38"就选 Kafka,会踩到什么?

展开答案(先停 10 秒再点)

会选错维度。这个场景的约束是低延迟不是高吞吐,几千 QPS 远在 RabbitMQ 的舒适区内(30K msg/s 才撞墙)。在这个负载下 RabbitMQ 的 p99 更低,而 Kafka 的批量/linger 反而给每条请求加了延迟。加上 RPC 需要请求-响应、临时回复队列这类 RabbitMQ 原生擅长的模式。结论:延迟敏感 + 中等吞吐 + 请求-响应 → RabbitMQ。605 这个数字在这里根本不相关。

4.5选型决策:RabbitMQ、Kafka、还是 Streams

复杂路由 / 每消息工作流 / RPC / 低延迟中等吞吐 → RabbitMQ;高吞吐事件流 / 重放 / 事件溯源 / 多独立消费者读同一份流 → Kafka;想要重放与扇出但留在 RabbitMQ 生态且吞吐要求不极端 → RabbitMQ Streams。

表 4.5 · 场景 → 选型
场景选谁为什么
复杂路由 / 按条件扇出到每个应用各自的队列RabbitMQexchange(fanout/direct/topic/headers/一致性哈希)做 broker 侧过滤;Kafka 无 broker 侧路由(01 §1.3)
每消息任务/工作队列、priority、TTL、DLX 工作流RabbitMQ这些是 RabbitMQ 的原生语义;Kafka 没有逐消息优先级/过期/死信
RPC / 请求-响应、临时回复队列RabbitMQ低延迟 push + 灵活路由天然契合;Kafka 的日志模型不为此设计
低延迟、中等吞吐RabbitMQ低负载下 p99 更低(4.4);没撞到 30K msg/s 的墙
高吞吐事件流、海量写入Kafka顺序写 + page cache + zero-copy + 批量,吞吐高一个量级(4.4)
重放 / 重新处理、事件溯源Kafka消息留存、offset 可任意倒带;classic/quorum 队列 ack 即删
多个独立消费者各自读同一份流、流处理Kafka不可变日志天然支持多读者互不干扰 + Kafka Streams
想要重放/扇出给多读者,但要留在 RabbitMQ 生态、吞吐不极端RabbitMQ Streamsappend-only 可重放日志、非破坏性消费,免迁 Kafka

底层机制(比文档深一层):第三个选项 RabbitMQ Streams 是 RabbitMQ 对 Kafka 的正面回应(GA 自 3.9,super-streams 分区在 4.x 稳定)。它是一个仅追加、可复制的日志,消费是非破坏性的——消费者按 x-stream-offset 附着到日志的某个位置(first / last / 某个 offset / 某个 timestamp),读过不删,可重放,多个消费者读同一份。4.2 进一步加了服务端 AMQP 1.0 的 SQL 风格 filter expression(服务端按表达式过滤,只下发命中的消息)。这把 Kafka 的几项核心能力搬进了 RabbitMQ:留存、重放、扇出给多读者。但它不等于 Kafka:一个 stream 仍然没有 classic/quorum 队列的 TTL / DLX / priority / 逐消息持久化语义;RabbitMQ 官方文档也直言——没有任何持久化队列类型能在吞吐上匹敌一个日志系统。所以 Streams 是从 RabbitMQ 这一端向 Kafka 收敛,而非完全取代;与此同时 Kafka 在 2024 年补上了 tiered storage(分层存储,热数据在本地盘、冷数据下沉到对象存储,便宜地长期留存并重放)。两个系统正从相反的出身互相靠拢:一个加日志语义,一个加廉价长留存。

选型起点 需要复杂路由 / RPC / 每消息工作流? 是 RabbitMQ 否 需要重放 / 事件溯源 / 超高吞吐? 否 RabbitMQ 低延迟·中吞吐 是 想留在 RabbitMQ 生态、 吞吐不极端? 是 RabbitMQ Streams 否 Kafka
图 4.3三问决策树:先问路由/RPC/逐消息工作流(是→RabbitMQ),再问重放/溯源/超高吞吐(否→RabbitMQ 低延迟中吞吐场景),最后在"要重放"里分流——留在生态且吞吐不极端→Streams,否则→Kafka。注意:超高吞吐这一支,Streams 仍不及 Kafka,吞吐压满时走右路。
边界 · 别把 Streams 当 Kafka 平替

RabbitMQ Streams 补上了留存与重放,但它不继承 classic/quorum 队列的 TTL、DLX、priority、逐消息持久化语义,吞吐也不及真正的日志系统;4.2 的 SQL filter 还只走 AMQP 1.0,原生 Stream 协议与传统 0.9.1 客户端都用不上。需要 Kafka 级吞吐、成熟流处理生态(Kafka Streams / ksqlDB / Connect)或分层存储长留存时,仍然是 Kafka。Streams 的定位是"不想为了重放就整体迁去 Kafka"的折中。

§本章 self-check

先合上教程,把答案写在纸上或编辑器里。写完再点开对照——直接点开等于把这一节当再读一遍。

  1. Kafka 的 broker 对"一条消息该投给谁、要不要过滤"做了什么?这和 RabbitMQ 的 exchange 有什么本质区别?
  2. "Kafka 比 RabbitMQ 快"这句话,在吞吐和单条延迟两个维度上分别成立吗?给出 Confluent 基准里的关键数字。
  3. 一个 RabbitMQ 队列接了多个 competing consumers,为什么"queue 是 FIFO 的"不等于"消息被 FIFO 地处理"?两个断裂点是什么?
  4. 场景判断:系统要把一份订单变更流喂给 5 个互相独立的下游(风控、报表、搜索索引、对账、归档),每个下游都要能从历史某点重新消费一遍,吞吐中等偏上。选 RabbitMQ、Kafka 还是 RabbitMQ Streams?为什么?
答案(先做完再展开)
  1. Kafka 的 broker 什么都不做——它只按 partition 追加字节,路由(选 partition)发生在 producer 的 key 哈希里,过滤/扇出发生在消费端或 Kafka Streams 里。RabbitMQ 的 exchange 把路由与过滤放进了 broker(topic/headers/一致性哈希等),代价是逐消息匹配带来的吞吐上限。
  2. 吞吐维度成立:Kafka ~605 MB/s vs RabbitMQ ~38 MB/s(i3en.2xlarge、1KB、3 副本),约 16 倍。单条延迟维度翻转:~30 MB/s 低负载下 RabbitMQ p99 约 1 ms,低于 Kafka;Kafka 的批量/linger 给单条加了延迟。"快"指吞吐高,不指每条都更快。
  3. 两个断裂点:① competing consumers——消息一人一条派给不同消费者,谁先处理完和谁先拿到无关,完成顺序与到达顺序脱钩;② requeue——nack 或断连退回的消息被重新插入队列(常在队首),越过后面的消息。要恢复有序并行,需一致性哈希 exchange + single active consumer。
  4. 选 Kafka。判据:多个独立消费者各自读同一份流(不可变日志天然支持,互不干扰)+ 每个都要从历史某点重放(offset 可倒带)+ 吞吐中等偏上。RabbitMQ classic/quorum 队列 ack 即删、无法重放,且扇出给 5 个独立读者要建 5 套队列。Streams 能满足重放与多读者,但题面吞吐偏上且无"必须留在 RabbitMQ 生态"的约束,Kafka 是更直接的归属;若题面强约束"已重度使用 RabbitMQ、不愿引入 Kafka",则 Streams 是折中解。
进阶挑战 · 刚好够不着

给一个"混合需求"系统拆分消息中间件

一个电商后台同时有三类消息流:(a) 下单后触发的履约工作流——发券、扣库存、通知,每步要重试、要死信、要按业务优先级插队;(b) 全站行为埋点事件——每秒数十万条,要长期留存供离线分析与模型训练随时重跑;(c) 一份订单状态变更流,要实时喂给 3 个内部服务,偶尔需要回放最近一天补数据。请为三类流各选一种方案(RabbitMQ / Kafka / RabbitMQ Streams),并各用一句话给出决定性判据。

提示(卡住再展开)

(a) → RabbitMQ:重试/DLX/priority/逐消息工作流是它的原生语义,吞吐不是约束。(b) → Kafka:每秒数十万条 + 长期留存 + 随时重跑,正是日志 + tiered storage 的主场,RabbitMQ 30K msg/s 就撞墙。(c) → RabbitMQ Streams 或 Kafka 皆可:要重放 + 多读者,若团队已重度用 RabbitMQ 且只需"最近一天回放"这种不极端的留存,Streams 免去引入 Kafka;若已有 Kafka 或要对接流处理生态,则 Kafka。决定性判据分别是:逐消息工作流语义、超高吞吐+长留存、重放+生态归属。