Skip to content

第 8 章 Spring Boot 消息服务 ​

写给新手的话:

前面我们学了数据库、缓存、安全管理,都是同步的:发请求 → 处理 → 返回结果。

但是有些场景,同步处理不太合适:

  • 注册完要发邮件、发短信,用户要等半天
  • 秒杀的时候,一下子几万请求过来,数据库扛不住
  • 系统之间调用,一个挂了其他的也受影响

这时候就需要消息队列了!

消息队列可以帮我们:异步处理、系统解耦、流量削峰。

这一章我们学最流行的消息中间件之一:RabbitMQ。

学习建议:

  • 先理解核心概念:队列、交换机、生产者、消费者
  • 搞清楚几种工作模式的区别和应用场景
  • 多动手,把每种模式都跑一遍
  • Spring Boot 整合是重点,实际项目中用得最多

准备好了吗?我们开始吧!


8.1 先搞明白:消息队列那些事儿 ​

8.1.1 什么是消息队列 ​

消息队列(Message Queue,简称 MQ),顾名思义,就是存放消息的队列。

打个比方:

  • 你去餐厅吃饭,点完菜,服务员把单子放到后厨的队列里
  • 厨师按顺序做菜,做好一个出一个
  • 你不用一直在厨房等,可以坐着玩手机
  • 菜做好了,服务员给你端过来

这就是消息队列的思想:

  • 你(生产者)把订单(消息)放到队列里
  • 厨师(消费者)从队列里拿订单处理
  • 你不用等,该干嘛干嘛(异步)

核心角色:

  • 生产者(Producer):发消息的一方
  • 消费者(Consumer):收消息、处理消息的一方
  • 消息代理(Broker):消息队列服务器(比如 RabbitMQ)
  • 队列(Queue):存放消息的地方

8.1.2 为什么要用消息队列 ​

消息队列主要解决三个问题:异步、解耦、削峰。

1. 异步处理 ​

同步的问题:

用户注册,要做三件事:

  1. 写数据库(50ms)
  2. 发注册邮件(300ms)
  3. 发短信(300ms)

同步的话,总共要 650ms,用户要等很久。

用了消息队列之后:

  1. 写数据库(50ms)
  2. 发消息到 MQ(10ms)
  3. 直接返回成功

总共 60ms,用户体验好很多。 邮件和短信由消费者慢慢处理,用户不用等。

同步:注册 → 写库 → 发邮件 → 发短信 → 返回(650ms)

异步:注册 → 写库 → 发MQ → 返回(60ms)
                    ↓
                消费者慢慢发邮件、发短信

2. 系统解耦 ​

没有 MQ 的时候:

订单系统直接调用库存系统、物流系统、积分系统。

订单系统 → 库存系统
        → 物流系统
        → 积分系统

问题:

  • 订单系统要知道其他所有系统的接口
  • 任何一个系统挂了,订单都受影响
  • 加一个新系统,要改订单系统的代码

有了 MQ 之后:

订单系统只管发消息到 MQ,其他系统自己订阅。

订单系统 → MQ ← 库存系统
            ← 物流系统
            ← 积分系统
            ← 以后加新系统,直接订阅就行

好处:

  • 系统之间不直接依赖
  • 加新功能不用改老代码
  • 一个系统挂了,不影响其他的

3. 流量削峰 ​

场景:秒杀活动

平时 QPS 100,秒杀的时候突然到 10000。 数据库只能扛 1000,直接就挂了。

用 MQ 削峰:

  • 请求都先放到 MQ 里
  • 消费者按数据库能承受的速度慢慢处理
  • 不会一下子把数据库打挂
10000 请求 → MQ → 消费者(1000/s)→ 数据库

就像水库一样,洪水来了先存着,慢慢放。

总结一下消息队列的好处:

  • 异步:提高响应速度
  • 解耦:系统之间不直接依赖
  • 削峰:扛住突发流量

这三个是最核心的,面试经常问。

8.1.3 常用消息中间件对比 ​

市面上有很多消息中间件,最常见的几个:

中间件开发语言特点适用场景
RabbitMQErlang功能丰富,可靠性高,社区活跃中小项目、业务消息
KafkaScala/Java高吞吐,分布式,日志场景大数据、日志收集、流处理
RocketMQJava阿里开源,功能全,金融级电商、金融、大公司
ActiveMQJava老牌,功能全,但性能一般传统企业、老项目
PulsarJava下一代,存算分离新兴,云原生

简单对比:

  • 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 路由)→ 队列 → 消费者
  1. 生产者把消息发给交换机
  2. 交换机收到消息,根据 Routing Key 和绑定规则,把消息路由到对应的队列
  3. 消息存在队列里
  4. 消费者从队列里取消息,处理

记住这个流程,后面学各种模式就好理解了。

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:

bash
brew install rabbitmq

Linux(Ubuntu/Debian):

bash
sudo apt install rabbitmq-server

Docker(推荐):

bash
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 常用端口 ​

端口作用
5672AMQP 协议端口(客户端连接用)
15672管理界面端口
25672集群通信端口

客户端连接用 5672,浏览器访问管理界面用 15672。


8.4 原生 Java 客户端入门 ​

我们先用原生的 Java 客户端(amqp-client)来感受一下 RabbitMQ 的基本用法。

8.4.1 加依赖 ​

xml
<dependency>
    <groupId>com.rabbitmq</groupId>
    <artifactId>amqp-client</artifactId>
    <version>5.6.0</version>
</dependency>

8.4.2 HelloWorld 模式(简单模式) ​

最简单的模式:一个生产者,一个队列,一个消费者。

生产者 → 队列 → 消费者

生产者(发消息):

java
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消息体(字节数组)

消费者(收消息):

java
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("消费者启动,等待消息...");
        // 消费者不要关闭连接,要一直监听
    }
}

测试一下:

  1. 先启动消费者
  2. 再运行生产者发消息
  3. 看看消费者有没有收到

简单模式是最基础的,后面的模式都是在这个基础上变化的。

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 → C2

X 是 Fanout 交换机(扇出交换机)。

说明:

  • 生产者把消息发到 Fanout 交换机
  • 交换机把消息广播到所有绑定的队列
  • 每个队列都收到完整的消息
  • 每个队列的消费者各自处理

应用场景:

  • 注册成功后,同时发邮件、发短信、加积分
  • 一条消息多个系统都要处理

关键点:

  • 交换机类型是 fanout
  • 路由键没用(因为是广播)
  • 所有绑定的队列都能收到

发布订阅模式非常常用,一定要掌握。

8.5.4 模式四:路由模式(Routing / Direct) ​

结构:

      → Q1(routingKey: error) → C1
P → X
      → Q2(routingKey: info, error, warning) → C2

X 是 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.#) → C2

X 是 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:加依赖

xml
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-amqp</artifactId>
</dependency>

步骤 2:配置 RabbitMQ

yaml
spring:
  rabbitmq:
    host: localhost        # 服务器地址
    port: 5672             # 端口
    virtual-host: /        # 虚拟主机
    username: guest        # 用户名
    password: guest        # 密码
    listener:
      simple:
        acknowledge-mode: manual  # 手动确认(后面讲)
        prefetch: 1               # 每次预取 1 条

步骤 3:启动类

java
@SpringBootApplication
public class RabbitmqApplication {
    public static void main(String[] args) {
        SpringApplication.run(RabbitmqApplication.class, args);
    }
}

就这么简单,Spring Boot 自动配置好了。

8.6.2 RabbitTemplate:发消息的工具 ​

Spring 提供了 RabbitTemplate 来发送消息,非常方便。

java
@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 序列化。

配置类:

java
@Configuration
public class RabbitMQConfig {

    /**
     * 消息转换器:用 JSON 序列化
     */
    @Bean
    public MessageConverter messageConverter() {
        return new Jackson2JsonMessageConverter();
    }
}

配置完之后,发对象的时候自动转成 JSON,消费者那边自动转回来。

强烈建议配置 JSON 序列化,调试方便,也支持跨语言。

8.6.4 @RabbitListener:监听消息 ​

Spring Boot 用 @RabbitListener 注解来监听消息,非常简单。

java
@Service
public class MessageConsumerService {

    /**
     * 监听队列
     */
    @RabbitListener(queues = "hello_world")
    public void receiveMessage(String message) {
        System.out.println("收到消息:" + message);
    }
}

就这么简单!加个注解,方法就变成消费者了。

还可以直接接收对象:

java
@RabbitListener(queues = "user.queue")
public void receiveUser(User user) {
    System.out.println("收到用户:" + user);
}

自动把 JSON 转成对象,太方便了!

8.6.5 用注解声明队列、交换机、绑定 ​

Spring Boot 还支持用注解直接声明队列、交换机、绑定,不用写配置类。

java
@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) ​

场景: 用户注册,同时发邮件和短信。

生产者:

java
@Service
public class MessageProducerService {

    @Autowired
    private RabbitTemplate rabbitTemplate;

    /**
     * 发布订阅模式:发消息到 fanout 交换机
     */
    public void psubPublisher(User user) {
        // 交换机名、路由键(fanout 模式下路由键没用)、消息
        rabbitTemplate.convertAndSend("pub/sub.exchange", "", user);
    }
}

消费者:

java
@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);
    }
}

测试:

java
@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 级别的存库,所有级别的都打印。

生产者:

java
/**
 * 路由模式:发消息到 direct 交换机
 */
public void routingPublisher(String level, String msg) {
    rabbitTemplate.convertAndSend("routing.exchange", "routingkey." + level, msg);
}

消费者:

java
/**
 * 只收 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) ​

场景: 订阅不同类型的通知。

生产者:

java
/**
 * 主题模式:发消息到 topic 交换机
 */
public void topicPublisher(String routingkey, String msg) {
    rabbitTemplate.convertAndSend("topic.exchange", routingkey, msg);
}

消费者:

java
/**
 * 订阅所有邮件相关的
 * 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 消息怎么会丢 ​

消息可能在三个地方丢:

  1. 生产者丢:生产者发消息,但是没到 Broker
  2. Broker 丢:消息到了 Broker,但是 Broker 挂了,没持久化
  3. 消费者丢:消费者收到了,但是没处理完就挂了

我们一个个来看怎么解决。

8.8.2 生产者确认(Publisher Confirm) ​

确保消息从生产者成功发到了 Broker。

两种确认机制:

  • Confirm:消息到了交换机,回调确认
  • Return:消息到了交换机,但没路由到队列,回调返回

配置:

yaml
spring:
  rabbitmq:
    publisher-confirm-type: correlated  # 开启确认
    publisher-returns: true             # 开启返回

代码:

java
@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 挂了,消息不丢。

要持久化三样东西:

  1. 交换机持久化:声明交换机时 durable=true
  2. 队列持久化:声明队列时 durable=true
  3. 消息持久化:消息的 deliveryMode=2(持久化)

Spring Boot 默认都是持久化的,一般不用特意配置。 但要知道有这么回事。

8.8.4 消费者确认(ACK) ​

确保消息被消费者成功处理了。

三种确认模式:

模式说明特点
none自动确认发出去就确认,可能丢消息
auto自动确认方法正常返回就确认,抛异常就不确认
manual手动确认自己调用 basicAck / basicNack

推荐用 manual(手动确认),最可靠。

配置:

yaml
spring:
  rabbitmq:
    listener:
      simple:
        acknowledge-mode: manual  # 手动确认

代码:

java
@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 消息幂等性 ​

什么是幂等性?

  • 同一条消息,处理一次和处理多次,结果是一样的
  • 不会因为重复消费导致数据错误

为什么要考虑幂等性?

  • 网络问题,消息可能重复发送
  • 消费者确认失败,消息重新入队
  • 各种原因导致消息重复消费

怎么保证幂等性?

  1. 唯一 ID:每条消息带唯一 ID,处理前查一下有没有处理过
  2. 数据库唯一约束:利用数据库的唯一索引
  3. 乐观锁:带版本号
  4. Redis 去重:用 Redis 的 SETNX

重要的业务一定要考虑幂等性。 不然重复消费可能导致数据错乱。

8.9.2 消息顺序性 ​

消息有序吗?

  • 单个队列,单个消费者:有序
  • 单个队列,多个消费者:无序(因为轮询)
  • 多个队列:各队列之间无序

怎么保证顺序?

  • 同一个业务的消息发到同一个队列
  • 这个队列只有一个消费者
  • 牺牲性能换顺序

大部分场景不需要严格顺序。 真需要的话,按上面的方式做。

8.9.3 消息积压怎么办 ​

消息积压的原因:

  • 消费者挂了
  • 消费者处理太慢
  • 流量突然变大

怎么解决:

  1. 紧急扩容:增加消费者数量(但要注意队列数量限制)
  2. 临时队列:把消息先转存,慢慢处理
  3. 优化消费逻辑:提高处理速度
  4. 限流:生产者那边限流,少发点

线上要监控队列的消息数量,积压了及时告警。

8.9.4 队列命名规范 ​

建议队列命名有规范,方便管理:

业务名.模块名.功能名.queue

例子:

  • user.register.email.queue
  • order.create.sms.queue
  • log.error.db.queue

交换机命名:

业务名.类型.exchange

例子:

  • user.fanout.exchange
  • log.direct.exchange
  • notify.topic.exchange

规范的命名,排查问题的时候方便很多。

8.9.5 监控与告警 ​

生产环境一定要监控:

  • 队列消息数量
  • 消费者数量
  • 消息处理速度
  • 连接数
  • 内存、磁盘使用情况

告警:

  • 队列积压超过阈值
  • 消费者掉线
  • 磁盘快满了
  • 内存不足

RabbitMQ 有很多监控工具,比如 Prometheus + Grafana。 生产环境一定要做好监控。


8.10 新手常见问题排查 ​

问题 1:连接不上 RabbitMQ ​

可能原因:

  1. RabbitMQ 没启动
  2. 地址端口不对
  3. 用户名密码错了
  4. 防火墙没开
  5. guest 用户不能远程访问

排查:

  • 先确认 RabbitMQ 启动了
  • 检查配置的 host、port、username、password
  • 远程访问的话,新建用户,不要用 guest
  • 检查防火墙端口

问题 2:消息发出去了,但消费者收不到 ​

可能原因:

  1. 交换机名不对
  2. 路由键不对
  3. 队列没绑定到交换机
  4. 消费者监听的队列名不对
  5. 消息被其他消费者抢走了

排查:

  • 去管理界面看队列里有没有消息
  • 检查交换机、队列、绑定是否正确
  • 检查路由键拼写
  • 确认消费者监听的队列名

问题 3:消息内容是乱码 ​

原因: 默认 JDK 序列化

解决: 配置 Jackson2JsonMessageConverter,用 JSON 序列化

问题 4:消息重复消费 ​

原因:

  • 消费者处理完没确认
  • 确认的时候网络出问题
  • 消息重新入队

解决:

  • 做好幂等性
  • 检查确认逻辑有没有问题

问题 5:消息积压 ​

原因: 消费者处理太慢,或者消费者挂了

解决:

  • 增加消费者
  • 优化消费逻辑
  • 排查消费者是不是报错了

问题 6:@RabbitListener 不生效 ​

可能原因:

  1. 类没被 Spring 管理(没加 @Service 等)
  2. 方法不是 public 的
  3. 参数类型不对
  4. 队列名不对

排查:

  • 检查类上有没有 @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☐
知道什么是死信队列☐
理解消息幂等性☐

面试常问 ​

  1. 为什么要用消息队列?

    • 异步处理:提高响应速度
    • 系统解耦:系统之间不直接依赖
    • 流量削峰:扛住突发流量
  2. 消息队列有什么缺点?

    • 系统复杂度增加
    • 消息一致性问题
    • 运维成本增加
    • 有消息丢失、重复、积压等问题
  3. RabbitMQ 有哪些工作模式?

    • 简单模式、工作队列模式
    • 发布订阅模式(Fanout)
    • 路由模式(Direct)
    • 主题模式(Topic)
    • 头模式(Headers)
  4. RabbitMQ 怎么保证消息不丢?

    • 生产者确认:Confirm / Return
    • 消息持久化:交换机、队列、消息都持久化
    • 消费者确认:手动 ACK
    • 死信队列:处理失败的消息
  5. 如何保证消息的幂等性?

    • 唯一 ID + 去重表
    • 数据库唯一约束
    • 乐观锁
    • Redis 去重
  6. 消息积压了怎么办?

    • 增加消费者
    • 优化消费逻辑
    • 临时转存
    • 生产者限流
  7. RabbitMQ 和 Kafka 有什么区别?

    • RabbitMQ:功能丰富,可靠性高,适合业务消息,吞吐中等
    • Kafka:高吞吐,分布式,适合日志、大数据,功能相对简单
    • 场景不同,没有绝对的好坏
  8. 死信队列是什么?有什么用?

    • 死信:被拒绝的、过期的、队列满了的消息
    • 死信队列:存放死信的队列
    • 用处:失败消息处理、延迟队列等

给实习生的建议 ​

  1. 先把核心概念搞明白:队列、交换机、绑定、Routing Key,这几个是基础。

  2. 五种模式都要会:尤其是发布订阅、路由、主题这三种,最常用。

  3. Spring Boot 整合是重点:实际项目都用这个,一定要熟练。

  4. 消息可靠性很重要:生产环境消息不能丢,ACK、持久化这些要掌握。

  5. 幂等性一定要考虑:不要假设消息只会来一次,要按会来多次来设计。

  6. 做好监控:队列积压、消费者掉线这些要能及时发现。

  7. 不要滥用 MQ:不是什么都要放 MQ,简单的同步调用就够了。MQ 增加了系统复杂度,该用的时候才用。

下一章我们学习任务调度和邮件发送。加油!


基于 Vite 强力驱动 | 纯静态轻量托管