消息队列
消息队列(Message Queue,简称 MQ)是分布式系统中非常常见的基础组件。它的核心作用不是“让系统变复杂”,而是把原来必须同步完成的一段调用,拆成“先把消息可靠写进去,再由下游慢慢处理”的异步链路。
通俗理解:系统 A 不直接调用系统 B,而是把一张“任务单”放到队列里;系统 B 有空时从队列取任务并处理。这样可以降低系统之间的耦合,也能在流量突然变大时保护下游服务。
为什么需要消息队列
1. 异步解耦
在没有 MQ 的情况下,订单服务创建订单后可能要同步调用库存、积分、短信、物流等服务。任何一个服务慢或者失败,都会影响创建订单接口。
使用 MQ 后,订单服务只需要把“订单已创建”事件发送出去,下游系统各自订阅并处理,主流程更短,也更容易扩展。
flowchart TD
A["订单服务"] --> B["消息队列"]
B --> C["库存服务"]
B --> D["积分服务"]
B --> E["短信服务"]
B --> F["物流服务"]2. 削峰填谷
秒杀、活动、支付回调这类场景会出现瞬时高峰。如果请求全部直接打到数据库或下游服务,很容易把系统压垮。MQ 可以先承接高峰流量,下游按自己的处理能力稳定消费。
flowchart TD
A[突发请求] --> B[写入 MQ]
B --> C[消费者按固定速率处理]
C --> D[数据库压力平稳]3. 最终一致性
分布式系统里很难让多个服务同时成功或同时失败。MQ 更常见的做法是保证主业务先成功,再通过可靠消息、重试、补偿等机制让下游最终处理成功。
核心概念
| 概念 | 说明 |
|---|---|
| Producer | 生产者,负责发送消息 |
| Consumer | 消费者,负责接收并处理消息 |
| Broker | 消息服务器,负责接收、存储、投递消息 |
| Topic / Queue | 消息分类或队列,不同产品命名略有差异 |
| Offset | 消费进度,记录消费者处理到哪里 |
| ACK | 消费确认,表示消息已经被成功处理 |
| Retry | 重试,消费失败后再次投递 |
| DLQ | 死信队列,多次失败后进入的异常队列 |
消息队列的整体流程
flowchart TD
P["生产者构建消息"] --> S["发送到 Broker"]
S --> W["Broker 持久化"]
W --> R["返回发送结果"]
W --> Q["等待消费者拉取或推送"]
Q --> C["消费者处理业务"]
C --> A{"处理成功?"}
A -->|是| ACK["提交消费确认"]
A -->|否| Retry["进入重试或死信"]常见问题
消息会不会丢
消息是否会丢,取决于三个阶段:
- 生产阶段:生产者发送失败是否重试,是否等待 Broker 确认。
- 存储阶段:Broker 是否刷盘,是否有主从复制。
- 消费阶段:消费者是否在业务真正成功后再确认。
所以“用 MQ 就不会丢消息”是不准确的。正确做法是根据业务选择可靠发送、持久化、手动 ACK、消费幂等等策略。
消息会不会重复
多数 MQ 更容易保证“至少投递一次”,也就是说消息可能重复。网络超时、消费者处理成功但 ACK 失败、Broker 重试等情况都会造成重复投递。
因此消费端必须设计幂等。常见方式包括唯一业务号去重、数据库唯一索引、状态机判断、Redis setnx 等。
消息顺序如何保证
顺序消息通常只能在一个局部范围内保证,例如同一个订单的消息进入同一个队列,由同一个消费者顺序处理。全局顺序会牺牲吞吐量,实际业务中很少使用。
消息堆积怎么办
消息堆积不是简单“加消费者”就能解决。要先判断生产速度、消费速度、消费者错误率、重试量、下游数据库或接口耗时、分区/队列是否热点,再决定限流、扩容、批处理、失败隔离、死信、补偿还是调整分区队列。
详细学习:消息堆积与背压。
目录
从零到生产级掌握
MQ 主线课程。把生产者、Broker、消费者、ACK、Offset、可靠性、重复消费、幂等、顺序消息、事务消息、重试、死信、堆积、背压、RabbitMQ/Kafka/RocketMQ 选型串成完整体系。
从零到精通验收清单
用可验证任务判断自己是不是真的学会 MQ。它按作用、发送、Broker、消费、可靠性、幂等、顺序、事务、死信、堆积、Lag、背压、三大 MQ 选型和商业场景拆成验收标准。
商业场景训练营
用订单创建事件、Outbox 本地消息表、消费者幂等、业务成功后 ACK、ES 同步失败补偿、消息堆积、扩容反弹、Lag 分布、RabbitMQ Ready/Unacked 和三大 MQ 选型,把原理放到商业项目里跑起来。
消息堆积与背压
生产排障高频知识点。重点学习消息积压、消费 Lag、Ready/Unacked、重试堆积、背压、限流、扩容、死信和补偿。
扩容后再次堆积
专门解释“消费者扩容后短暂有效,后面又开始堆积”的原理。重点学习有效消费 TPS、分区/队列并行度、下游瓶颈、重试风暴、热点 key、线程池和连接池排查。
Lag 分布与堆积定位
专门解释“Lag 分布是什么”。重点学习为什么不能只看总 Lag,而要按 Kafka Partition、RocketMQ MessageQueue、RabbitMQ Queue/Unacked 分布判断整体消费慢、局部热点、慢消息、重试风暴和顺序阻塞。
MQ面试题
刷题复习入口。详细原理回到对应知识点页学习。
RabbitMQ
适合传统企业应用、复杂路由、任务分发、延迟/死信处理、对 AMQP 生态有要求的场景。
Kafka
适合高吞吐事件流、日志采集、用户行为埋点、CDC 数据同步、实时计算、数据管道和可回放场景。
RocketMQ
适合高吞吐、金融级可靠消息、事务消息、顺序消息、消费重试等场景。
RabbitMQ、Kafka 和 RocketMQ 怎么选
| 对比项 | RabbitMQ | Kafka | RocketMQ |
|---|---|---|---|
| 核心模型 | Exchange + Queue + Binding | Topic + Partition + Log | Topic + Queue + ConsumerGroup |
| 吞吐能力 | 中高吞吐,路由能力强 | 高吞吐,适合连续事件流 | 高吞吐,适合业务消息堆积 |
| 消息保留 | 偏队列模型,消费后可删除 | 按时间或大小保留,天然支持回放 | 支持堆积、重试和死信 |
| 顺序消息 | 单队列可做局部顺序 | 分区内顺序 | 队列内顺序,业务语义更直接 |
| 事务消息 | 通常靠发布确认、业务表、补偿实现 | 支持 Kafka 内部事务,外部库需 Outbox/CDC | 原生事务消息能力较强 |
| 路由能力 | 很强,Exchange 类型丰富 | 较弱,主要靠 Topic 和 key | 中等,靠 Topic、Tag、Key |
| 适合场景 | 企业集成、任务队列、复杂路由、延迟/死信 | 日志、埋点、数据同步、实时计算、事件总线 | 电商交易、金融异步链路、事务消息、顺序消息 |
学习建议:先理解本页的 MQ 通用概念,再学习 RabbitMQ 核心模型、Kafka 架构与存储原理 或 RocketMQ 领域模型。如果重点关注分布式事务,直接看 RocketMQ 事务消息;如果重点关注日志、数据同步和回放,直接看 Kafka 商业常用场景。
代码 Demo:订单创建后发送消息
下面用伪代码演示 MQ 在业务里的位置。重点是:主业务先落库,再发送事件消息,下游异步处理。
@Service
public class OrderService {
private final OrderRepository orderRepository;
private final MessageProducer messageProducer;
public OrderService(OrderRepository orderRepository, MessageProducer messageProducer) {
this.orderRepository = orderRepository;
this.messageProducer = messageProducer;
}
@Transactional
public Long createOrder(CreateOrderCommand command) {
Order order = new Order(command.userId(), command.productId(), command.count());
orderRepository.save(order);
OrderCreatedEvent event = new OrderCreatedEvent(order.getId(), order.getUserId());
messageProducer.send("order.created", event);
return order.getId();
}
}消费端必须做幂等:
@Component
public class StockConsumer {
public void onMessage(OrderCreatedEvent event) {
if (alreadyProcessed(event.orderId())) {
return;
}
reduceStock(event.orderId());
markProcessed(event.orderId());
}
}真实项目可以用 RabbitMQ、Kafka、RocketMQ 等产品实现 MessageProducer,但“可靠发送、业务幂等、失败重试”这三个原则都不能省。
