04-RabbitMQ 与消息队列高级
对应原始资料:
项目一前置课b/day13-RabbitMQ、高级02-MQ高级
一、为什么要用 MQ
消息队列(Message Queue)三大作用:
- 异步:注册后异步发短信、邮件,提升响应速度。
- 解耦:订单系统和积分系统通过 MQ 通信,互不影响。
- 削峰:高并发时消息先入队,消费者按能力处理,保护下游。
二、常见 MQ 对比
| MQ | 语言 | 特点 |
|---|---|---|
| RabbitMQ | Erlang | 功能全、可靠性高、AMQP 协议、互联网公司常用 |
| Kafka | Scala | 吞吐量极高、日志/大数据场景 |
| RocketMQ | Java | 阿里出品,事务消息,电商场景 |
| ActiveMQ | Java | 老牌,逐渐减少 |
三、RabbitMQ 核心概念
Producer → Exchange(交换机)→ Queue(队列)→ Consumer
↑ Binding(绑定)- Broker:MQ 服务器。
- VirtualHost:虚拟主机(隔离)。
- Exchange:路由消息到队列。
- Queue:存储消息。
- Binding:交换机和队列的绑定关系。
- RoutingKey:路由键。
四、四种交换机类型
| 类型 | 路由规则 |
|---|---|
| Direct | RoutingKey 完全匹配 |
| Fanout | 广播,发给所有绑定的队列 |
| Topic | 按模式匹配(* 一个词,# 多个词) |
| Headers | 按消息头匹配(少用) |
五、安装
bash
docker run -d --name rabbitmq -p 5672:5672 -p 15672:15672 rabbitmq:management
# 5672 通信端口,15672 管理控制台默认账号 guest / guest。
六、SpringBoot 整合
1. 依赖
xml
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>2. 配置
yaml
spring:
rabbitmq:
host: 127.0.0.1
port: 5672
username: guest
password: guest
listener:
simple:
acknowledge-mode: manual # 手动 ACK
prefetch: 1 # 每次拉取数量3. 发送消息
java
@Autowired private RabbitTemplate rabbit;
rabbit.convertAndSend("exchange", "routing.key", "hello");
// 对象会自动用 Jackson 转 JSON4. 接收消息
java
@RabbitListener(bindings = @QueueBinding(
value = @Queue("order.queue"),
exchange = @Exchange("order.exchange"),
key = "order.create"
))
public void receive(String msg, Channel channel, Message message) throws IOException {
System.out.println("收到:" + msg);
channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
}七、可靠性(MQ 高级重点)
消息从生产者到消费者,可能丢失的三个环节:
- 生产者 → MQ:发送失败。
- MQ 自身:宕机丢消息。
- MQ → 消费者:消费失败。
1. 生产者确认
- Publisher Confirm:消息成功到达 Exchange 的回调。
- Publisher Return:消息从 Exchange 到 Queue 失败的回调。
2. 持久化
- 交换机、队列、消息都设为
durable,重启不丢。
3. 消费者手动 ACK
java
channel.basicAck(tag, false); // 确认
channel.basicNack(tag, false, true);// 否认,第三个参数 requeue=true 重回队列业务异常 → 进死信队列;处理失败 → 重试 N 次后进死信。
八、死信队列(DLX)
消息变成"死信"的情况:
- 消息被
basicNack/reject且 requeue=false。 - 消息 TTL 过期。
- 队列长度超限。
给队列绑定死信交换机:
java
@Queue(value = "order.queue",
arguments = {
@Argument(name="x-dead-letter-exchange", value="dlx.exchange"),
@Argument(name="x-dead-letter-routing-key", value="dlx.key")
})九、延迟队列(典型应用:订单超时取消)
RabbitMQ 本身没延迟队列,常见实现:
- TTL + 死信队列:消息设过期时间,过期转死信 → 消费者监听死信队列。
- rabbitmq-delayed-message-exchange 插件:直接发延迟消息。
十、消息重复消费与幂等性
网络抖动可能导致同一条消息被消费多次。解决:
- 业务层保证幂等性:唯一 ID + 数据库唯一约束 / Redis 记录已处理 ID。
十一、消息堆积与顺序性
- 堆积:增加消费者、扩队列、临时扩容消费线程。
- 顺序:把同一业务键的消息路由到同一队列,单消费者顺序消费。
练习建议
- 用 Docker 启动 RabbitMQ,在管理控制台手动发一条消息。
- SpringBoot 实现「简单队列」收发。
- 用 Fanout 交换机实现「一条消息多消费者处理」。
- 实现订单超时取消:TTL + 死信队列。
- 实现手动 ACK + 消费失败重试。