微服务03 - RabbitMQ 快速入门
本文参考 B 站黑马:https://www.bilibili.com/video/BV1mN4y1Z7t9
1、MQ 介绍
1.1、同步调用
例如:调用支付服务(用户付款成功)时,需要依次调用多个服务(订单服务(更新订单状态)、短信服务(短信通知用户)、积分服务(增加用户积分)等)。
缺点
- 性能下降。消费者需要等待所有提供者依次执行完成。
- 级联失败。如果提供者出现故障,则消费者同样出现故障。
- 耦合度高。如果新增业务需求,则需修改原有代码。
优点
- 时效性高。可以立即得到结果。
适用场景
- 对时效性要求高的场景。例如:在查询订单时,同时查询用户信息。
1.2、异步调用
例如:调用支付服务(用户付款成功)时,不用关心后续服务执行情况(订单服务(更新订单状态)、短信服务(短信通知用户)、积分服务(增加用户积分)等)。
优点
- 性能提升。
- 故障隔离。
- 耦合度低。
- 流量削峰。
缺点
- 时效性差。不能立即得到结果。
- 不能确定下游业务是否执行成功。
- 依赖 Broker(中间件,代理)的可靠性、吞吐量等。
适用场景
- 高性能、高并发、高可用,对时效性要求不高的场景。
1.3、MQ 常见技术
MQ(Message Queue,消息队列)是消息中间件,用于在分布式系统中实现异步通信。MQ 是异步调用中的 Broker(中间件,代理)。
MQ 常见技术
- RabbitMQ、RocketMQ:用于可靠性要求高的场景。例如:适合小型项目、或针对具体业务等。
- ActiveMQ:性能一般,使用较少。
- Kafka:用于吞吐量要求高、海量数据但可靠性要求低的场景。例如:针对日志、大数据等。
| RabbitMQ | ActiveMQ | RocketMQ | Kafka | |
|---|---|---|---|---|
| 公司/社区 | Rabbit | Apache | 阿里开发(现为 Apache) | Apache |
| 开发语言 | Erlang | Java | Java | Scala 与 Java |
| 协议支持(功能) | AMQP、XMPP、SMTP、STOMP | OpenWire、REST、AMQP、XMPP、STOMP | 自定义协议 | 自定义协议 |
| 可用性 | 高 | 一般 | 高 | 高 |
| 单机吞吐量(并发) | 一般 | 差 | 高 | 非常高 |
| 消息延迟 | 微秒级 | 毫秒级 | 毫秒级 | 毫秒级 |
| 消息可靠性 | 高 | 一般 | 高 | 一般 |
2、RabbitMQ 入门
2.1、RabbitMQ 安装
RabbitMQ 是基于 Erlang 语言开发的消息中间件,用于在分布式系统中实现异步通信。
官网:https://www.rabbitmq.com/。
安装(Linux 版,使用 Docker)
- 一、上传 mq.tar 到 /tmp 目录,执行以下命令。
docker load -i mq.tar # 加载镜像
docker images # 查看镜像
# 运行
docker run \
--name mq \
--hostname mq1 \
-p 15672:15672 \
-p 5672:5672 \
-d \
rabbitmq:3.8-management
# 查看
docker ps
----------------
docker run \ # 运行
--name mq \ # 名称
--hostname mq1 \ # 主机名,以后集群使用
-p 15672:15672 \ # 控制台端口,用于控制台登录
-p 5672:5672 \ # 监听端口,用于消息通信
-d \ # 后台运行
rabbitmq:3.8-management # 镜像名称
- 二、登录。访问 http://192.168.112.131:15672。账号/密码:guest/guest。
整体架构与核心概念
- Vitual Host:虚拟主机,用于数据隔离。
- Publisher:消息发布者。
- Consumer:消息消费者。
- Queue:队列,用于存储消息。
- Exchange:交换机,用于路由消息。
【RabbitMQ Server】
→→→→【Exchange→→→→Queue】→→→→Consumer
Publisher→→→→【Exchange→→→→Queue】→→→→Consumer
→→→→【Exchange→→→→Queue】→→→→Consumer
2.2、RabbitMQ 控制台操作
需求
在控制台中,创建 Queue、使用默认交换机 amq.fanout 发布消息。
步骤
- 一、在控制台的 Queue 中,添加一个 Queue。(输入 Name 然后点击 【Add queue】)
Name: hello.queue1
Name: hello.queue2
Name: hello.queue3
- 二、在控制台的 Exchange 中,绑定 Queue,发布消息。(在默认交换机 amq.fanout 中,输入 Name 然后点击【Bind】,输入 Payload 然后点击【Publish message】)
Payload: Hello,RabbitMQ!
- 三、在控制台的 Queue 中,查看消息。(点击新添加的 Queue,点击 【Get messages】)
2.3、数据隔离机制
每个 User 都有自己的 Vitual Host。
每个 Vitual Host 都有自己的 Exchange、Queue。
通常不同项目或服务,使用自己的 User 与 Vitual Host,用于数据隔离。
需求
- 一、创建一个 User。
- 二、退出当前 User,使用新 User 创建一个 Vitual Host。
- 三、测试不同 Vitual Host 之间的数据隔离现象。
示例
- 一、创建一个 User。
Username: xiaofeng
Password: xiaofeng
xiaofeng
Tag: administrator // 权限标签
- 二、退出当前 User,使用新 User 的账号密码登录,并为该 User 创建一个 Vitual Host。
Name: /xiaofeng
- 三、测试不同 Vitual Host 之间的数据隔离现象。
3、Spring AMQP 入门
3.1、Spring AMQP 介绍
AMQP(Advanced Message Queuing Protocol,高级消息队列协议)是微服务之间传递业务消息的开放标准(规范)。该协议与语言或平台无关,符合微服务对技术独立性的要求。
Spring AMQP 是基于 AMQP 定义的一套 API 规范,提供模板(RabbitTemplate)用于发布与接收消息。包括:Spring AMQP 是接口规范,Spring Rabbit 是底层实现(封装 RabbitMQ)。
官网:https://spring.io/projects/spring-amqp。
特点
- 提供 RabbitListener 消息监听器,用于异步处理消息。
- 提供 RabbitTemplate,用于发布与接收消息。
- 提供 RabbitAdmin,用于自动声明队列、交换与绑定。
3.2、环境准备
步骤
- 一、创建父项目 rabbitmq,导入 Spring AMQP 依赖。
- 二、创建子模块 consumer,创建启动类,配置 YML。
- 三、创建子模块 publisher,创建启动类,配置 YML。
示例
- 一、创建父项目 rabbitmq,导入依赖。
<properties>
<java.version>1.8</java.version>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<project.reporting.outputEncoding>UTF-8</project.reporting.outputEncoding>
<spring-boot.version>2.6.13</spring-boot.version>
</properties>
<dependencies>
<!--Spring AMQP-->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
<!--Test-->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
</dependencies>
<dependencyManagement>
<dependencies>
<!--Spring Boot 管理依赖-->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-dependencies</artifactId>
<version>${spring-boot.version}</version>
<type>pom</type>
<scope>import</scope>
</dependency>
</dependencies>
</dependencyManagement>
- 二、创建子模块 consumer,创建启动类,配置 YML。
@SpringBootApplication
public class ConsumerApplication {
public static void main(String[] args) {
SpringApplication.run(ConsumerApplication.class, args);
}
}
spring:
rabbitmq: # RabbitMQ 配置
host: 192.168.112.131 # 主机
port: 5672 # 监听端口
virtual-host: /xiaofeng # 虚拟主机
username: xiaofeng # 用户名
password: xiaofeng # 密码
- 三、创建子模块 publisher,创建启动类,配置 YML。
@SpringBootApplication
public class PublisherApplication {
public static void main(String[] args) {
SpringApplication.run(PublisherApplication.class, args);
}
}
spring:
rabbitmq: # RabbitMQ 配置
host: 192.168.112.131 # 主机
port: 5672 # 监听端口
virtual-host: /xiaofeng # 虚拟主机
username: xiaofeng # 用户名
password: xiaofeng # 密码
3.3、消息队列模式
官网教程:https://www.rabbitmq.com/tutorials。
| 模式 | 交换机 | 说明 |
|---|---|---|
| Hello World 简单模式 | 无 | 一个发布者,一个消费者。 |
| Work Queues 工作模式 | 无 | 一个发布者,多个消费者互相竞争,可以提高消息处理速度。 |
| Publish/Subscribe 发布/订阅模式 | Fanout | 一个发布者,多个消费者收到相同消息。 |
| Routing 路由模式 | Direct | 一个发布者,多个消费者选择性的接收消息。 |
| Topics 主题模式 | Topic | 一个发布者,多个消费根据主题选择性的接收消息。 |
| RPC 请求/回复模式 | ||
| Publisher Confirms 发布者确认模式 |
3.3、简单模式
需求
- 一、在控制台中,创建队列 simple.queue。
- 二、在 consumer 中,使用 @RabbitListener 监听 simple.queue,然后启动 Spring Boot。
- 三、在 publisher 中,使用 RabbitTemplate 向 simple.queue 发布消息。
示例
- 一、在控制台中,创建队列 simple.queue。
- 二、在 consumer 中,使用 @RabbitListener 监听 simple.queue,然后启动 Spring Boot。
@Component
public class SpringRabbitListener {
// Simple
@RabbitListener(queues = "simple.queue")
public void listenSimpleQueue(String msg) {
System.out.println("【收到消息】:" + msg);
}
}
- 三、在 publisher 中,使用 RabbitTemplate 向 simple.queue 发布消息,然后运行测试方法。
@SpringBootTest
public class PublisherTest {
@Autowired
private RabbitTemplate rabbitTemplate;
@Test
void send2SimpleQueue() {
String queueName = "simple.queue";// 队列名称
String msg = "hello,world";// 消息
rabbitTemplate.convertAndSend(queueName, msg);// 发布消息
}
}
3.4、工作模式
需求
- 一、在控制台中,创建队列 work.queue。
- 二、在 consumer 中,定义两个方法,同时监听 work.queue。consumer1 每秒处理 50 条消息,consumer2 每秒处理 5 条消息。
- 三、在 publisher 中,1 秒发布 50 条消息到 work.queue。
示例
- 一、在控制台中,创建队列 work.queue。
- 二、在 consumer 中,定义两个方法,同时监听 work.queue。consumer1 每秒处理 50 条消息,consumer2 每秒处理 5 条消息。
// Work
@RabbitListener(queues = "work.queue")
public void listenWorkQueue1(String msg) throws InterruptedException {
Thread.sleep(20);
System.out.println("【Consumer1 收到消息】:" + msg);
}
@RabbitListener(queues = "work.queue")
public void listenWorkQueue2(String msg) throws InterruptedException {
Thread.sleep(200);
System.out.println("【Consumer2 收到消息】========:" + msg);
}
- 三、在 publisher 中,1 秒发布 50 条消息到 work.queue。
@Test
void send2WorkQueue() throws InterruptedException {
String queueName = "work.queue";
for (int i = 1; i <= 50; i++) {
String message = "消息" + i;
rabbitTemplate.convertAndSend("work.queue", message);
Thread.sleep(20);
}
}
3.5、限制消息预取
在默认情况下,RabbitMQ 会将消息轮询投递给与队列绑定的每个消费者,并未考虑每个消费者的处理消息能力。例如:Consumer1、Consumer2 都是处理 25 条消息。
限制消息预取,设置 prefetch 为 1,可以确保每次仅投递给每个消费者 1 条消息,处理完成后再从队列中获取。例如:Consumer1、Consumer2 根据处理消息能力,能者多劳。
示例
在 consumer 中,配置 prefetch 的值为 1。
spring:
rabbitmq: # RabbitMQ 配置
host: 192.168.112.131 # 主机
port: 5672 # 监听端口
virtual-host: /xiaofeng # 虚拟主机
username: xiaofeng # 用户名
password: xiaofeng # 密码
listener:
simple:
prefetch: 1 # 限制消息预取,默认没有限制
3.6、发布/订阅模式(Fanout 交换机 )
Fanout ExChange 会将接收到的消息发布到所有与其绑定的队列,也称为广播模式。
需求
- 一、在控制台中,创建队列 fanout.queue1、fanout.queue2,创建交换机 xiaofeng.fanout 并绑定这个两个队列。
- 二、在 consumer 中,创建两个方法,分别监听 fanout.queue1、fanout.queue2。
- 三、在 publisher 中,创建测试方法,向 xiaofeng.fanout 发布消息。
示例
-
一、在控制台中,创建队列 fanout.queue1、fanout.queue2,创建交换机 xiaofeng.fanout 并绑定这个两个队列。
-
二、在 consumer 中,创建两个方法,分别监听 fanout.queue1、fanout.queue2。
// Fanout
@RabbitListener(queues = "fanout.queue1")
public void listenFanout1(String msg) {
System.out.println("【Consumer1 收到消息】:" + msg);
}
@RabbitListener(queues = "fanout.queue2")
public void listenFanout2(String msg) {
System.out.println("【Consumer2 收到消息】:" + msg);
}
- 三、在 publisher 中,创建测试方法,向 xiaofeng.fanout 发布消息。
@Test
void sendFanout() {
String exChangeName = "xiaofeng.fanout";
String msg = "重要通知,明天所有人加班";
rabbitTemplate.convertAndSend(exChangeName, null, msg);
}
3.7、路由模式(Direct 交换机)
Direct Exchange 会将接收到的消息根据规则发布到指定的队列,也称为定向模式。
需求
- 一、在控制台中,创建队列 direct.queue1、direct.queue2,创建交换机 xiaofeng.direct 并绑定这个两个队列。
- 二、在 consumer 中,创建两个方法,分别监听 direct.queue1、direct.queue2。
- 三、在 publisher 中,创建测试方法,根据不同的 Routing Key 向 xiaofeng.direct 发布消息。
示例
- 一、在控制台中,创建队列 direct.queue1、direct.queue2,创建交换机 xiaofeng.direct 并绑定这个两个队列。
direct.queue1 blue
direct.queue1 red
direct.queue2 red
direct.queue2 yellow
- 二、在 consumer 中,创建两个方法,分别监听 direct.queue1、direct.queue2。
// Direct
@RabbitListener(queues = "direct.queue1")
public void listenDirect1(String msg) {
System.out.println("【Consumer1 收到消息】:" + msg);
}
@RabbitListener(queues = "direct.queue2")
public void listenDirect2(String msg) {
System.out.println("【Consumer2 收到消息】:" + msg);
}
- 三、在 publisher 中,创建测试方法,根据不同的 Routing Key 向 xiaofeng.direct 发布消息。
@Test
void sendDirect() {
String exChangeName = "xiaofeng.direct";
String msg = "西安暴雨预警,注意防范";
// rabbitTemplate.convertAndSend(exChangeName, "red", msg);
// rabbitTemplate.convertAndSend(exChangeName, "blue", msg);
rabbitTemplate.convertAndSend(exChangeName, "yellow", msg);
}
3.8、主题模式(Topic 交换机)
Topic Exchange 与 Direct Exchange 类似,区别是 Routing Key 可以使用通配符。
- “#”:表示零或多个单词。例如:Routing Key 为 china.# 时, 可以匹配到 china.weather、china.news,不能匹配到 japan.weather、japan.news。
- “*”:表示一个单词。
需求
- 一、在控制台中,创建队列 topic.queue1、topic.queue2,创建交换机 xiaofeng.topic 并绑定这个两个队列。
- 二、在 consumer 中,创建两个方法,分别监听 topic.queue1、topic.queue2。
- 三、在 publisher 中,创建测试方法,根据不同的 Routing Key 向 xiaofeng.topic 发布消息。
示例
- 一、在控制台中,创建队列 topic.queue1、topic.queue2,创建交换机 xiaofeng.topic 并绑定这个两个队列。
topic.queue1 china.#
topic.queue2 #.news
- 二、在 consumer 中,创建两个方法,分别监听 topic.queue1、topic.queue2。
// Topic
@RabbitListener(queues = "topic.queue1")
public void listenTopic1(String msg) {
System.out.println("【Consumer1 收到消息】:" + msg);
}
@RabbitListener(queues = "topic.queue2")
public void listenTopic2(String msg) {
System.out.println("【Consumer2 收到消息】:" + msg);
}
- 三、在 publisher 中,创建测试方法,根据不同的 Routing Key 向 xiaofeng.topic 发布消息。
@Test
void sendTopic() {
String exChangeName = "xiaofeng.topic";
String msg = "今天天气晴,平均温度 28°C";
// rabbitTemplate.convertAndSend(exChangeName, "china.weather", msg);
// rabbitTemplate.convertAndSend(exChangeName, "china.news", msg);
rabbitTemplate.convertAndSend(exChangeName, "japan.news", msg);
}
4、创建队列与交换机
4.1、使用配置类(不太建议)
Spring AMQP 提供工厂类 ,用于创建队列与交换机、及其绑定关系。
- Queue:用于创建队列,可以使用工厂类 QueueBuilder。
- Exchange:用于创建交换机,可以使用工厂类 ExchangeBuilder。
- Binding:用于绑定队列与交换机,可以使用工厂类 BindingBuilder。
需求
使用配置类,实现 Fan 交换机中的示例。
示例
- 一、在控制台中,删除相关队列、交换机,及其绑定关系。
- 二、在 Consumer 中,使用配置类。
@Configuration
public class FanoutConfig {
// 队列
@Bean
public Queue fanoutQueue1() {
return QueueBuilder.durable("fanout.queue1").build();
// return new Queue("fanout.queue1");
}
@Bean
public Queue fanoutQueue2() {
return QueueBuilder.durable("fanout.queue2").build();
// return new Queue("fanout.queue2");
}
// 交换机
@Bean
public FanoutExchange xiaofengFanout() {
return ExchangeBuilder.fanoutExchange("xiaofeng.fanout").build();
// return new FanoutExchange("xiaofeng.fanout");
}
// 绑定关系
// @Bean
// public Binding binding1() {
// return BindingBuilder.bind(fanoutQueue1()).to(xiaofengFanout());
// }
// @Bean
// public Binding binding2() {
// return BindingBuilder.bind(fanoutQueue2()).to(xiaofengFanout());
// }
@Bean
public Binding binding1(Queue fanoutQueue1, FanoutExchange xiaofengFanout) {
return BindingBuilder.bind(fanoutQueue1).to(xiaofengFanout);
}
@Bean
public Binding binding2(Queue fanoutQueue2, FanoutExchange xiaofengFanout) {
return BindingBuilder.bind(fanoutQueue2).to(xiaofengFanout);
}
}
4.2、使用注解(建议)
需求
使用注解,实现 Topic 交换机中的示例。
示例
- 一、在控制台中,删除相关队列、交换机,及其绑定关系。
- 二、在 Consumer 中,使用注解。
// Topic
// @RabbitListener(queues = "topic.queue1")
// public void listenTopic1(String msg) {
// System.out.println("【Consumer1 收到消息】:" + msg);
// }
// @RabbitListener(queues = "topic.queue2")
// public void listenTopic2(String msg) {
// System.out.println("【Consumer2 收到消息】:" + msg);
// }
// Topic
@RabbitListener(bindings = @QueueBinding(
value = @Queue(name = "topic.queue1"),// 队列
exchange = @Exchange(name = "xiaofeng.topic", type = ExchangeTypes.TOPIC),// 交换机
key = {"china.#"}// Routing key
))
public void listenTopic1(String msg) {
System.out.println("【Consumer1 收到消息】:" + msg);
}
@RabbitListener(bindings = @QueueBinding(
value = @Queue(name = "topic.queue2"),// 队列
exchange = @Exchange(name = "xiaofeng.topic", type = ExchangeTypes.TOPIC),// 交换机
key = {"#.news"}// Routing key
))
public void listenTopic2(String msg) {
System.out.println("【Consumer2 收到消息】:" + msg);
}
5、消息转换器
5.1、发布 Object 类型的消息
需求
使用 RabbitTemplate 的 convertAndSend() 方法发布 Map 类型的消息,在控制台中查看消息。
示例
- 一、在 consumer 中,创建队列 object.queue。
@Component
public class QueueConfig {
@Bean
public Queue objectQueue() {
return QueueBuilder.durable("object.queue").build();
}
}
- 二、在 publisher 中,发布 Map 类型的消息。
@Test
void send2SimpleQueue() {
String queueName = "object.queue";
Map<String, String> msg = new HashMap<>();
msg.put("name", "张三");
msg.put("age", "22");
rabbitTemplate.convertAndSend(queueName, msg);
}
- 三、在控制台中查看消息。
content_type: application/x-java-serialized-object
Payload: rO0ABXNyABFqYXZhLnV0aWwuSGFzaE1hcAUH2sHDFmDRAwACRgAKbG9hZEZhY3RvckkACXRocmVzaG9sZHhwP0AAAAAAAAx3CAAAABAAAAACdAAEbmFtZXQA
BuW8oOS4iXQAA2FnZXQAAjIyeA==
这里使用的是 JDK 序列化。
Spring 使用 org.springframework.amqp.support.converter.MessageConverter 处理对象,默认实现是 SimpleMessageConverter,最终使用 ObjectOutputStream 完成序列化。
5.2、消息转换器
示例
- 一、在控制台中,删除队列 object.queue。
- 二、在父项目中,导入 JSON 依赖。(jackson-dataformat-xml 包含 jackson-databind,两者均可)
<!--JSON-->
<dependency>
<groupId>com.fasterxml.jackson.core</groupId>
<artifactId>jackson-databind</artifactId>
</dependency>
<!-- <dependency>-->
<!-- <groupId>com.fasterxml.jackson.dataformat</groupId>-->
<!-- <artifactId>jackson-dataformat-xml</artifactId>-->
<!-- </dependency>-->
- 三、在 consumer 与 publisher 的启动类中,注入 Jackson2JsonMessageConverter。(或使用配置类)
@Bean
public MessageConverter messageConverter() {
return new Jackson2JsonMessageConverter();
}
- 四、在控制台中查看消息。
content_type: application/json
Payload: {"name":"张三","age":"22"}
更多推荐


所有评论(0)