后端工程(03):消息中间件——Kafka/RabbitMQ 模型对比、可靠投递、消费组、延迟权衡
更新时间:2026-09-01。本文是
backend/engineering/后端工程第 03 篇,接 鉴权体系。消息中间件是后端系统解耦、削峰填谷、异步处理的核心组件。Kafka 和 RabbitMQ 是最主流的两个消息系统,但设计理念完全不同。
本文要回答的问题
- Kafka 和 RabbitMQ 的区别是什么?什么时候选哪个?
- 消息可靠投递怎么保证?at-least-once、exactly-once 是什么?
- 消费组是什么?怎么实现消息广播和点对点?
- 延迟和吞吐量的权衡怎么选?
一、Kafka vs RabbitMQ 对比
| 对比 | Kafka | RabbitMQ |
|---|---|---|
| 模型 | 发布-订阅(日志) | 队列 + 交换机 |
| 存储 | 持久化到磁盘,顺序读写 | 内存优先,可持久化 |
| 吞吐量 | 极高(百万/秒) | 中等(万/秒) |
| 延迟 | 略高(ms 级) | 极低(μs 级) |
| 消息顺序 | 分区内有序 | 队列内有序 |
| 路由 | 基于 topic + partition key | 交换机 + 路由键 |
| 消费模型 | 拉取(pull) | 推送(push),可拉取 |
| 消息保留 | 可配置(按时间/大小) | 消费后删除 |
| 适合场景 | 日志、流处理、大数据 | 任务队列、事件驱动 |
选型建议
Kafka 适合:
- 日志收集、流处理
- 大数据管道
- 高吞吐量场景
- 消息持久化需要长期保留
RabbitMQ 适合:
- 任务队列、延迟队列
- 复杂的路由模式
- 低延迟场景
- 需要灵活的交换机配置二、Kafka 核心概念
text
Kafka 架构:
┌─────────┐ ┌─────────┐ ┌─────────┐
│ Producer │ │ Producer │ │ Producer │
└────┬────┘ └────┬────┘ └────┬────┘
│ │ │
└────────────┼────────────┘
│
▼
┌───────────────┐
│ Kafka Cluster │
├─────────────────┤
│ Topic: orders │
│ ├ Partition 0 │
│ ├ Partition 1 │
│ └ Partition 2 │
└───────────────┘
│
┌────────────┼────────────┐
│ │ │
▼ ▼ ▼
┌─────────┐ ┌─────────┐ ┌─────────┐
│ Consumer│ │ Consumer│ │ Consumer│
│ Group A │ │ Group A │ │ Group B │
└─────────┘ └─────────┘ └─────────┘核心概念:
- Topic:消息分类,类似数据库表
- Partition:topic 的分区,每个分区是一个有序的日志(append-only)
- Offset:消息在分区中的偏移量,消费者记录 offset 表示消费到哪里
- Consumer Group:消费组,组内消费者共同消费 topic 的分区,一个分区只能被组内一个消费者消费
- Broker:Kafka 服务器节点
三、RabbitMQ 核心概念
text
RabbitMQ 架构:
┌──────────┐
│ Producer │
└────┬─────┘
│
▼
┌──────────┐
│ Exchange │ ← 交换机:决定消息路由到哪个队列
└────┬─────┘
│
▼
┌──────────┐
│ Queue │ ← 队列:存储消息
└────┬─────┘
│
▼
┌──────────┐
│ Consumer │
└──────────┘交换机类型:
direct:精确匹配路由键,消息发送到 routing_key 完全匹配的队列topic:模式匹配路由键,orders.*匹配orders.created、orders.updatedfanout:广播,消息发送到所有绑定的队列headers:按头匹配,不常用
四、消息可靠投递
三种语义
| 语义 | 说明 | 实现 |
|---|---|---|
| at-most-once | 最多一次,可能丢失 | 不重试,适合不重要通知 |
| at-least-once | 至少一次,可能重复 | 生产端确认 + 消费端确认,重试 |
| exactly-once | 精确一次,不丢不重 | 事务 + 幂等消费,成本高 |
Kafka 的可靠投递
text
# 生产端:acks 参数
acks=0 # 不等待确认,最快,可能丢
acks=1 # 等待 leader 确认,不丢(如果 leader 挂了可能丢)
acks=all # 等待所有副本确认,最安全,不丢
# 消费端:enable.auto.commit
enable.auto.commit=false # 手动提交 offset,确保处理完再提交
# 处理完再提交,即使崩溃也不会丢失RabbitMQ 的可靠投递
text
# 生产端:publisher confirm
channel.confirmSelect() # 开启发布确认
channel.waitForConfirms() # 等待确认,确保消息到达 broker
# 消费端:手动 ack
channel.basicConsume(queue, autoAck=false, callback)
callback 中处理完消息后:channel.basicAck(deliveryTag, false)
# 确保消息处理完再确认五、消息重复消费问题
text
# at-least-once 语义下,消息可能重复消费
# 消费者处理完消息,但提交 offset 前崩溃
# 重平衡后,相同的消息会被重新消费
# 解决:幂等消费
# 1. 消息带唯一 ID
# 2. 消费端用唯一 ID 去重
# 例子:数据库唯一键
INSERT INTO processed_messages (message_id, ...)
VALUES ('msg-123', ...)
ON CONFLICT (message_id) DO NOTHING;
# 重复消息不会重复处理六、延迟和吞吐量权衡
| 策略 | 延迟 | 吞吐量 | 说明 |
|---|---|---|---|
| Kafka acks=0 | 极低 | 极高 | 不确认,可能丢数据 |
| Kafka acks=1 | 低 | 高 | 确认 leader 写入 |
| Kafka acks=all | 高 | 中 | 确认所有副本,最安全 |
| RabbitMQ 内存 | 极低 | 高 | 不持久化,重启丢失 |
| RabbitMQ 持久化 | 中 | 中 | 持久化到磁盘 |
经验:
- 日志、监控:acks=0 或 1,允许丢失少量数据
- 订单、支付、核心业务:acks=all,手动提交 offset,确保不丢
七、常见坑对照
| 坑 | 现象 | 对策 |
|---|---|---|
| Kafka 分区数太多 | 文件句柄不够,性能下降 | 分区数 = 消费者数 * 2 |
| 消费组重平衡 | 所有消费者暂停 | 减少重平衡,session.timeout.ms 调大 |
| 消息堆积来不及消费 | 延迟越来越高 | 增加消费者数量,分区数要足够 |
| 重复消费 | 业务数据重复 | 幂等消费,消息去重 |
相关与延伸
下一篇:微服务架构——服务拆分、注册发现、配置中心、网关、链路追踪;API 设计规范,见 API 设计规范。
一句话总结
消息中间件:Kafka 高吞吐量(百万/秒),pull 模型,分区 + 消费组,适合日志和流处理;RabbitMQ 低延迟,push 模型,交换机 + 路由键,适合任务队列和事件驱动;可靠投递:生产端 acks=all/confirm,消费端手动提交 offset/ack;at-least-once 保证不丢,但需要幂等消费处理重复消息;延迟和吞吐量权衡:核心业务用 acks=all + 手动提交,日志用 acks=0/1。