RabbitMQ (Java)学习笔记
RabbitMQ (Java)学习笔记
原创 已于 2025-03-14 18:35:45 修改 · 粉丝可见 · 1.2k 阅读 · 20 · 25 GEO检测 · 编辑
文章链接:https://blog.csdn.net/hacker_51/article/details/146189775
目录
[TOC]
一、概述
RabbitMQ是一种开源的消息代理软件,基于AMQP(高级消息队列协议)实现。它充当消息中间件的角色,允许应用程序通过消息队列进行异步通信。RabbitMQ的主要功能是接收、存储和转发消息,从而解耦应用程序组件,提高系统的可扩展性和可靠性。
①核心组件
生产者(Producer) :负责创建消息并将其发送到RabbitMQ服务器。生产者可以是任何应用程序,它们将消息封装成特定的格式,然后通过网络发送到消息代理。
交换器(Exchange) :交换器是消息传递的枢纽,它接收生产者发送的消息,并根据一定的规则将消息路由到一个或多个队列中。RabbitMQ提供了多种类型的交换器,如直接交换器(direct)、扇形交换器(fanout)、主题交换器(topic)等,每种交换器的路由规则不同。
直接交换器 :根据消息的路由键(routing key)直接将消息发送到与该路由键绑定的队列。
扇形交换器 :将消息广播到所有绑定的队列,不考虑路由键。
主题交换器 :根据路由键的模式匹配规则将消息发送到匹配的队列。
队列(Queue) :队列是存储消息的容器,它是一个先进先出(FIFO)的数据结构。队列中的消息会被消费者(Consumer)依次消费。队列可以持久化,即使RabbitMQ服务器重启,队列中的消息也不会丢失。
消费者(Consumer) :消费者从队列中获取消息并进行处理。消费者可以是应用程序的另一个组件,也可以是独立的进程。消费者通过订阅队列来接收消息,一旦队列中有消息到达,消费者就会按照一定的策略(如轮询、负载均衡等)进行消费。
②工作原理
消息发送过程
生产者创建消息,并指定交换器和路由键。
生产者将消息发送到RabbitMQ服务器。
交换器根据路由键将消息路由到一个或多个队列。
如果队列不存在,交换器会根据配置决定是否丢弃消息或返回错误。
消息接收过程
消费者向RabbitMQ服务器发送订阅请求,指定要消费的队列。
当队列中有消息到达时,RabbitMQ服务器将消息推送给消费者。
消费者接收到消息后进行处理,处理完成后向服务器发送确认消息(ack),告知服务器该消息已被成功消费。如果消费者在处理消息过程中失败或崩溃,RabbitMQ服务器会将消息重新放入队列,等待其他消费者消费。
③优势
高可靠性 :RabbitMQ支持消息持久化,即使服务器出现故障,消息也不会丢失。同时,它还支持镜像队列,可以在多个节点之间复制队列,提高系统的可用性。
高可用性 :RabbitMQ可以部署在多个节点上,形成集群。集群中的节点可以相互备份,当某个节点出现故障时,其他节点可以接管其工作,保证系统的正常运行。
灵活的路由策略 :通过不同类型的交换器和路由键,可以实现复杂的消息路由逻辑,满足各种业务场景的需求。
易于扩展 :RabbitMQ支持水平扩展,可以通过增加节点来提高系统的处理能力。同时,它还支持多种编程语言的客户端库,方便开发者进行集成。
④应用场景
异步任务处理 :当应用程序需要执行耗时的任务时,可以将任务封装成消息发送到RabbitMQ,然后由消费者在后台进行处理,从而提高系统的响应速度。
服务间通信 :在微服务架构中,不同的服务可以通过RabbitMQ进行通信,实现服务的解耦和异步交互。
事件驱动架构 :RabbitMQ可以作为事件总线,将事件消息传递给不同的消费者,实现事件驱动的业务逻辑。
日志收集 :将日志消息发送到RabbitMQ,然后由专门的日志处理程序进行消费和分析,实现日志的集中管理和分析。
二、入门
MQ技术选型参考图如下:

1、docker 安装MQ
1 | |
2、SpringAMQP
SpringAmqp的官方地址 : SpringAMQP官方网址

3、代码实现
pom 依赖
1 | |
配置RabbitMQ服务端信息
在每个微服务中引入MQ服务端信息,这样微服务才能连接到RabbitMQ
1 | |
发送消息
SpringAMQP提供了RabbitTemplate工具类,方便我们发送消息。样例发送消息代码如下
1 | |
接收消息
SpringAMQP提供声明式的消息监听,我们只需要通过 注解 在方法上声明要监听的队列名称,将来SpringAMQP就会把消息传递给当前方法:
1 | |
三、基础
work Queue
Workqueues,任务模型。简单来说就是 让多个消费者绑定到一个队列,共同消费队列中的消息

案例
模拟WorkQueue,实现一个队列绑定多个消费者
基本思路如下:
在RabbitMQ的控制台创建一个队列,名为work.queue
在publisher服务中定义测试方法,在1秒内产生50条消息,发送到work.queue
在consumer服务中定义两个消息监听者,都监听work.queue队列
消费者1每秒处理50条消息,消费者2每秒处理5条消息
通过Thread.sleep 方式 模拟性能不同的时候接受到的数据
1 | |
最后发现,消费者接受到的消息是默认的一人一半效果,而处理消息方面可以显示出效率比较低下

消费者消息推送限制(解决消息堆积方案之一)
默认情况下,RabbitMQ的会将消息依次 轮询投递 给绑定在队列上的每一个消费者。但这并没有考虑到消费者是否已经处理完消息,可能出现 消息堆积 。
因此我们需要修改application.yml,设置preFetch值为1,确保同一时刻最多投递给消费者1条消息:
1 | |
我们重新发送消息可以发现

这样我们可以充分消费每一条消息,能者多劳效果。
Work模型的使用
多个消费者绑定到一个队列,可以加快消息处理速度
同一条消息只会被一个消费者处理
通过设置prefetch来控制消费者预取的消息数量,处理完一条再处理下一条,实现能者多劳
Fanout交换机
真正生产环境都会经过exchange来发送消息,而不是直接发送到队列,交换机的类型有以下三种:
Fanout : 广播
Direct : 定向
Topic : 话题

Fanout 广播模式
Fanout Exchange 会将接收到的消息广播到每一个跟其绑定的queue,所以也叫广播模式
通过 exchange 广播消息给队列,多个消费者 处理绑定的一个队列中的不同订单消息,并行操作,提高效率。

案例
利用SpringAMQP演示FanoutExchange的使用
实现思路如下:
在RabbitMQ控制台中,声明队列fanout.queue1和fanout.queue2
在RabbitMO控制台中,声明交换机hmall.fanout,将两个队列与其绑定
在consumer服务中,编写两个消费者方法,分别监听fanout.queue1和fanout.queue2
在publisher中编写测试方法,向hmall.fanout发送消息.
消费者方法
1 | |
发送消息方法
1 | |
Fanout Exchange交换机的作用是什么?
接收publisher发送的消息
将消息按照规则路由到与之绑定的队列
Fanout Exchange的会将消息路由到每个绑定的队列
Direct 交换机
Direct Exchange 会将接收到的消息根据规则路由到指定的Queue,因此称为 定向 路由。
每一个Queue都与Exchange设置一个BindingKey
发布者发送消息时,指定消息的RoutingKey
Exchange将消息路由到BindingKey与消息RoutingKey一致的队列

觉得不错的话,给点打赏吧 ୧(๑•̀⌄•́๑)૭
wechat pay
ali pay