欢迎光临
我们一直在努力

Spring Cloud + RabbitMQ 在 ruoyi-cloud 中的保姆级实战(含 Docker Compose 部署与工具模块新建)

说明:
这篇文章主要是我在折腾 ruoyi-cloud 集成 RabbitMQ 过程中做的一些记录
我本身也还在学习,难免会有理解不够到位的地方,如果哪位大佬发现问题,欢迎指出,一起交流、一块进步文中涉及的集成代码建议直接跳到第 7 部分查看,会更直观

一、RabbitMQ 是什么? 

RabbitMQ 是一个消息队列系统,本质是一个 “消息中转站”

它负责: 

  • 接收消息 

  • 存储消息 

  • 分发消息 

就像快递中转站一样,把消息可靠地从一个系统转到另一个系统


二、RabbitMQ 能解决什么问题? 

1)削峰(扛高并发) 

高峰期请求很多,系统处理不过来? 
→ 先丢到 MQ 排队 

系统压力瞬间降低

2)解耦(模块不互相依赖) 

下单后: 

  • 发送短信 

  • 扣库存 

  • 写日志 

  • 加积分 

如果全部写在一起会非常乱

用 MQ: 

  • 下单系统发一条“订单创建”的消息 

  • 后续系统都自己监听自己要的消息 

互不干扰 自己做自己的事情

3)异步处理 

短信、推送、记录日志这种不需要立即完成的任务 丢 MQ 异步处理 接口响应更快

4)防止消息丢失 

RabbitMQ 有各种确认机制,能确保: 

  • 生产者发到了 

  • MQ 收到了 

  • 消费者处理成功了 

保证 至少一次送达 


三、RabbitMQ 的结构 

1)Exchange(交换机) 

决定消息如何路由到队列

常用的: 

  • direct(精确匹配 routing key) 

  • topic(模糊匹配) 

  • fanout(广播) 

2)Queue(队列) 

消息最终存储的位置

3)Routing Key(路由键) 

交换机根据它把消息放到正确的队列

4)Binding(绑定)

队列 + 交换机 + routing key 的关系


 四、RabbitMQ 的工作流程

流程图

1. 生产者发送消息到 Exchange(带 routingKey) 
2. Exchange 根据 routingKey 找到对应的 Queue 
3. 消息进入 Queue(等待处理) 
4. 消费者监听 Queue,取出消息执行处理 
5. 消费成功 → ack 确认 
6. 未 ack → MQ 会重新投递

RabbitMQ 的核心流程只有这 4 个: 生产者 → 交换机 → 队列 → 消费者

类比

像快递一样:

组件 

类比 

作用 

Producer 

寄快递的人 

发送消息 

Exchange 

中转站 

判断消息送去哪 

Queue 

快递货架 

存消息 

Consumer 

收快递的人 

处理消息 


五、RabbitMQ 的运行原理

1. 交换机不是用来存消息的 

它只负责把消息分发到正确的队列

2. 真正存消息的是队列 

队列就像一个“排队的数组”

3. 消费者监听队列 

只要有消息进来,它就会被自动触发

4. 消费者必须 ack 

代表我处理完了

如果不 ack: 

  • MQ 会重发 

  • 或者丢到死信队列(DLX) 

5. MQ 会持久化消息 

重启后消息不会丢(前提是 配置了durable=true) 

6. MQ 支持多消费者 

多个消费端可以 抢任务,提升处理速度


六、交换机

交换机的作用就一个

决定消息要送到哪些队列(按什么规则分发)

RabbitMQ 支持多种交换机类型,就是不同的分发规则

1)Direct Exchange(精确匹配) 

核心特征:routing key 必须精确一致

绑定: 

  • 队列绑定 key:system.info 

  • 队列绑定 key:system.warn 

如果发送: 

  • routing key = system.info → 只能到 info 队列 

  • routing key = system.warn → 只能到 warn 队列 

错一个字母都收不到

适用场景 

  • 一个消息只让某个队列接收 

  • 普通业务模块分发 

  • 简单路由 

优点 

  • 简单 

  • 可控 

  • 精准送达 

2)Topic Exchange(模糊匹配) 

核心特征:支持通配符 * 和 #

通配符 

含义 

匹配一个单词 

匹配零个或多个单词 

 类比 

搜索“order.*” 
它能匹配 

  • order.create 

  • order.cancel 

  • order.refund 

如果订阅 “user.#” 
它能匹配 

  • user 

  • user.add 

  • user.delete.phone 

  • user.update.email.name……任何开头是 user 的

示例 

绑定: 

  • 队列 A:order.* → 只收订单一级事件 

  • 队列 B:order.# → 收所有订单相关事件 

  • 队列 C:*.refund → 收所有退款事件 

发送: 

  • order.create → A、B 收到 

  • order.refund → A、B、C 收到 

  • order.refund.wx → B、C 收到 

适用场景 

  • 大型系统 

  • 复杂业务分类 

  • 多级别业务事件 

优点 

灵活强大,非常适合微服务

3)Fanout Exchange(广播) 

核心特征:不看 routing key,所有绑定队列都收到

 类比 

喇叭广播: 

“各位注意,明天停电!” 

无论你是谁,都听得到

 示例 

exchange 绑定了 3 个队列: 

  • system.log 

  • system.user 

  • system.cache 

发消息后 这 3 个都能收到

 适用场景 

  • 通知全体服务 

  • 刷新缓存 

  • 清除本地缓存 

  • 配置更新 

  • 广播类消息(如 WebSocket 大范围推送) 

优点 

简单粗暴、保证所有队列都收到


4)三者总结对比 

类型 

分发方式 

routing key 

场景 

direct 

精确匹配 

必须完全一致 

普通业务路由 

topic 

模糊匹配 

支持 *、# 

大型系统分类 

fanout 

广播 

忽略 

刷缓存、通知所有人 

七、开始实操

1)下载rabbitMq

我这边还是使用的是docker-compose下载的rabbitMq

version: "3.8"
services:

#RabbitMQ
rabbitmq:
image: rabbitmq:3.12-management
container_name: rabbitmq
restart: always
ports:
– "5672:5672" # 应用连接端口
– "15672:15672" # Web控制台端口
environment:
RABBITMQ_DEFAULT_USER: admin
RABBITMQ_DEFAULT_PASS: admin
volumes:
– ./rabbitMQ-data:/var/lib/rabbitmq
networks:
– config_network

networks:
config_network:
driver: bridge

2)新建一份工具模块

我是直接copy原来ruoyi自带的redis模块

copy出来以后改个名字 我这边叫做**-common-rabbitMQ看个人习惯怎么取名 没有这么多讲究

然后我是把这里的所有文件都删除了只留下了

这些个文件:配置类 常量类 (记得改名字 我这边是从redis改成了rabbit)

记得改一下这里初始化的文件内容,改成后面我们写的配置类:

替换成

com.test.common.rabbit.configure.RabbitMQConfig

打开common父pom

找到这个pom 打开 ,打开后可以看到这个redis是怎么放在这里的 我们ctrl+c  +v一波:

再打开最顶部的父pom

打开后看到这样的内容我们依旧是找到redis使用一波ctrl c v

接着我们来到需要使用rabbitmq的模块打开pom文件贴入我们的模块

<dependency>
<groupId>com.ruoyi</groupId>
<artifactId>ruoyi-common-rabbitMQ</artifactId>
</dependency>

一个可以复用的模块差不多就完成了


2)导入依赖

<!– RabbitMQ –>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>

<!– JSON 序列化 –>
<dependency>
<groupId>com.fasterxml.jackson.core</groupId>
<artifactId>jackson-databind</artifactId>
</dependency>

<!– Lombok –>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<optional>true</optional>
</dependency>

3)配置application.yml

在需要集成进去的模块中配置文件如下

spring:
rabbitmq:
host: 127.0.0.1
port: 5672
username: admin
password: admin
virtual-host: /
publisher-confirm-type: correlated
publisher-returns: true
template:
mandatory: true

 virtual-host: /   

虚拟主机(vhost)
RabbitMQ 可以通过 vhost 隔离不同项目的交换机、队列
类似 MySQL 的不同数据库

publisher-confirm-type: correlated

这是 生产者确认机制(Confirm 确认)
用来保证:

消息有没有成功到达 RabbitMQ 的 交换机(Exchange)

correlated

启用 回调,可以在代码中写

rabbitTemplate.setConfirmCallback((correlationData, ack, cause) -> {

});

  • ack = true → 消息成功到达交换机

  • ack = false → 消息没到交换机(MQ 可能挂了)

这是消息可靠传递的第一层保障

publisher-returns: true

这是配合 mandatory 用的

它负责:

消息到达交换机了,但交换机无法把消息路由到队列时触发回调

例如交换机没绑定队列,routingKey 写错等

回调方法:

rabbitTemplate.setReturnsCallback(returned -> { … });

template.mandatory: true

默认情况下,如果消息无法路由,RabbitMQ 会 直接丢弃消息

mandatory = true 之后:

不丢!走 ReturnCallback 回调!

 避免消息无声无息消失了 而这套组合就是 RabbitMQ 的可靠消息机制;

配置作用
publisher-confirm-type = correlated 保证消息能否到达交换机
publisher-returns = true 保证消息能否到达队列
mandatory = true 消息无法路由时一定触发 returns

这段配置开启了:

  • 连接 RabbitMQ 的基础配置

  • 消息确认(生产者确认)

  • 消息路由失败回调

  • 消息不丢失机制

  • RabbitMQ 的生产端可靠性三段式:

  • 发送成功 → 交换机 Confirm 回调

  • 交换机能否路由到队列 → Return 回调

  • 保证消息不会被悄悄吞掉 → mandatory


  • 4)编写配置类

    1.RabbitConstant.java

    public class RabbitConstant {
    public static final String EXCHANGE_SYSTEM = "system.exchange";
    public static final String QUEUE_SYSTEM_MESSAGE = "system.message.queue";
    public static final String ROUTING_SYSTEM_MESSAGE = "system.message.route";
    }

    这就是个放置常量的地方 后续业务上来可以将一些值(队列 / route key / 交换机)的常量放在这里统一管理;


    2.RabbitMQConfig.java

    这是个配置类

    import com.ruoyi.common.rabbit.constant.RabbitConstant;
    import lombok.extern.slf4j.Slf4j;
    import org.springframework.amqp.core.Binding;
    import org.springframework.amqp.core.BindingBuilder;
    import org.springframework.amqp.core.Queue;
    import org.springframework.amqp.core.TopicExchange;
    import org.springframework.amqp.rabbit.connection.ConnectionFactory;
    import org.springframework.amqp.rabbit.core.RabbitTemplate;
    import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter;
    import org.springframework.context.annotation.Bean;
    import org.springframework.context.annotation.Configuration;

    @Slf4j
    @Configuration
    public class RabbitMQConfig {

    /**
    * 队列声明
    */
    @Bean
    public Queue systemMessageQueue() {
    return new Queue(RabbitConstant.QUEUE_SYSTEM_MESSAGE, true);
    }

    /**
    * 交换机
    */
    @Bean
    public TopicExchange systemExchange() {
    return new TopicExchange(RabbitConstant.EXCHANGE_SYSTEM, true, false);
    }

    /**
    * 绑定
    */
    @Bean
    public Binding systemBinding() {
    return BindingBuilder.bind(systemMessageQueue())
    .to(systemExchange())
    .with(RabbitConstant.ROUTING_SYSTEM_MESSAGE);
    }

    /**
    * JSON 序列化
    */
    @Bean
    public Jackson2JsonMessageConverter messageConverter() {
    return new Jackson2JsonMessageConverter();
    }

    /**
    * RabbitTemplate
    */
    @Bean
    public RabbitTemplate rabbitTemplate(ConnectionFactory connectionFactory) {
    RabbitTemplate template = new RabbitTemplate(connectionFactory);
    template.setMessageConverter(messageConverter());

    // success/fail at exchange
    template.setConfirmCallback((correlationData, ack, cause) -> {
    if (!ack) {
    log.error("交换机未接收消息:{}", cause);
    }
    });

    //路由失败
    template.setReturnsCallback(returnedMessage ->
    log.error("路由失败:{}", returnedMessage)
    );

    // ReturnCallback 必须
    template.setMandatory(true);

    return template;
    }
    }


    ① 声明队列

    @Bean
    public Queue systemMessageQueue() {
    return new Queue(RabbitConstant.QUEUE_SYSTEM_MESSAGE, true);
    }

    含义:

    • 创建一个 持久化队列(durable = true)

    • 队列名称从 RabbitConstant 常量里来(统一管理)

    作用:让 MQ 里有一条可用的队列


    ② 声明交换机

    @Bean
    public TopicExchange systemExchange() {
    return new TopicExchange(RabbitConstant.EXCHANGE_SYSTEM, true, false);
    }

    参数含义:

    • true = durable,交换机持久化

    • false = autoDelete,关闭自动删除

    作用:创建用于消息分发的 Topic 类型交换机


    ③ 队列绑定到交换机

    指定 routingKey

    @Bean
    public Binding systemBinding() {
    return BindingBuilder.bind(systemMessageQueue())
    .to(systemExchange())
    .with(RabbitConstant.ROUTING_SYSTEM_MESSAGE);
    }

    意义:

    • 绑定队列

    • 指定交换机

    • 指定 routingKey

    作用:消息发送时指定 routingKey → 交换机 → 队列

    没有这一步,你发消息时第一件事就是丢失路由


    ④ 消息序列化为 JSON 

    @Bean
    public Jackson2JsonMessageConverter messageConverter() {
    return new Jackson2JsonMessageConverter();
    }

    默认 RabbitMQ 用的是:

    Java 序列化(字节串)

    容易出错,也不通用

    换成 JSON:

    • 可读性强

    • 前端/后端/其他服务都能用

    • 解析方便

    作用:让发送和接收消息都走 JSON 格式


    ⑤ 配置 RabbitTemplate 

    这是保证“消息可靠投递”的核心代码

    (1) 设置 JSON 转换器

    template.setMessageConverter(messageConverter());

    保证发出去的消息已经是 JSON,而不是 ObjectOutputStream 的序列化


    (2) ConfirmCallback 

    消息到达交换机是否成功

    template.setConfirmCallback((correlationData, ack, cause) -> {
    if (!ack) {
    log.error("交换机未接收消息:{}", cause);
    }
    });

    触发条件:

    • 连接成功后,发送消息 → 交换机

    • MQ 返回 ack = true 或 ack = false

    作用:

    • ack = true → 投递到交换机成功

    • ack = false → MQ 拒收、网络问题、交换机不存在…

    你能第一时间知道 消息连交换机都到不了   


    (3) ReturnCallback

    交换机到队列失败

    template.setReturnsCallback(returnedMessage ->
    log.error("路由失败:{}", returnedMessage)
    );

    触发场景:

    • 交换机收到了消息

    • 但是 routingKey 不匹配任何队列

    典型问题:

    • routingKey 写错

    • 队列没绑定

    • 交换机类型不匹配

    ReturnCallback 会把失败原因详细打出来


    (4) mandatory = true

    必须执行 ReturnCallback

    template.setMandatory(true);

    这行非常重要!

    否则:

    路由失败的消息默认会悄悄被 RabbitMQ 丢弃!!

    mandatory = true 保证:

    • 消息无法投递时,你必须收到回调通知

    • 不会丢

    这就是一个标准的 带消息确认机制 的 RabbitMQ 模块基础配置

    组件作用
    Queue 创建队列
    Exchange 创建交换机
    Binding 绑定 routingKey
    JacksonConverter 消息 JSON 化
    ConfirmCallback 保障消息到达交换机
    ReturnCallback 保障消息路由队列成功
    mandatory=true 保障失败不丢消息

    八 、测试

    我们打开需要引入的模块

    1)编写一个接收消息类:

    @Component
    public class MsgListener {

    @RabbitListener(queues = RabbitConstant.QUEUE_SYSTEM_MESSAGE)
    public void receive(String msg) {
    System.err.println("【消费者收到消息】" + msg);
    }
    }

    2)编写一个service

    调用模拟我们生产者发送消息

    @Service
    public class MsgTestService {

    @Resource
    private RabbitTemplate rabbitTemplate;

    public void testSend() {
    rabbitTemplate.convertAndSend(
    RabbitConstant.EXCHANGE_SYSTEM,
    RabbitConstant.ROUTING_SYSTEM_MESSAGE,
    "Hello MQ!"
    );
    }
    }

    3)编写一个controller层开始测试

    (我是有现成的直接在这里测试一样的 你有也可以)

    @RestController
    @RequestMapping("/test")
    @Slf4j
    @RequiredArgsConstructor
    public class TestController extends BaseController {
    private final MsgTestService msgTestService;

    @GetMapping("/mq")
    public AjaxResult rabbit() {
    msgTestService.testSend();
    return AjaxResult.success();
    }
    }

    这个不需要在意细节 你用@Autowired也行 看个人习惯 用到哪个就用什么;

    4)调用结果


    九、结语

    关于 RabbitMQ 在 ruoyi-cloud 里的接入流程,就先分享到这里
    这些内容都是我自己在实际开发中一点点摸出来的,有些做法也许不是最完美,但至少都是走过、踩过、验证过的路
    如果你能从这里面找到对自己有用的东西,那我写这篇就算值了

    当然,每个人的项目场景都不一样,如果你在使用 RabbitMQ 的时候遇到新的问题、或者有更好的方案,也欢迎一起讨论
    技术不是闭门造车,大家互相补充,才能越走越稳

    那就这样吧,感谢你看到最后————完结撒花~ 下次见! *★,°*:.☆( ̄▽ ̄)/$:*.°★* 。

    赞(0)
    未经允许不得转载:171主机测评 » Spring Cloud + RabbitMQ 在 ruoyi-cloud 中的保姆级实战(含 Docker Compose 部署与工具模块新建)
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址