系统设计:消息队列架构
消息队列是分布式系统的异步通信中枢,解耦、削峰、异步三大核心价值。
为什么需要消息队列?
场景一:系统解耦
无 MQ:
订单服务 → 库存服务 → 物流服务 → 通知服务
(强依赖,一环失败全部回滚)
有 MQ:
订单服务 → MQ「订单创建」→ 库存服务
→ 物流服务
→ 通知服务
(各消费者独立处理,互不影响)
场景二:流量削峰
秒杀场景:
瞬时 100 万请求 → 订单服务只能处理 1 万/秒
无 MQ:99% 请求超时/失败
有 MQ:请求写入 MQ,订单服务匀速消费,MQ 缓冲峰值
场景三:异步处理
用户注册:
无 MQ:注册 → 写DB(50ms) → 发邮件(200ms) → 发短信(100ms) = 350ms
有 MQ:注册 → 写DB(50ms) → 发MQ(5ms) = 55ms,邮件/短信异步消费
两种消息模型
点对点(Point-to-Point / Queue)
Producer ──→ Queue ──→ Consumer A
└──→ Consumer B
特点:
- 一条消息只被一个消费者消费
- 消费者竞争消费(类似线程池任务分发)
- 典型:RabbitMQ 的 Queue、RocketMQ 的普通消息
发布-订阅(Publish-Subscribe / Topic)
Producer ──→ Topic ──→ Subscriber A
├──→ Subscriber B
└──→ Subscriber C
特点:
- 一条消息被所有订阅者消费
- 支持消息回溯(按 offset 重新消费)
- 典型:Kafka 的 Topic、RocketMQ 的广播消息
Kafka 核心原理
架构
┌──────────┐ ┌──────────────────────────────────────┐
│ Producer │────→│ Kafka Cluster │
└──────────┘ │ ┌─────────┐ ┌─────────┐ ┌─────────┐ │
│ │Broker 1 │ │Broker 2 │ │Broker 3 │ │
│ │ ┌─────┐ │ │ ┌─────┐ │ │ ┌─────┐ │ │
│ │ │ P0 │ │ │ │ P1 │ │ │ │ P2 │ │ │
│ │ │ P3 │ │ │ │ P4 │ │ │ │ P5 │ │ │
│ │ └──┬──┘ │ │ └──┬──┘ │ │ └──┬──┘ │ │
│ └────┼────┘ └────┼────┘ └────┼────┘ │
│ └───────────┼───────────┘ │
└───────────────────┼───────────────────┘
│
┌──────────┐ ┌───▼────┐
│Consumer A│──→│ Group │
│Consumer B│──→│ "g1" │
└──────────┘ └────────┘
关键概念:
| 概念 | 说明 |
|---|---|
| Topic | 消息主题,逻辑上的消息集合 |
| Partition | Topic 的物理分片,每个 partition 是有序的 |
| Broker | Kafka 服务器节点 |
| Consumer Group | 消费者组,组内消费者共同消费一个 Topic 的所有 partition |
| Offset | 消息在 partition 中的位置(位移),消费者用 offset 标记消费进度 |
| Replication | 副本机制,每个 partition 有 leader + 若干 follower |
Partition 与 Consumer Group 的关系
Topic T1 有 6 个 partition:P0, P1, P2, P3, P4, P5
Consumer Group G1 有 3 个消费者:
C1 消费 P0, P1
C2 消费 P2, P3
C3 消费 P4, P5
Consumer Group G2 也有 3 个消费者:
独立消费同样的 6 个 partition
(每个 group 维护自己的 offset)
关键规则:一个 partition 同一时间只能被 group 内的一个消费者消费。
Kafka 的高吞吐设计
- 顺序写磁盘:Kafka 顺序写磁盘的速度比随机写内存还快(磁盘预读 + OS 页缓存)
- 零拷贝(Zero-Copy):
sendfile()系统调用,数据从磁盘直接通过网卡发送,不经过用户态 - 批量处理:Producer 批量发送、Consumer 批量拉取
- 压缩:Producer 端压缩(GZIP/Snappy/LZ4),Broker 直接存储压缩数据
Kafka vs RabbitMQ 选型对比
| 维度 | Kafka | RabbitMQ |
|---|---|---|
| 设计目标 | 高吞吐日志流处理 | 通用消息路由 |
| 吞吐量 | 百万级/秒 | 万级/秒 |
| 消息模型 | 发布-订阅 | 点对点 + 路由 |
| 消息顺序 | Partition 内有序 | Queue 内有序 |
| 消息回溯 | ✅ 支持(按 offset) | ❌ 消费即删除 |
| 消息持久化 | ✅ 默认持久化 | 可选 |
| 延迟 | 毫秒级(批量) | 微秒级 |
| 复杂度 | 较高(Zookeeper/KRaft) | 较低 |
| 适用场景 | 日志收集、实时流处理 | 任务队列、RPC 替代 |
选型建议:
- 大数据、日志、事件流 → Kafka
- 任务分发、延时队列、复杂路由 → RabbitMQ
消息可靠性三大保障
1. 生产端:ACK 确认机制
from kafka import KafkaProducer
producer = KafkaProducer(
bootstrap_servers=['kafka1:9092'],
acks='all', # 0: 不等待 | 1: leader确认 | all: 所有ISR确认
retries=3, # 发送失败重试次数
retry_backoff_ms=1000
)
# acks='all' 保证消息至少被一个 ISR(In-Sync Replica)确认才返回成功
2. 消费端:手动提交 Offset
from kafka import KafkaConsumer
consumer = KafkaConsumer(
'my-topic',
bootstrap_servers=['kafka1:9092'],
group_id='my-group',
enable_auto_commit=False # 关闭自动提交
)
for msg in consumer:
try:
process(msg)
consumer.commit_sync() # 处理成功后手动提交 offset
except Exception as e:
# 不提交 offset,消息会被重新消费
log_error(e)
至少一次 vs 恰好一次 vs 最多一次:
| 语义 | 实现方式 | 特点 |
|---|---|---|
| 最多一次 | 先提交 offset 再处理 | 可能丢消息 |
| 至少一次 | 先处理再提交 offset | 可能重复消费 |
| 恰好一次 | 幂等消费 + 去重表/事务 | 实现复杂,推荐 |
3. 幂等性:防止重复消费
业务层面实现幂等:
# 方案一:数据库唯一约束
# 消费消息时,将消息 ID 作为业务表的唯一键
INSERT INTO orders (msg_id, order_id, ...) VALUES (?, ?, ...)
-- 重复消费时因唯一键冲突自动失败
# 方案二:Redis SETNX
if redis.setnx(f"consumed:{msg_id}", "1", ex=86400):
process(msg)
else:
skip(msg) # 已消费过
# 方案三:Kafka 内置幂等(Producer 端)
# enable.idempotence=true
# Kafka 自动去重因网络重试导致的重复发送
消息顺序性保证
Kafka 的顺序保证
分区级别有序,全局无序:
# 如果业务需要全局有序,可以:
# 方案一:单 partition(牺牲吞吐量)
producer.send('orders', key='all_orders', value=order_data)
# 方案二:按业务键分区(如用户级别有序)
# 相同 user_id 的消息进入同一 partition
producer.send('orders', key=str(user_id), value=order_data)
跨分区顺序问题
Partition 0: 订单创建(1) → 订单支付(3)
Partition 1: 订单修改(2)
消费者看到顺序可能是:创建 → 修改 → 支付(实际应该是:创建 → 修改 → 支付)
解决:按订单 ID 分派到同一个 partition
高可用设计
Kafka 的 Replication
Partition P0:
Leader: Broker 1(读写)
Follower: Broker 2(同步复制)
Follower: Broker 3(同步复制)
ISR(In-Sync Replicas)= [Broker1, Broker2, Broker3]
如果 Broker 1 宕机:
- 从 ISR 中选举新 Leader(如 Broker 2)
- Producer 和 Consumer 自动切换到新 Leader
副本同步机制
- Leader 收到消息 → 写入本地 log → 发送 ACK
- Follower 拉取消息 → 写入本地 log → 更新 fetch offset
- ISR 列表动态维护:落后太多的 follower 被移出 ISR
面试答题框架
第一步:明确场景(30秒)
我需要确认:吞吐量要求?是否需要消息回溯?延迟容忍度?消费者数量?
第二步:选型(1分钟)
高吞吐日志流 → Kafka;通用任务队列 → RabbitMQ;云原生 → Pulsar。
第三步:Kafka 架构(2分钟)
Topic → Partition → Broker → Consumer Group → Offset。重点讲 partition 与 consumer 的映射关系。
第四步:可靠性(2分钟)
Producer ACK(acks=all)+ Consumer 手动提交 offset + 业务幂等(唯一键/Redis SETNX)。
第五步:高可用(1分钟)
Replication + ISR + Leader 选举,保证 N-1 个节点故障时仍可服务。
第六步:扩展问题(面试官追问)
- 消息积压怎么办?(扩容消费者、增加 partition、降低处理延迟)
- 如何实现延时消息?(Kafka 本身不支持,可用 RabbitMQ 死信队列或自定义实现)
- Kafka 与 Pulsar 的区别?(Pulsar 存储计算分离、更好的多租户、内置 geo-replication)
常见问题
Q:Kafka 为什么比 RabbitMQ 吞吐高?
- 顺序写磁盘 + OS 页缓存;2. 零拷贝;3. 批量处理;4. 更简单的协议。
Q:Consumer Group 重平衡(Rebalance)有什么问题?
消费者加入/退出时触发重平衡,期间全部消费者停止消费,造成延迟尖刺。Kafka 2.4+ 引入 Incremental Rebalance 优化。
Q:Kafka 的 offset 存储在哪里?
Kafka 0.9+ 存储在内部的 __consumer_offsets Topic 中,之前版本存储在 Zookeeper。
Q:消息队列和 RPC(如 gRPC)怎么选?
需要实时响应、强一致性 → RPC;需要解耦、削峰、异步 → MQ。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。