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 旁边逐项展开:同一个"消息从生产者到消费者"的过程,两套系统把哪些职责放在了哪一侧。
| 职责 | 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 的状态机活在消费者里。
对一个 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 长轮询拉取,速率由消费者自己定。
| 维度 | 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(背压,下游处理不过来时让上游慢下来)在两套系统里是两种东西。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 摆在一起看。
| 维度 | RabbitMQ | Kafka |
|---|---|---|
| 并行模型 | 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)给出一组常被引用的数字:
| 系统 | 峰值吞吐 | 延迟特征 |
|---|---|---|
| Kafka | ~605 MB/s | 200K msg/s 时 p99 约 5 ms;吞吐压倒性领先 |
| Pulsar | ~305 MB/s | 200K 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 墙。
"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。
| 场景 | 选谁 | 为什么 |
|---|---|---|
| 复杂路由 / 按条件扇出到每个应用各自的队列 | RabbitMQ | exchange(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 Streams | append-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(分层存储,热数据在本地盘、冷数据下沉到对象存储,便宜地长期留存并重放)。两个系统正从相反的出身互相靠拢:一个加日志语义,一个加廉价长留存。
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
先合上教程,把答案写在纸上或编辑器里。写完再点开对照——直接点开等于把这一节当再读一遍。
- Kafka 的 broker 对"一条消息该投给谁、要不要过滤"做了什么?这和 RabbitMQ 的 exchange 有什么本质区别?
- "Kafka 比 RabbitMQ 快"这句话,在吞吐和单条延迟两个维度上分别成立吗?给出 Confluent 基准里的关键数字。
- 一个 RabbitMQ 队列接了多个 competing consumers,为什么"queue 是 FIFO 的"不等于"消息被 FIFO 地处理"?两个断裂点是什么?
- 场景判断:系统要把一份订单变更流喂给 5 个互相独立的下游(风控、报表、搜索索引、对账、归档),每个下游都要能从历史某点重新消费一遍,吞吐中等偏上。选 RabbitMQ、Kafka 还是 RabbitMQ Streams?为什么?
答案(先做完再展开)
- Kafka 的 broker 什么都不做——它只按 partition 追加字节,路由(选 partition)发生在 producer 的 key 哈希里,过滤/扇出发生在消费端或 Kafka Streams 里。RabbitMQ 的 exchange 把路由与过滤放进了 broker(topic/headers/一致性哈希等),代价是逐消息匹配带来的吞吐上限。
- 吞吐维度成立:Kafka ~605 MB/s vs RabbitMQ ~38 MB/s(i3en.2xlarge、1KB、3 副本),约 16 倍。单条延迟维度翻转:~30 MB/s 低负载下 RabbitMQ p99 约 1 ms,低于 Kafka;Kafka 的批量/linger 给单条加了延迟。"快"指吞吐高,不指每条都更快。
- 两个断裂点:① competing consumers——消息一人一条派给不同消费者,谁先处理完和谁先拿到无关,完成顺序与到达顺序脱钩;② requeue——nack 或断连退回的消息被重新插入队列(常在队首),越过后面的消息。要恢复有序并行,需一致性哈希 exchange + single active consumer。
- 选 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。决定性判据分别是:逐消息工作流语义、超高吞吐+长留存、重放+生态归属。