设计一个消息队列¶
一、需求¶
设计一个 MQ,要支持:
- 生产/消费消息。
- 持久化,宕机不丢。
- 消费失败可重试。
- 顺序消息 / 广播 / 集群消费。
- 高可用。
二、核心模型¶
1. Topic / Queue / Message¶
- Topic:一类消息的逻辑名。
- Queue:Topic 下的物理队列(通常多个以并发)。
- Message:消息体 + 元数据(ID、Tag、Key、时间戳)。
2. 消费模式¶
- 集群消费:同一条消息只被一个消费组内一个消费者处理。
- 广播消费:同一条消息被所有消费者处理。
三、存储设计¶
方案 1:文件存储(Kafka / RocketMQ 思路)¶
- 每个 Topic 多 Queue,每个 Queue 对应一个顺序写文件。
- 消息 append 写入,性能极高(顺序写磁盘比随机写内存还快)。
- 用索引文件记录 offset → 物理位置。
- 定期清理过期文件。
方案 2:数据库存储¶
简单但性能差,只适合低吞吐场景。
方案 3:Kafka 的分区日志¶
四、投递语义¶
| 语义 | 含义 | 实现 |
|---|---|---|
| At most once | 最多一次,可能丢 | 发完不重试 |
| At least once | 至少一次,可能重复 | 失败重试,消费端 ACK |
| Exactly once | 精确一次 | 事务消息、去重 |
工程上一般做到 At least once + 消费端幂等。
五、关键机制¶
1. ACK 确认¶
消费者处理完消息后回复 ACK;未 ACK 的消息超时后重投。
2. 重试与死信队列¶
消费失败 N 次后,投递到 DLQ(Dead Letter Queue),人工处理。
3. 顺序消息¶
同一业务 key(如 orderId)的消息路由到同一 Queue,单 Queue 单消费者消费。
4. 延迟消息¶
RocketMQ 支持 18 个延迟级别;Kafka 用时间轮或外部实现。
5. 回溯消费¶
记录 offset,支持按时间 / offset 重新消费。
六、高可用¶
- Broker 主从:主写从同步,主挂从接管。
- Controller / NameServer:元数据管理。
- 生产端重试:发送失败切换 broker。
- 消费端集群:一个消费组多实例,Rebalance 分配 Queue。
七、为什么不直接用 Redis List 做 MQ¶
- 没有 ACK、重试、死信。
- 持久化不可靠。
- 不能水平扩展。
- 没有顺序、过滤、延迟等高级特性。
面试答题框架
设计 MQ 按:模型(Topic/Queue)→ 存储(顺序写文件)→ 投递语义 → ACK → 重试/死信 → 顺序 → 高可用 → 监控。