第 8 章 Spring Boot 消息服务
写给新手的话:
前面我们学了数据库、缓存、安全管理,都是同步的:发请求 → 处理 → 返回结果。
但是有些场景,同步处理不太合适:
- 注册完要发邮件、发短信,用户要等半天
- 秒杀的时候,一下子几万请求过来,数据库扛不住
- 系统之间调用,一个挂了其他的也受影响
这时候就需要消息队列了!
消息队列可以帮我们:异步处理、系统解耦、流量削峰。
这一章我们学最流行的消息中间件之一:RabbitMQ。
学习建议:
- 先理解核心概念:队列、交换机、生产者、消费者
- 搞清楚几种工作模式的区别和应用场景
- 多动手,把每种模式都跑一遍
- Spring Boot 整合是重点,实际项目中用得最多
准备好了吗?我们开始吧!
8.1 先搞明白:消息队列那些事儿
8.1.1 什么是消息队列
消息队列(Message Queue,简称 MQ),顾名思义,就是存放消息的队列。
打个比方:
- 你去餐厅吃饭,点完菜,服务员把单子放到后厨的队列里
- 厨师按顺序做菜,做好一个出一个
- 你不用一直在厨房等,可以坐着玩手机
- 菜做好了,服务员给你端过来
这就是消息队列的思想:
- 你(生产者)把订单(消息)放到队列里
- 厨师(消费者)从队列里拿订单处理
- 你不用等,该干嘛干嘛(异步)
核心角色:
- 生产者(Producer):发消息的一方
- 消费者(Consumer):收消息、处理消息的一方
- 消息代理(Broker):消息队列服务器(比如 RabbitMQ)
- 队列(Queue):存放消息的地方
8.1.2 为什么要用消息队列
消息队列主要解决三个问题:异步、解耦、削峰。
1. 异步处理
同步的问题:
用户注册,要做三件事:
- 写数据库(50ms)
- 发注册邮件(300ms)
- 发短信(300ms)
同步的话,总共要 650ms,用户要等很久。
用了消息队列之后:
- 写数据库(50ms)
- 发消息到 MQ(10ms)
- 直接返回成功
总共 60ms,用户体验好很多。 邮件和短信由消费者慢慢处理,用户不用等。
同步:注册 → 写库 → 发邮件 → 发短信 → 返回(650ms)
异步:注册 → 写库 → 发MQ → 返回(60ms)
↓
消费者慢慢发邮件、发短信2. 系统解耦
没有 MQ 的时候:
订单系统直接调用库存系统、物流系统、积分系统。
订单系统 → 库存系统
→ 物流系统
→ 积分系统问题:
- 订单系统要知道其他所有系统的接口
- 任何一个系统挂了,订单都受影响
- 加一个新系统,要改订单系统的代码
有了 MQ 之后:
订单系统只管发消息到 MQ,其他系统自己订阅。
订单系统 → MQ ← 库存系统
← 物流系统
← 积分系统
← 以后加新系统,直接订阅就行好处:
- 系统之间不直接依赖
- 加新功能不用改老代码
- 一个系统挂了,不影响其他的
3. 流量削峰
场景:秒杀活动
平时 QPS 100,秒杀的时候突然到 10000。 数据库只能扛 1000,直接就挂了。
用 MQ 削峰:
- 请求都先放到 MQ 里
- 消费者按数据库能承受的速度慢慢处理
- 不会一下子把数据库打挂
10000 请求 → MQ → 消费者(1000/s)→ 数据库就像水库一样,洪水来了先存着,慢慢放。
总结一下消息队列的好处:
- 异步:提高响应速度
- 解耦:系统之间不直接依赖
- 削峰:扛住突发流量
这三个是最核心的,面试经常问。
8.1.3 常用消息中间件对比
市面上有很多消息中间件,最常见的几个:
| 中间件 | 开发语言 | 特点 | 适用场景 |
|---|---|---|---|
| RabbitMQ | Erlang | 功能丰富,可靠性高,社区活跃 | 中小项目、业务消息 |
| Kafka | Scala/Java | 高吞吐,分布式,日志场景 | 大数据、日志收集、流处理 |
| RocketMQ | Java | 阿里开源,功能全,金融级 | 电商、金融、大公司 |
| ActiveMQ | Java | 老牌,功能全,但性能一般 | 传统企业、老项目 |
| Pulsar | Java | 下一代,存算分离 | 新兴,云原生 |
简单对比:
- RabbitMQ:功能最丰富,交换机模式灵活,适合业务消息。性能中等。
- Kafka:吞吐量最高,适合日志、大数据。功能相对简单。
- RocketMQ:阿里出品,功能全,性能好,适合大规模业务。
- ActiveMQ:老牌,用的人越来越少了。
这一章我们学 RabbitMQ。 因为它功能丰富,概念清晰,适合入门。 学会了 RabbitMQ,再学其他的也快。
8.1.4 RabbitMQ 简介
RabbitMQ 是一个开源的消息代理,基于 Erlang 语言开发,实现了 AMQP 协议。
特点:
- 可靠性高:消息持久化、确认机制
- 功能丰富:多种交换机模式、死信队列、延迟队列等
- 社区活跃:文档多,问题好查
- 管理界面:自带 Web 管理界面,方便查看
- 支持多种协议:AMQP、MQTT、STOMP 等
AMQP 是什么?
AMQP(Advanced Message Queuing Protocol)高级消息队列协议,是一个开放的标准协议。
就像 HTTP 是网页的标准协议一样,AMQP 是消息队列的标准协议。 RabbitMQ 就是这个协议的一个实现。
8.2 RabbitMQ 核心概念
学 RabbitMQ,先把这几个核心概念搞明白。
8.2.1 核心角色
| 角色 | 英文 | 作用 |
|---|---|---|
| 生产者 | Producer | 发消息的 |
| 消费者 | Consumer | 收消息、处理消息的 |
| 代理 | Broker | 消息队列服务器(RabbitMQ 本身) |
| 虚拟主机 | Virtual Host | 类似命名空间,隔离不同项目的队列 |
8.2.2 队列、交换机、绑定
这三个是 RabbitMQ 最核心的概念,一定要搞清楚。
队列(Queue)
队列就是存放消息的地方。
- 消息存在队列里
- 消费者从队列里拿消息
- 队列是 FIFO(先进先出)的
- 队列有名字,通过名字来识别
队列就像邮箱,邮件(消息)存在里面,等人来取。
交换机(Exchange)
交换机负责接收生产者发来的消息,然后根据规则路由到队列。
生产者不直接把消息发到队列,而是发到交换机!
打个比方:
- 生产者是寄信的人
- 交换机是邮局
- 队列是收件人的邮箱
- 你把信交给邮局,邮局根据地址把信送到对应的邮箱
- 你不用管信怎么送的,交给邮局就行
为什么要有交换机?直接发队列不行吗?
因为交换机可以实现灵活的路由:
- 一条消息可以发到多个队列
- 可以根据条件路由到不同队列
- 可以实现各种模式(发布订阅、路由、主题等)
绑定(Binding)
绑定就是把队列和交换机关联起来。
绑定的时候可以指定 Routing Key(路由键),交换机根据 Routing Key 来决定消息发到哪个队列。
绑定就像告诉邮局:"地址是北京的信,送到这个邮箱"。
8.2.3 消息流转过程
完整的消息流转过程:
生产者 → 交换机 → (根据 Routing Key 路由)→ 队列 → 消费者- 生产者把消息发给交换机
- 交换机收到消息,根据 Routing Key 和绑定规则,把消息路由到对应的队列
- 消息存在队列里
- 消费者从队列里取消息,处理
记住这个流程,后面学各种模式就好理解了。
8.2.4 交换机的类型
RabbitMQ 有几种不同类型的交换机,对应不同的路由规则:
| 交换机类型 | 英文名 | 路由规则 | 对应模式 |
|---|---|---|---|
| 直连交换机 | Direct | 精确匹配 Routing Key | 简单模式、路由模式 |
| 扇出交换机 | Fanout | 广播到所有绑定的队列 | 发布订阅模式 |
| 主题交换机 | Topic | 通配符匹配 Routing Key | 主题模式 |
| 头交换机 | Headers | 根据消息头匹配 | 头模式(少用) |
后面会详细讲每种模式。
8.2.5 Virtual Host(虚拟主机)
Virtual Host(vhost)是 RabbitMQ 的虚拟主机,用来隔离不同的环境。
- 每个 vhost 有自己的队列、交换机、绑定
- 不同 vhost 之间完全隔离
- 类似数据库的 database 概念
默认 vhost 是 /。
比如一个公司有多个项目,可以每个项目一个 vhost,互不影响。
8.3 RabbitMQ 安装与管理界面
8.3.1 安装 RabbitMQ
RabbitMQ 基于 Erlang,所以要先装 Erlang,再装 RabbitMQ。
Windows:
- 去官网下载 Erlang 和 RabbitMQ 的安装包
- 先装 Erlang,再装 RabbitMQ
- 一路下一步就行
Mac:
brew install rabbitmqLinux(Ubuntu/Debian):
sudo apt install rabbitmq-serverDocker(推荐):
docker run -d \
--name rabbitmq \
-p 5672:5672 \
-p 15672:15672 \
rabbitmq:3-management用 Docker 最方便,不用装环境,一条命令搞定。
3-management版本带管理界面。
8.3.2 管理界面
RabbitMQ 自带 Web 管理界面,非常方便。
访问地址: http://localhost:15672
默认账号密码: guest / guest
注意:guest 用户只能从 localhost 访问,远程访问需要新建用户。
管理界面能做什么:
- 查看连接、通道、队列、交换机
- 新建/删除队列、交换机
- 查看消息数量
- 手动发消息、收消息
- 管理用户和权限
- 查看统计信息
学 RabbitMQ 一定要多玩管理界面,直观易懂。
8.3.3 常用端口
| 端口 | 作用 |
|---|---|
| 5672 | AMQP 协议端口(客户端连接用) |
| 15672 | 管理界面端口 |
| 25672 | 集群通信端口 |
客户端连接用 5672,浏览器访问管理界面用 15672。
8.4 原生 Java 客户端入门
我们先用原生的 Java 客户端(amqp-client)来感受一下 RabbitMQ 的基本用法。
8.4.1 加依赖
<dependency>
<groupId>com.rabbitmq</groupId>
<artifactId>amqp-client</artifactId>
<version>5.6.0</version>
</dependency>8.4.2 HelloWorld 模式(简单模式)
最简单的模式:一个生产者,一个队列,一个消费者。
生产者 → 队列 → 消费者生产者(发消息):
public class Producer_HelloWorld {
public static void main(String[] args) throws IOException, TimeoutException {
// 1. 创建连接工厂
ConnectionFactory factory = new ConnectionFactory();
// 2. 设置连接参数
factory.setHost("localhost"); // 服务器地址
factory.setPort(5672); // 端口
factory.setVirtualHost("/"); // 虚拟主机
factory.setUsername("guest"); // 用户名
factory.setPassword("guest"); // 密码
// 3. 创建连接
Connection connection = factory.newConnection();
// 4. 创建通道(Channel)
Channel channel = connection.createChannel();
// 5. 声明队列(如果不存在就创建,存在就什么都不做)
// 参数:队列名, 是否持久化, 是否排他, 是否自动删除, 其他参数
channel.queueDeclare("hello_world", true, false, false, null);
// 6. 发送消息
String body = "hello rabbitmq~~~";
// 参数:交换机名, 路由键, 其他属性, 消息体
channel.basicPublish("", "hello_world", null, body.getBytes());
System.out.println("消息发送成功");
// 7. 关闭资源
channel.close();
connection.close();
}
}queueDeclare 参数说明:
| 参数 | 说明 |
|---|---|
queue | 队列名称 |
durable | 是否持久化(true=重启不丢) |
exclusive | 是否排他(只有这个连接能用) |
autoDelete | 是否自动删除(没人用了就删) |
arguments | 其他参数 |
basicPublish 参数说明:
| 参数 | 说明 |
|---|---|
exchange | 交换机名(空字符串表示默认交换机) |
routingKey | 路由键(简单模式下就是队列名) |
props | 消息属性 |
body | 消息体(字节数组) |
消费者(收消息):
public class Consumer_HelloWorld {
public static void main(String[] args) throws IOException, TimeoutException {
// 1. 创建连接工厂
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
factory.setPort(5672);
factory.setVirtualHost("/");
factory.setUsername("guest");
factory.setPassword("guest");
// 2. 创建连接
Connection connection = factory.newConnection();
// 3. 创建通道
Channel channel = connection.createChannel();
// 4. 声明队列
channel.queueDeclare("hello_world", true, false, false, null);
// 5. 创建消费者
Consumer consumer = new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag,
Envelope envelope,
AMQP.BasicProperties properties,
byte[] body) throws IOException {
System.out.println("收到消息:" + new String(body));
System.out.println("交换机:" + envelope.getExchange());
System.out.println("路由键:" + envelope.getRoutingKey());
}
};
// 6. 监听队列
// 参数:队列名, 是否自动确认, 消费者
channel.basicConsume("hello_world", true, consumer);
System.out.println("消费者启动,等待消息...");
// 消费者不要关闭连接,要一直监听
}
}测试一下:
- 先启动消费者
- 再运行生产者发消息
- 看看消费者有没有收到
简单模式是最基础的,后面的模式都是在这个基础上变化的。
8.4.3 Work 模式(工作队列)
一个队列,多个消费者。消息轮流分给消费者(轮询)。
生产者 → 队列 → 消费者1
→ 消费者2
→ 消费者3应用场景: 任务处理,多个 worker 分担压力。
特点:
- 一个队列,多个消费者
- 一条消息只会被一个消费者处理
- 默认轮询分发,一个一条
这个模式也很常用,比如后台任务处理。 代码跟简单模式差不多,只是多起几个消费者就行。
8.5 RabbitMQ 工作模式详解
RabbitMQ 有好几种工作模式,我们一个个来看。
8.5.1 模式一:简单模式(HelloWorld)
结构:
P → Q → C说明:
- 一个生产者
- 一个队列
- 一个消费者
- 用默认交换机(Direct 类型)
应用场景: 最简单的场景,一对一发消息。
8.5.2 模式二:工作队列模式(Work Queue)
结构:
P → Q → C1
→ C2
→ C3说明:
- 一个生产者
- 一个队列
- 多个消费者
- 消息轮流分配(轮询)
- 一条消息只被一个消费者处理
应用场景: 任务分发,多个 worker 提高处理速度。
8.5.3 模式三:发布订阅模式(Publish/Subscribe / Fanout)
结构:
→ Q1 → C1
P → X
→ Q2 → C2X 是 Fanout 交换机(扇出交换机)。
说明:
- 生产者把消息发到 Fanout 交换机
- 交换机把消息广播到所有绑定的队列
- 每个队列都收到完整的消息
- 每个队列的消费者各自处理
应用场景:
- 注册成功后,同时发邮件、发短信、加积分
- 一条消息多个系统都要处理
关键点:
- 交换机类型是 fanout
- 路由键没用(因为是广播)
- 所有绑定的队列都能收到
发布订阅模式非常常用,一定要掌握。
8.5.4 模式四:路由模式(Routing / Direct)
结构:
→ Q1(routingKey: error) → C1
P → X
→ Q2(routingKey: info, error, warning) → C2X 是 Direct 交换机(直连交换机)。
说明:
- 生产者发消息的时候指定 Routing Key
- 交换机根据 Routing Key 精确匹配,发到对应的队列
- 一个队列可以绑定多个 Routing Key
应用场景:
- 日志分级:error 级别的存数据库,所有级别的都打印
- 不同级别的消息不同处理
例子:
- Q1 绑定了
error,只收到 error 级别的日志 - Q2 绑定了
info、error、warning,收到所有级别的日志
路由模式也很常用,比发布订阅更灵活。
8.5.5 模式五:主题模式(Topics / Topic)
结构:
→ Q1(pattern: *.email.*) → C1
P → X
→ Q2(pattern: #.sms.#) → C2X 是 Topic 交换机(主题交换机)。
说明:
- 跟路由模式类似,但 Routing Key 可以用通配符
- 更灵活,可以实现复杂的路由规则
通配符规则:
*(星号):匹配一个单词#(井号):匹配 0 个或多个单词- 单词之间用
.分隔
例子:
info.email.admin→ 匹配*.email.*,也匹配info.#info.sms→ 匹配#.sms.#,不匹配*.email.*
应用场景:
- 更灵活的路由
- 订阅不同主题的消息
- 比如:订阅所有邮件相关的,订阅所有 admin 相关的
主题模式是最灵活的,功能最强大。
8.5.6 模式六:头模式(Headers)
根据消息的 Header 来路由,不用 Routing Key。
用得比较少,了解一下就行。
8.5.7 模式对比总结
| 模式 | 交换机类型 | 路由方式 | 特点 |
|---|---|---|---|
| 简单模式 | 默认(Direct) | 精确匹配队列名 | 一对一 |
| 工作队列 | 默认(Direct) | 精确匹配队列名 | 一对多,轮询 |
| 发布订阅 | Fanout | 广播 | 一条消息多队列 |
| 路由模式 | Direct | 精确匹配 Routing Key | 按 key 路由 |
| 主题模式 | Topic | 通配符匹配 Routing Key | 最灵活 |
| 头模式 | Headers | 按 Header 匹配 | 少用 |
8.6 Spring Boot 整合 RabbitMQ(重点)
实际项目中,我们一般不会用原生客户端,而是用 Spring Boot 整合。
Spring Boot 对 RabbitMQ 有很好的支持,用起来非常方便。
8.6.1 环境搭建
步骤 1:加依赖
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>步骤 2:配置 RabbitMQ
spring:
rabbitmq:
host: localhost # 服务器地址
port: 5672 # 端口
virtual-host: / # 虚拟主机
username: guest # 用户名
password: guest # 密码
listener:
simple:
acknowledge-mode: manual # 手动确认(后面讲)
prefetch: 1 # 每次预取 1 条步骤 3:启动类
@SpringBootApplication
public class RabbitmqApplication {
public static void main(String[] args) {
SpringApplication.run(RabbitmqApplication.class, args);
}
}就这么简单,Spring Boot 自动配置好了。
8.6.2 RabbitTemplate:发消息的工具
Spring 提供了 RabbitTemplate 来发送消息,非常方便。
@Service
public class MessageProducerService {
@Autowired
private RabbitTemplate rabbitTemplate;
/**
* 发送消息
*/
public void sendMessage(String exchange, String routingKey, Object message) {
rabbitTemplate.convertAndSend(exchange, routingKey, message);
}
}常用方法:
| 方法 | 说明 |
|---|---|
convertAndSend(routingKey, message) | 发到默认交换机 |
convertAndSend(exchange, routingKey, message) | 发到指定交换机 |
convertAndSend(exchange, routingKey, message, postProcessor) | 带消息处理 |
receiveAndConvert(queueName) | 接收消息(拉模式) |
convertAndSend 会自动把对象转换成消息,不用自己转字节数组。
8.6.3 消息序列化:JSON
默认情况下,RabbitTemplate 用 JDK 序列化,存进去是二进制,看着乱码。
一般我们配置成 JSON 序列化。
配置类:
@Configuration
public class RabbitMQConfig {
/**
* 消息转换器:用 JSON 序列化
*/
@Bean
public MessageConverter messageConverter() {
return new Jackson2JsonMessageConverter();
}
}配置完之后,发对象的时候自动转成 JSON,消费者那边自动转回来。
强烈建议配置 JSON 序列化,调试方便,也支持跨语言。
8.6.4 @RabbitListener:监听消息
Spring Boot 用 @RabbitListener 注解来监听消息,非常简单。
@Service
public class MessageConsumerService {
/**
* 监听队列
*/
@RabbitListener(queues = "hello_world")
public void receiveMessage(String message) {
System.out.println("收到消息:" + message);
}
}就这么简单!加个注解,方法就变成消费者了。
还可以直接接收对象:
@RabbitListener(queues = "user.queue")
public void receiveUser(User user) {
System.out.println("收到用户:" + user);
}自动把 JSON 转成对象,太方便了!
8.6.5 用注解声明队列、交换机、绑定
Spring Boot 还支持用注解直接声明队列、交换机、绑定,不用写配置类。
@RabbitListener(
bindings = @QueueBinding(
value = @Queue("email.queue"), // 队列
exchange = @Exchange(value = "pubsub.exchange", type = "fanout"), // 交换机
key = "" // 路由键
)
)
public void psubConsumerEmail(User user) {
System.out.println("邮件业务收到消息:" + user);
}注解说明:
@QueueBinding:绑定队列和交换机@Queue:声明队列@Exchange:声明交换机key:路由键
用注解的好处:消费者自己声明需要的队列和绑定,代码集中,好维护。 实际项目中推荐这种方式。
8.7 三种模式的 Spring Boot 实现
我们用 Spring Boot 来实现前面讲的三种模式。
8.7.1 发布订阅模式(Fanout)
场景: 用户注册,同时发邮件和短信。
生产者:
@Service
public class MessageProducerService {
@Autowired
private RabbitTemplate rabbitTemplate;
/**
* 发布订阅模式:发消息到 fanout 交换机
*/
public void psubPublisher(User user) {
// 交换机名、路由键(fanout 模式下路由键没用)、消息
rabbitTemplate.convertAndSend("pub/sub.exchange", "", user);
}
}消费者:
@Service
public class MessageConsumerService {
/**
* 邮件消费者
*/
@RabbitListener(bindings = @QueueBinding(
value = @Queue("email.queue"),
exchange = @Exchange(value = "pub/sub.exchange", type = "fanout")
))
public void psubConsumerEmail(User user) {
System.out.println("邮件业务收到消息:" + user);
}
/**
* 短信消费者
*/
@RabbitListener(bindings = @QueueBinding(
value = @Queue("sms.queue"),
exchange = @Exchange(value = "pub/sub.exchange", type = "fanout")
))
public void psubConsumerSms(User user) {
System.out.println("短信业务收到消息:" + user);
}
}测试:
@RestController
public class UserController {
@Autowired
private MessageProducerService messageProducerService;
@RequestMapping("/user/{id}/{username}")
public void register(@PathVariable Integer id,
@PathVariable String username) {
User user = new User(id, username);
messageProducerService.psubPublisher(user);
}
}访问 http://localhost:8080/user/1/zhangsan,两个消费者都会收到消息。
8.7.2 路由模式(Direct)
场景: 日志处理,error 级别的存库,所有级别的都打印。
生产者:
/**
* 路由模式:发消息到 direct 交换机
*/
public void routingPublisher(String level, String msg) {
rabbitTemplate.convertAndSend("routing.exchange", "routingkey." + level, msg);
}消费者:
/**
* 只收 error 级别
*/
@RabbitListener(bindings = @QueueBinding(
value = @Queue("routing.error.queue"),
exchange = @Exchange(value = "routing.exchange", type = "direct"),
key = "routingkey.error"
))
public void routingConsumerError(String message) {
System.out.println("收到 error 级别日志:" + message);
}
/**
* 收 info、error、warning 级别
*/
@RabbitListener(bindings = @QueueBinding(
value = @Queue("routing.all.queue"),
exchange = @Exchange(value = "routing.exchange", type = "direct"),
key = {"routingkey.error", "routingkey.info", "routingkey.warning"}
))
public void routingConsumerAll(String message) {
System.out.println("收到日志:" + message);
}测试:
- 发 error 消息:两个消费者都收到
- 发 info 消息:只有第二个消费者收到
8.7.3 主题模式(Topic)
场景: 订阅不同类型的通知。
生产者:
/**
* 主题模式:发消息到 topic 交换机
*/
public void topicPublisher(String routingkey, String msg) {
rabbitTemplate.convertAndSend("topic.exchange", routingkey, msg);
}消费者:
/**
* 订阅所有邮件相关的
* info.#.email.# 表示:info 开头,中间任意,email,后面任意
*/
@RabbitListener(bindings = @QueueBinding(
value = @Queue("topic.email.queue"),
exchange = @Exchange(value = "topic.exchange", type = "topic"),
key = "info.#.email.#"
))
public void topicConsumerEmail(String message) {
System.out.println("邮件订阅收到消息:" + message);
}
/**
* 订阅所有短信相关的
*/
@RabbitListener(bindings = @QueueBinding(
value = @Queue("topic.sms.queue"),
exchange = @Exchange(value = "topic.exchange", type = "topic"),
key = "info.#.sms.#"
))
public void topicConsumerSms(String message) {
System.out.println("短信订阅收到消息:" + message);
}测试:
info.email.admin→ 邮件消费者收到info.user.sms→ 短信消费者收到info.email.sms→ 两个都收到
8.8 消息可靠性
消息队列虽然好用,但也要考虑可靠性:消息会不会丢?
8.8.1 消息怎么会丢
消息可能在三个地方丢:
- 生产者丢:生产者发消息,但是没到 Broker
- Broker 丢:消息到了 Broker,但是 Broker 挂了,没持久化
- 消费者丢:消费者收到了,但是没处理完就挂了
我们一个个来看怎么解决。
8.8.2 生产者确认(Publisher Confirm)
确保消息从生产者成功发到了 Broker。
两种确认机制:
- Confirm:消息到了交换机,回调确认
- Return:消息到了交换机,但没路由到队列,回调返回
配置:
spring:
rabbitmq:
publisher-confirm-type: correlated # 开启确认
publisher-returns: true # 开启返回代码:
@Service
public class MessageProducerService {
@Autowired
private RabbitTemplate rabbitTemplate;
@PostConstruct
public void init() {
// 确认回调:消息到了交换机
rabbitTemplate.setConfirmCallback((correlationData, ack, cause) -> {
if (ack) {
System.out.println("消息成功到达交换机");
} else {
System.out.println("消息没到达交换机,原因:" + cause);
// 可以做重发、告警等处理
}
});
// 返回回调:消息没路由到队列
rabbitTemplate.setReturnsCallback(returned -> {
System.out.println("消息没路由到队列:" + new String(returned.getMessage().getBody()));
});
}
}重要的消息一定要开生产者确认,保证消息不丢。
8.8.3 消息持久化
确保 Broker 挂了,消息不丢。
要持久化三样东西:
- 交换机持久化:声明交换机时 durable=true
- 队列持久化:声明队列时 durable=true
- 消息持久化:消息的 deliveryMode=2(持久化)
Spring Boot 默认都是持久化的,一般不用特意配置。 但要知道有这么回事。
8.8.4 消费者确认(ACK)
确保消息被消费者成功处理了。
三种确认模式:
| 模式 | 说明 | 特点 |
|---|---|---|
none | 自动确认 | 发出去就确认,可能丢消息 |
auto | 自动确认 | 方法正常返回就确认,抛异常就不确认 |
manual | 手动确认 | 自己调用 basicAck / basicNack |
推荐用 manual(手动确认),最可靠。
配置:
spring:
rabbitmq:
listener:
simple:
acknowledge-mode: manual # 手动确认代码:
@RabbitListener(queues = "hello_world")
public void receiveMessage(String message,
Channel channel,
@Header(AmqpHeaders.DELIVERY_TAG) long tag) throws IOException {
try {
System.out.println("收到消息:" + message);
// 处理业务...
// 手动确认
// 参数:deliveryTag, 是否批量确认
channel.basicAck(tag, false);
} catch (Exception e) {
// 处理失败,拒绝消息
// 参数:deliveryTag, 是否批量, 是否重新入队
channel.basicNack(tag, false, true);
// 或者丢弃:channel.basicReject(tag, false);
}
}三个方法:
basicAck:确认(消息处理成功,删掉)basicNack:拒绝(可以重新入队)basicReject:拒绝(一次一条,跟 Nack 类似)
手动确认是保证消息可靠性的关键。 处理成功才确认,失败了可以重新入队或者丢弃。
8.8.5 死信队列(DLX)
什么是死信?
- 消息被拒绝(basicNack / basicReject)且 requeue=false
- 消息过期了(TTL 到了)
- 队列满了
死信消息可以放到另一个队列里,叫死信队列。
应用场景:
- 失败的消息存起来,后面人工处理
- 延迟队列(利用 TTL + 死信队列)
配置方式:
- 给普通队列设置
x-dead-letter-exchange参数 - 死信交换机绑定死信队列
- 死信消息会自动发到死信队列
死信队列是 RabbitMQ 的高级功能,很实用。 面试可能会问,知道有这么个东西就行。
8.9 企业最佳实践
8.9.1 消息幂等性
什么是幂等性?
- 同一条消息,处理一次和处理多次,结果是一样的
- 不会因为重复消费导致数据错误
为什么要考虑幂等性?
- 网络问题,消息可能重复发送
- 消费者确认失败,消息重新入队
- 各种原因导致消息重复消费
怎么保证幂等性?
- 唯一 ID:每条消息带唯一 ID,处理前查一下有没有处理过
- 数据库唯一约束:利用数据库的唯一索引
- 乐观锁:带版本号
- Redis 去重:用 Redis 的 SETNX
重要的业务一定要考虑幂等性。 不然重复消费可能导致数据错乱。
8.9.2 消息顺序性
消息有序吗?
- 单个队列,单个消费者:有序
- 单个队列,多个消费者:无序(因为轮询)
- 多个队列:各队列之间无序
怎么保证顺序?
- 同一个业务的消息发到同一个队列
- 这个队列只有一个消费者
- 牺牲性能换顺序
大部分场景不需要严格顺序。 真需要的话,按上面的方式做。
8.9.3 消息积压怎么办
消息积压的原因:
- 消费者挂了
- 消费者处理太慢
- 流量突然变大
怎么解决:
- 紧急扩容:增加消费者数量(但要注意队列数量限制)
- 临时队列:把消息先转存,慢慢处理
- 优化消费逻辑:提高处理速度
- 限流:生产者那边限流,少发点
线上要监控队列的消息数量,积压了及时告警。
8.9.4 队列命名规范
建议队列命名有规范,方便管理:
业务名.模块名.功能名.queue例子:
user.register.email.queueorder.create.sms.queuelog.error.db.queue
交换机命名:
业务名.类型.exchange例子:
user.fanout.exchangelog.direct.exchangenotify.topic.exchange
规范的命名,排查问题的时候方便很多。
8.9.5 监控与告警
生产环境一定要监控:
- 队列消息数量
- 消费者数量
- 消息处理速度
- 连接数
- 内存、磁盘使用情况
告警:
- 队列积压超过阈值
- 消费者掉线
- 磁盘快满了
- 内存不足
RabbitMQ 有很多监控工具,比如 Prometheus + Grafana。 生产环境一定要做好监控。
8.10 新手常见问题排查
问题 1:连接不上 RabbitMQ
可能原因:
- RabbitMQ 没启动
- 地址端口不对
- 用户名密码错了
- 防火墙没开
- guest 用户不能远程访问
排查:
- 先确认 RabbitMQ 启动了
- 检查配置的 host、port、username、password
- 远程访问的话,新建用户,不要用 guest
- 检查防火墙端口
问题 2:消息发出去了,但消费者收不到
可能原因:
- 交换机名不对
- 路由键不对
- 队列没绑定到交换机
- 消费者监听的队列名不对
- 消息被其他消费者抢走了
排查:
- 去管理界面看队列里有没有消息
- 检查交换机、队列、绑定是否正确
- 检查路由键拼写
- 确认消费者监听的队列名
问题 3:消息内容是乱码
原因: 默认 JDK 序列化
解决: 配置 Jackson2JsonMessageConverter,用 JSON 序列化
问题 4:消息重复消费
原因:
- 消费者处理完没确认
- 确认的时候网络出问题
- 消息重新入队
解决:
- 做好幂等性
- 检查确认逻辑有没有问题
问题 5:消息积压
原因: 消费者处理太慢,或者消费者挂了
解决:
- 增加消费者
- 优化消费逻辑
- 排查消费者是不是报错了
问题 6:@RabbitListener 不生效
可能原因:
- 类没被 Spring 管理(没加 @Service 等)
- 方法不是 public 的
- 参数类型不对
- 队列名不对
排查:
- 检查类上有没有 @Service / @Component
- 检查方法是不是 public
- 检查参数类型能不能正确转换
- 检查队列名
8.11 本章小结
恭喜你!消息服务这一章学完了!
这一章内容也不少,核心是理解 RabbitMQ 的工作模式和 Spring Boot 整合。
你都学了什么
消息队列基础:
- 什么是消息队列
- 为什么要用(异步、解耦、削峰)
- 常用消息中间件对比
- RabbitMQ 简介
核心概念:
- 生产者、消费者、Broker
- 队列、交换机、绑定
- Routing Key
- Virtual Host
- 交换机类型(Direct、Fanout、Topic、Headers)
安装与管理:
- RabbitMQ 安装
- 管理界面使用
- 常用端口
原生客户端:
- HelloWorld 模式
- Work 模式
工作模式(重点):
- 简单模式
- 工作队列模式
- 发布订阅模式(Fanout)
- 路由模式(Direct)
- 主题模式(Topic)
- 头模式
Spring Boot 整合(重点):
- 环境搭建
- RabbitTemplate
- JSON 消息序列化
- @RabbitListener 注解
- 注解声明队列和绑定
- 三种模式的整合实现
消息可靠性:
- 生产者确认(Confirm / Return)
- 消息持久化
- 消费者确认(ACK)
- 死信队列
企业最佳实践:
- 消息幂等性
- 消息顺序性
- 消息积压处理
- 命名规范
- 监控与告警
动手实践清单
| 实践项 | 做完打勾 |
|---|---|
| 能安装 RabbitMQ,打开管理界面 | ☐ |
| 能用原生客户端实现简单模式 | ☐ |
| 理解五种工作模式的区别 | ☐ |
| 能搭建 Spring Boot + RabbitMQ 环境 | ☐ |
| 会用 RabbitTemplate 发消息 | ☐ |
| 会用 @RabbitListener 监听消息 | ☐ |
| 会配置 JSON 序列化 | ☐ |
| 能实现发布订阅模式 | ☐ |
| 能实现路由模式 | ☐ |
| 能实现主题模式 | ☐ |
| 理解消息可靠性的三个层面 | ☐ |
| 会用手动 ACK | ☐ |
| 知道什么是死信队列 | ☐ |
| 理解消息幂等性 | ☐ |
面试常问
为什么要用消息队列?
- 异步处理:提高响应速度
- 系统解耦:系统之间不直接依赖
- 流量削峰:扛住突发流量
消息队列有什么缺点?
- 系统复杂度增加
- 消息一致性问题
- 运维成本增加
- 有消息丢失、重复、积压等问题
RabbitMQ 有哪些工作模式?
- 简单模式、工作队列模式
- 发布订阅模式(Fanout)
- 路由模式(Direct)
- 主题模式(Topic)
- 头模式(Headers)
RabbitMQ 怎么保证消息不丢?
- 生产者确认:Confirm / Return
- 消息持久化:交换机、队列、消息都持久化
- 消费者确认:手动 ACK
- 死信队列:处理失败的消息
如何保证消息的幂等性?
- 唯一 ID + 去重表
- 数据库唯一约束
- 乐观锁
- Redis 去重
消息积压了怎么办?
- 增加消费者
- 优化消费逻辑
- 临时转存
- 生产者限流
RabbitMQ 和 Kafka 有什么区别?
- RabbitMQ:功能丰富,可靠性高,适合业务消息,吞吐中等
- Kafka:高吞吐,分布式,适合日志、大数据,功能相对简单
- 场景不同,没有绝对的好坏
死信队列是什么?有什么用?
- 死信:被拒绝的、过期的、队列满了的消息
- 死信队列:存放死信的队列
- 用处:失败消息处理、延迟队列等
给实习生的建议
先把核心概念搞明白:队列、交换机、绑定、Routing Key,这几个是基础。
五种模式都要会:尤其是发布订阅、路由、主题这三种,最常用。
Spring Boot 整合是重点:实际项目都用这个,一定要熟练。
消息可靠性很重要:生产环境消息不能丢,ACK、持久化这些要掌握。
幂等性一定要考虑:不要假设消息只会来一次,要按会来多次来设计。
做好监控:队列积压、消费者掉线这些要能及时发现。
不要滥用 MQ:不是什么都要放 MQ,简单的同步调用就够了。MQ 增加了系统复杂度,该用的时候才用。
下一章我们学习任务调度和邮件发送。加油!