智能BI项目-引入RabbitMQ(七)
往期开发(已经完成)
- 智能 BI 项目-介绍(一)
- 智能 BI 项目-初始化(二)
- 智能 BI 项目-初学 AI 分析(三)
- 智能 BI 项目-AI 接口调用(四)
- 智能 BI 项目-接口优化(五)
- 智能 BI 项目-接口的异步化(六)
- 智能 BI 项目-引入 RabbitMQ(七)
今日后端开发
分析目前系统现状
现状: 目前实现异步的方式是通过本地的从线程池实现的
无法集中的去限制,只能够单机限制
假设 AI 服务限制只能有 2-3 个用户同时去调用,目前可以通过设置线程池最大核心线程数为 2-3 来实现。
假设系统用量极具增大,要求将系统进行分布式改造,多台服务器,每个服务器都要有两个最大核心线程数 2-3
就会有 2N-3N 个线程数这样总数加起来就超过 AI 服务最大的限制
解决方案:将所有的服务统一的发送到一个消息管理地方,然后统一的把任务分发给各个分布式系统(集中地去存储并分发要执行的项目)通过本地线城池实现的异步化,由于任务是放在内存中执行的,丢失的可能性非常大,假设系统宕机就会导致系统丢失本次任务
虽然任务会丢失,但是可以通过人工手段,从数据库中取出,通过某些手段,重新执行此任务,但是开发维护成本非常高
以上就是很经典的重试场景,可以通过定时任务去解决,是不需要我们开发者过于关心任务丢失重试的场景
解决方案:把任务存储到一个持久化的存储硬盘来进行存储。优化
如果系统功能的越来越多,长耗时任务越来越长,系统越来越复杂,例如要开多个线程池,资源可能会出现资源竞争的情况
此时我们可以把系统进行拆分,也就是服务拆分,把长耗时任务,消耗很多资源的任务,单独的抽成一个程序,不影响主业务的执行
解决方案:可以使用一个中间件(中间人),让中间件去帮我们去连接俩个系统(程序),比如核心系统和智能生成业务
引入中间件(RabbitMQ)
连接多个系统的桥梁,或者帮助多个系统进行通信、协作
例如 Redis、消息队列、分布式存储 Etcd
- 概念:存储消息的队列
- 关键字:存储、消息、队列
- 存储:存储数据
- 消息:某种数据结构,比如字符串、对象、JSON 字符串、二进制数据等等
- 队列:先进先出的数据结构
- 见解:消息队列是一种特殊的数据库吗?应该是可以这样理解吧
- 应用场景
在多个不同的系统,或者是应用之间实现消息的传输,也可以进行存储
不需要考虑传输应用的编程语言、系统、框架等。可以对应用进行解耦
- 生产者:Producer,可以比作是快递员,发送消息的人(客户端)
- 消费者:Consumer,可以比作是取快递的人,接受读取消息的人(客户端)
- 消息:Message,可以比作是快递,就是快递员派送的快递也是给消费者传递的消息
- 消息队列:Queue,可以比作是快递用快递车来进行派送,存储到车里面进行派送,消息派送的队列
为什么不直接进行传输,而要用消息队列?
这是因为生产者不需要消费者什么时候去消息,要不要消费,生产者只需要把东西给消息队列,业务就完成了
这就将生产者和消费者进行了解耦,俩者互相不影响,自己完成任务即可
为什么使用消息队列
异步处理
- 生产者发送完消息之后,就可以继续去接受其他任务,不需要等待,
- 消费者想什么时候去取出这个消息去消费都可以,不会产生堵塞
削峰填谷
- 先把用户请求,让生产者发送到消息队列中,消费者可以按照自己的需求,慢慢的去取消息
- 例如:某一时刻来了十万个请求,原本情况下十万个请求要直接到内部进行处理,很快系统就扛不住压力就会宕机了
- 解决:请求进来时候,先把消息放到消息队列中,系统按照自己最大的处理能力,去定量的取出消息恒定速率的去处理消息,从而降低了系统的压力,稳定的去处理
消息队列的优势
- 数据持久化:可以把消息几种的存储到硬盘里,服务器重启也不会丢失
- 可扩展性:可以根据需求,随时的增加或者减少节点,继续保持稳定的服务
- 应用解耦:可以连接各个不同语言、框架开发系统,让这些系统能够灵活的传输读取数据等等。
应用解耦的优点
- 最早的开发是把所有的功能都放到同一个项目中,调用多个子系统是,一个环节出错或者导致系统挂掉,系统整体就要出错
- 使用消息队列进行解耦。一个系统挂了就不会影响另一个系统的运行
- 系统挂了在后续恢复之后,仍然可以继续取出消息,重试之前的业务逻辑
- 生产者把消息发送到消息队列,就可以立即返回,不用同步调用所有系统,性能也会提高
订阅模式
- 如果一个非常大的系统要给其他子系统发送通知,最简单的方式就是一个个的去通知子系统,调用相关的系统去通知。
- 这样会产生很大的问题:每次更新发通知就要一个个调用,效率很低,速度还慢,还一直在占用资源。还有可能通知过程中有的调用失败没有通知到
- 解决方案:大的核心系统每次通知把消息只发送到一个地方,其他的系统都去订阅这个地方也就是消息队列,去读取这个消息队列的通知信息
消息队列的缺点
俗话说没有绝对完美的设计,总归这么多好处的消息队列,也有它的缺点
如果要给系统引入额外的消息中间件,系统就会变得更加复杂,庞大,而且还要花费成本去维护这个中间件,额外的费用去部署
还要面临,消息丢失,处理的顺序,重复消费,数据的一致性,也就是分布式系统要考虑的问题
主流的消息队列选型
主流技术
- activemq
- rabbitmq
- kafka
- rocketmq
- zeromq
- pulsar
- Apache Inlong (Tube)
技术的选型指标
- 吞吐量:IO、并发
- 时效性:类似延迟、消息的发送、到达时间
- 可用性:系统可用的比率 宕机的可能性
- 可靠性:消息不丢失,功能正常完成

RabbitMQ 入门实战
特点
生态好、好学习、易于理解、时效性强、支持很多不同语言的客户端、可扩展性、可用性都不错。
学习性价比高的消息队列,适用于绝大多数中小规模的分布式系统
官网:https://www.rabbitmq.com基本概念
AMQP 协议: https://www.rabbitmq.com/tutorials/amqp-concepts.html
高级消息队列协议 (Advanced Message Queue Protocol)- 生产者:发消息到某个交换机
- 消费者: 从某个队列中取消息
- 交换机 (Exchange): 负责把消息 转发到对应的队列
- 队列 (Queue) : 存储消息的
- 路由(Routes): 转发,就是怎么把消息从一个地方转到另一个地方 (比如从生产者转发到某个队列)

引入依赖
1
2
3
4
5
6<!-- https://mvnrepository.com/artifact/com.rabbitmq/amqp-client -->
<dependency>
<groupId>com.rabbitmq</groupId>
<artifactId>amqp-client</artifactId>
<version>5.17.0</version>
</dependency>单向发送
一个生产者给一个队列发送消息,一个消费者从这个队列中取消息 1<->1

- 生产者代码
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
package com.bi.spring.demomq;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import lombok.extern.slf4j.Slf4j;
import java.nio.charset.StandardCharsets;
public class Producer {
private static final String QUEUE_NAME = "hello_mq";
public static void main(String[] args) {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
try {
// 从工厂中创建连接
Connection connection = factory.newConnection();
// 获取通信通道
Channel channel = connection.createChannel();
// 声明一个队列
channel.queueDeclare(QUEUE_NAME, false, false, false, null);
String message = "Hello,World";
// 发送消息
channel.basicPublish("", QUEUE_NAME, null, message.getBytes(StandardCharsets.UTF_8));
// 打印信息
System.out.println("Producer send message: " + message);
} catch (Exception e) {
log.info("创建链接失败!E:" + e);
}
}
}Channel 频道:可以理解为消息队列的 Client(比如 jdbcClient,redisClient),提供了和消息队列简历 server 建立的通信的传输方法为了复用链接,提高传输效率。程序通过操作 Channel 来操作 rabbitmq(收发消息操作)
创建消息队列的参数
- queueName:消息队列名称 同名称和参数的消息队列只能创建一次
- durabale:消息队列重启后,消息是否丢失
- exclusive:是否只允许当前创建的消息队列的连接操作消息队列
- autoDelete:没有任何东西使用消息队列后,是否要删除队列
- 消费者代码
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36package com.bi.spring.demomq;
import com.rabbitmq.client.*;
import lombok.extern.slf4j.Slf4j;
import java.nio.charset.StandardCharsets;
public class Consumer {
private static final String QUEUE_NAME = "hello_mq";
public static void main(String[] args) {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
try {
// 从工厂中创建连接
Connection connection = factory.newConnection();
// 获取通信通道
Channel channel = connection.createChannel();
// 声明一个队列 没有就会创建 有的话就不再创建
channel.queueDeclare(QUEUE_NAME, false, false, false, null);
System.out.println("Consumer waiting for producer to send message");
// 定义如何处理消息 这个是函数式接口
DeliverCallback deliverCallback = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
System.out.println("Consumer received message from producer,message:" + message);
};
// 消费消息
channel.basicConsume(QUEUE_NAME, true, deliverCallback, consumer -> {
});
} catch (Exception e) {
log.info("创建链接失败!E:" + e);
}
}
}多消费者
场景:多个系统同时去接受并处理任务,尤其是每个系统的处理能力有限
一个生产者给一个队列发消息,多个消费者从这个队列取消息。1<->n
- 队列持久化
1
channel.queueDeclare(TASK_QUEUE_NAME, true, false, false, null);
- 消息持久化
指定 MessageProperties.PERSISTENT_TEXT_PLAIN 参数
1
2channel.basicPublish("", QUEUE_NAME, MessageProperties.PERSISTENT_TEXT_PLAIN,
message.getBytes(StandardCharsets.UTF_8));- 生产者代码
使用 Scanner 来模拟多次输入
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41package com.bi.spring.demomq;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.MessageProperties;
import lombok.extern.slf4j.Slf4j;
import java.nio.charset.StandardCharsets;
import java.util.Scanner;
public class MultiProducer {
private static final String QUEUE_NAME = "multi_queue";
public static void main(String[] args) {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
try {
// 从工厂中创建连接
Connection connection = factory.newConnection();
// 获取通信通道
Channel channel = connection.createChannel();
// 声明一个队列
channel.queueDeclare(QUEUE_NAME, true, false, false, null);
// 声明一个输入对象
Scanner scanner = new Scanner(System.in);
while (true) {
// 输入信息
String message = scanner.nextLine();
// 发送消息
channel.basicPublish("", QUEUE_NAME, MessageProperties.PERSISTENT_TEXT_PLAIN, message.getBytes(StandardCharsets.UTF_8));
// 打印信息
System.out.println("Producer send message: " + message);
}
} catch (Exception e) {
log.info("创建链接失败!E:" + e);
}
}
}- 消费者代码
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57package com.bi.spring.demomq;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.DeliverCallback;
import lombok.extern.slf4j.Slf4j;
import java.nio.charset.StandardCharsets;
public class MultiConsumer {
private static final String QUEUE_NAME = "multi_queue";
public static void main(String[] args) {
try {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
// 从工厂中创建连接
Connection connection = factory.newConnection();
// 循环创建
for (int i = 0; i < 2; i++) {
// 获取通信通道
Channel channel = connection.createChannel();
// 声明一个队列 没有就会创建 有的话就不再创建
channel.queueDeclare(QUEUE_NAME, true, false, false, null);
System.out.println("Consumer waiting for producer to send message");
// 定义如何处理消息 这个是函数式接口
int finalI1 = i;
DeliverCallback deliverCallback = (consumerTag, delivery) -> {
try {
String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
System.out.println("Consumer received message from producer,message:" + message + "-编号" + finalI1);
// 确认消息
channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
// 模拟处理能力
Thread.sleep(10000);
} catch (Exception e) {
e.printStackTrace();
// 消息失败策略
channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, false);
} finally {
System.out.println("Consumer is Done!");
// 如果没有失败最终都要确认消息
channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
}
};
// 消费消息
channel.basicConsume(QUEUE_NAME, true, deliverCallback, consumer -> {
});
}
} catch (Exception e) {
log.info("创建链接失败!E:" + e);
}
}
}控制单个消费的处理任务的积压数,每个消费者最多同时处理的任务数
1
channel.basicQos(1);
交换机
一个生产者给多个队列发送消息,一个生产者 -> 多个队列
交换机的作用:提供消息转发功能,类似与网络路由器
问题:怎么把消息转发到不同的队列上,好让消费者从不同的队列消费绑定:交换机可以与队列进行关联,通过路由 key 进行绑定
1
channel.queueBind(QUEUE_NAME,EXCHANGE_NAME,"key");

fanout 扇出、广播
特点:消息会被转发到所有绑定到该交换机的队列
场景:很适用于发布订阅的场景。比如写日志,可以多个系统间共享
注意点:生产者和消费者需要绑定同一个交换机,并且要先创建队列,才能绑定
- 生产者代码
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42package com.bi.spring.demomq;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.MessageProperties;
import lombok.extern.slf4j.Slf4j;
import java.nio.charset.StandardCharsets;
import java.util.Scanner;
public class FanoutProducer {
private static final String EXCHANGE_NAME = "fanout-exchange";
public static void main(String[] args) {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
try {
// 从工厂中创建连接
Connection connection = factory.newConnection();
// 获取通信通道
Channel channel = connection.createChannel();
// 声明一个交换机
channel.exchangeDeclare(EXCHANGE_NAME, "fanout");
// 声明一个输入对象
Scanner scanner = new Scanner(System.in);
while (scanner.hasNext()) {
// 输入信息
String message = scanner.nextLine();
// 发送消息
channel.basicPublish(EXCHANGE_NAME, "", MessageProperties.PERSISTENT_TEXT_PLAIN, message.getBytes(StandardCharsets.UTF_8));
// 打印信息
System.out.println("Producer send message: " + message);
}
} catch (Exception e) {
log.info("创建链接失败!E:" + e);
}
}
}- 消费者代码
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60package com.bi.spring.demomq;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.DeliverCallback;
import lombok.extern.slf4j.Slf4j;
import java.nio.charset.StandardCharsets;
public class FanoutConsumer {
private static final String EXCHANGE_NAME = "fanout-exchange";
private static final String QUEUE_NAME1 = "fanout-queue-one";
private static final String QUEUE_NAME2 = "fanout-queue-two";
public static void main(String[] args) {
try {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
// 从工厂中创建连接
Connection connection = factory.newConnection();
// 创建频道
Channel channel1 = connection.createChannel();
Channel channel2 = connection.createChannel();
// 声明交换机
channel1.exchangeDeclare(EXCHANGE_NAME, "fanout");
// 创建队列 1
channel1.queueDeclare(QUEUE_NAME1, true, false, false, null);
channel1.queueBind(QUEUE_NAME1, EXCHANGE_NAME, "");
// 创建队列2
channel2.queueDeclare(QUEUE_NAME2, true, false, false, null);
channel2.queueBind(QUEUE_NAME2, EXCHANGE_NAME, "");
System.out.println("Consumer waiting for producer to send message");
// 声明消息处理回调函数1
DeliverCallback deliverCallback1 = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
System.out.println("[1] Received " + message);
};
// 声明消息处理回调函数2
DeliverCallback deliverCallback2 = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
System.out.println("[2] Received " + message);
};
// 开启消费监听1
channel1.basicConsume(QUEUE_NAME1, true, deliverCallback1, consumerTag -> {
});
// 开启消费监听2
channel2.basicConsume(QUEUE_NAME2, true, deliverCallback2, consumerTag -> {
});
} catch (Exception e) {
log.info("创建链接失败!E:" + e);
}
}
}效果:所有的消费者都能够接收到消息
direct 交给特定的队列
routingkey:可以让交换机和队列进行关联,可以指定让交换机给那个队列发送消息,路由键,控制消息要给那个队列
特点:消息会根据路由键转发到指定的队列
场景:特定的消息只交给特定的系统来处理
绑定关系:完全匹配字符串
多个队列可以绑定同一个路由键

- 生产者代码
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46package com.bi.spring.demomq;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.MessageProperties;
import lombok.extern.slf4j.Slf4j;
import java.nio.charset.StandardCharsets;
import java.util.Scanner;
public class DirectProducer {
private static final String EXCHANGE_NAME = "direct-exchange";
public static void main(String[] args) {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
try {
// 从工厂中创建连接
Connection connection = factory.newConnection();
// 获取通信通道
Channel channel = connection.createChannel();
// 声明一个交换机
channel.exchangeDeclare(EXCHANGE_NAME, "direct");
// 声明一个输入对象
Scanner scanner = new Scanner(System.in);
while (scanner.hasNext()) {
// 输入信息
String userInput = scanner.nextLine();
String[] s = userInput.split(" ");
// 消息
String message = s[0];
// 路由键
String routingKey = s[1];
// 发送消息
channel.basicPublish(EXCHANGE_NAME, routingKey, MessageProperties.PERSISTENT_TEXT_PLAIN, message.getBytes(StandardCharsets.UTF_8));
// 打印信息
System.out.println("Producer send message: " + message + "- routingKey" + routingKey);
}
} catch (Exception e) {
log.info("创建链接失败!E:" + e);
}
}
}- 消费者代码
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62package com.bi.spring.demomq;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.DeliverCallback;
import lombok.extern.slf4j.Slf4j;
import java.nio.charset.StandardCharsets;
public class DirectConsumer {
private static final String EXCHANGE_NAME = "direct-exchange";
private static final String QUEUE_NAME1 = "direct-queue-one";
private static final String QUEUE_NAME2 = "direct-queue-two";
public static void main(String[] args) {
try {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
// 从工厂中创建连接
Connection connection = factory.newConnection();
// 创建频道
Channel channel1 = connection.createChannel();
Channel channel2 = connection.createChannel();
// 声明交换机
channel1.exchangeDeclare(EXCHANGE_NAME, "fanout");
// 创建队列 1
channel1.queueDeclare(QUEUE_NAME1, true, false, false, null);
// 绑定路由键1
channel1.queueBind(QUEUE_NAME1, EXCHANGE_NAME, "routing_key_one");
// 创建队列2
channel2.queueDeclare(QUEUE_NAME2, true, false, false, null);
// 绑定路由键2
channel2.queueBind(QUEUE_NAME2, EXCHANGE_NAME, "routing_key_two");
System.out.println("Consumer waiting for producer to send message");
// 声明消息处理回调函数1
DeliverCallback deliverCallback1 = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
System.out.println("[1] Received " + message);
};
// 声明消息处理回调函数2
DeliverCallback deliverCallback2 = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
System.out.println("[2] Received " + message);
};
// 开启消费监听1
channel1.basicConsume(QUEUE_NAME1, true, deliverCallback1, consumerTag -> {
});
// 开启消费监听2
channel2.basicConsume(QUEUE_NAME2, true, deliverCallback2, consumerTag -> {
});
} catch (Exception e) {
log.info("创建链接失败!E:" + e);
}
}
}topic 会模糊发送特定的队列
特点:消息会根据一个模糊的路由键转发到指定的队列
场景:特定的一类的消息可以交给特定的一类系统来处理
绑定关系:可以模糊匹配多个绑定- *:匹配一个单词,比如 *.orange,可匹配项有 a.orange、b.orange 缺点就是必须要有匹配项不然匹配不到
- #: 匹配零个或者多个单词,比如 a.#,可匹配项有 a.a、a.b、a.a.a 也可以不匹配,缺点就是太过灵活

示例图

- 生产者代码
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50package com.bi.spring.demomq;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.MessageProperties;
import lombok.extern.slf4j.Slf4j;
import java.nio.charset.StandardCharsets;
import java.util.Scanner;
public class TopicProducer {
private static final String EXCHANGE_NAME = "topic-exchange";
public static void main(String[] args) {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
try {
// 从工厂中创建连接
Connection connection = factory.newConnection();
// 获取通信通道
Channel channel = connection.createChannel();
// 声明一个交换机
channel.exchangeDeclare(EXCHANGE_NAME, "topic");
// 声明一个输入对象
Scanner scanner = new Scanner(System.in);
while (scanner.hasNext()) {
// 输入信息
String userInput = scanner.nextLine();
String[] s = userInput.split(" ");
if (s.length < 1) {
continue;
}
// 消息
String message = s[0];
// 路由键
String routingKey = s[1];
// 发送消息
channel.basicPublish(EXCHANGE_NAME, routingKey, MessageProperties.PERSISTENT_TEXT_PLAIN, message.getBytes(StandardCharsets.UTF_8));
// 打印信息
System.out.println("Producer send message: " + message + "- routingKey" + routingKey);
}
} catch (Exception e) {
log.info("创建链接失败!E:" + e);
}
}
}- 消费者代码
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
package com.bi.spring.demomq;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.DeliverCallback;
import lombok.extern.slf4j.Slf4j;
import java.nio.charset.StandardCharsets;
public class TopicConsumer {
private static final String EXCHANGE_NAME = "topic-exchange";
private static final String QUEUE_NAME1 = "topic-queue-one";
private static final String QUEUE_NAME2 = "topic-queue-two";
private static final String QUEUE_NAME3 = "topic-queue-three";
public static void main(String[] args) {
try {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
// 从工厂中创建连接
Connection connection = factory.newConnection();
// 创建频道
Channel channel = connection.createChannel();
// 声明交换机
channel.exchangeDeclare(EXCHANGE_NAME, "topic");
// 创建队列 1
channel.queueDeclare(QUEUE_NAME1, true, false, false, null);
// 绑定路由键1
channel.queueBind(QUEUE_NAME1, EXCHANGE_NAME, "#.队列1.#");
// 创建队列2
channel.queueDeclare(QUEUE_NAME2, true, false, false, null);
// 绑定路由键2
channel.queueBind(QUEUE_NAME2, EXCHANGE_NAME, "#.队列2.#");
// 创建队列3
channel.queueDeclare(QUEUE_NAME3, true, false, false, null);
// 绑定路由键3
channel.queueBind(QUEUE_NAME3, EXCHANGE_NAME, "#.队列3.#");
System.out.println("Consumer waiting for producer to send message");
// 声明消息处理回调函数1
DeliverCallback deliverCallback1 = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
System.out.println("[1] Received " + message);
};
// 声明消息处理回调函数2
DeliverCallback deliverCallback2 = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
System.out.println("[2] Received " + message);
};
// 声明消息处理回调函数2
DeliverCallback deliverCallback3 = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
System.out.println("[3] Received " + message);
};
// 开启消费监听1
channel.basicConsume(QUEUE_NAME1, true, deliverCallback1, consumerTag -> {
});
// 开启消费监听2
channel.basicConsume(QUEUE_NAME2, true, deliverCallback2, consumerTag -> {
});
// 开启消费监听3
channel.basicConsume(QUEUE_NAME3, true, deliverCallback3, consumerTag -> {
});
} catch (Exception e) {
log.info("创建链接失败!E:" + e);
}
}
}headers 了解即可
类似主题和直接交换机,可以根据 headers 中的内容来指定发送到那个队列
由于性能较差且复杂 不作使用
核心特性
消息过期机制
可以给每条消息指定一个有效期,一段时间内未被消费者处理,就过期了
使用场景:清理过期数据、模拟延迟队列的实现,专门让某个程序去处理过期请求
给队列中的所有消息指定过期时间
1
2
3
4
5// 创建队列,指定消息过期参数
Map<String,Object> args = new HasMap<>();
args.put("x-message-ttl",5000);
// args指定参数
channel.queueDeclare(QUEUE_NAME,false,false,false,args);
如果在过期时间内,还未收到消费者来取消息,消息才会过期
如果消息已经接收到,但是没有确认,是不会过期的如果消息处于待消费状态并且到达过期时间后,消息将会被标记为过期。但是,如果消息已经被消费者消费,并且正在被处理中,即使过期时间到了,消息依旧会被正常处理
消费者代码:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44package com.bi.spring.demomq;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.DeliverCallback;
import lombok.extern.slf4j.Slf4j;
import java.nio.charset.StandardCharsets;
import java.util.HashMap;
import java.util.Map;
public class TtlConsumer {
private static final String QUEUE_NAME = "ttl_queue";
public static void main(String[] args) {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
try {
// 从工厂中创建连接
Connection connection = factory.newConnection();
// 获取通信通道
Channel channel = connection.createChannel();
Map<String, Object> canshu = new HashMap<>();
canshu.put("x-message-ttl", 5000);
// 声明一个队列 没有就会创建 有的话就不再创建 放入参数
channel.queueDeclare(QUEUE_NAME, false, false, false, canshu);
System.out.println("Consumer waiting for producer to send message");
// 定义如何处理消息 这个是函数式接口
DeliverCallback deliverCallback = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
System.out.println("Consumer received message from producer,message:" + message);
};
// 消费消息
channel.basicConsume(QUEUE_NAME, true, deliverCallback, consumer -> {
});
} catch (Exception e) {
log.info("创建链接失败!E:" + e);
}
}
}生产者代码就正常取出消息消费即可,不作代码演示。
给某条消息指定过期时间
语法
1
2AMQP.BasicProperties properties = new AMQP.BasicProperties().builder().expiration("1000").build();
channel.basicPublish("", QUEUE_NAME, properties, message.getBytes(StandardCharsets.UTF_8));消息确认机制
为了保证消息成功被消费,rabbitMQ 提供了消息确认机制,当消费者接收到消息后,要给予反馈:
ack:消费成功
nack:消费失败
reject:拒绝只有告诉 rabbitMQ 服务器消费成功,服务器才会放心的移除消息
支持配置 autoack,会自动执行 ack 命令,接收到消息就立刻成功
一般情况下 建议 autoack 设置为 false 要根据实际情况,去手动确认指定确认某条消息:
1
channel.basicAck(delivery.getEnvelope().getDeliveryTag(),false);
指定拒绝某条消息
第三个参数表示是否要重试1
channel.basicNack(delivery.getEnvelope().getDeliveryTag(),false,false);
死信队列
官方文档: https://www.rabbitmq.com/dlx.html
为了保证消息的可靠性,比如每条消息都成功消费,需要提供一个容错机制,即: 失败的消息怎么处理?
死信: 过期的消息、拒收的消息、消息队列满了、处理失败的消息的统称死信队列:专门处理死信的队列
死信队列:专门处理死信的队列(注意,它就是一个普队列,只不过是专门用来处理死信的,你甚至可以理解这个队列的名称叫“死信队列”)
死信交换机:专门给死信队列转发消息的交换机(注意,它就是一个普通交换机,只不过是专门给死信队列发消息而已,理解为这个交换机的名称就叫 “死信交换机”)。也存在路由绑定
死信可以通过死信交换机绑定到死信队列。

1.创建死信交换机并绑定关系
2.给失败之后需要容错处理的消息队列绑定死信交换机
1
2
3
4
5
6
7
8
9
10// 指定死信队列参数
Map<String,Object> args = new HashMap<>();
// 要绑定到哪个交换机 等于额外绑定一个死信队列 处理特殊情况
args.put("x-dead-letter-exchange",DEAD EXCHANGE NAME);
// 指定死信要转发到哪个死信队列
args.put("x-dead-letter-routing-key","sixin_key");
// 创建队列,随机分配一个队列名称
String queueName = "sixi_queue";
channel.queueDeclare(queueName, true, false, false, args);
channel.queueBind(queueName, EXCHANGE NAME,"sixinkey_1");- 3.可以给要容错的队列指定死信之后的转发规则,死信应该再换发到那个死信队列
1
args.put("x-dead-letter-routing-key","sixin_de_queue")
- 4.可以通过程序来读取死信队列中的消息,从而进行处理
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33Connection connection = factory.newConnection();
// 获取通信通道
Channel channel = connection.createChannel();
// 声明死信交换机
channel.exchangeDeclare(DEAD_EXCHANGE_NAME, "direct");
// 创建死信队列
channel.queueDeclare(DEAD_QUEUE_NAME1, true, false, false, null);
// 将死信队列绑定到死信交换机并授予routingKey
channel.queueBind(DEAD_QUEUE_NAME1, DEAD_EXCHANGE_NAME, "routing-key-dead-queue-one");
// 创建死信队列2
channel.queueDeclare(DEAD_QUEUE_NAME2, true, false, false, null);
// 将死信队列绑定到死信交换机并授予routingKey 2
channel.queueBind(DEAD_QUEUE_NAME2, DEAD_EXCHANGE_NAME, "routing-key-dead-queue-two");
// 开启死信队列的消费监听 1
DeliverCallback deliverCallback1 = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
// 拒绝消息
channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, false);
// 打印接收消息
System.out.println("[1] Dead receive message:" + message + "-key:" + delivery.getEnvelope().getRoutingKey());
};
channel.basicConsume(DEAD_QUEUE_NAME1, false, deliverCallback1, consumerTag -> {
});
// 开启死信队列的消费监听 2
DeliverCallback deliverCallback2 = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
// 拒绝消息
channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, false);
// 打印接收消息
System.out.println("[2] Dead receive message:" + message + "-key:" + delivery.getEnvelope().getRoutingKey());
};
channel.basicConsume(DEAD_QUEUE_NAME2, false, deliverCallback2, consumerTag -> {
});- 5.生产者代码
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76package com.bi.spring.demomq;
import com.rabbitmq.client.*;
import lombok.extern.slf4j.Slf4j;
import java.nio.charset.StandardCharsets;
import java.util.Scanner;
public class DlxDirectProducer {
private static final String DEAD_EXCHANGE_NAME = "dlx-direct-exchange";
public static final String DEAD_QUEUE_NAME1 = "dlx-one-queue";
public static final String DEAD_QUEUE_NAME2 = "dlx-two-queue";
private static final String WORK_EXCHANGE_NAME = "work-direct-exchange";
public static void main(String[] args) {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
try {
// 从工厂中创建连接
Connection connection = factory.newConnection();
// 获取通信通道
Channel channel = connection.createChannel();
// 声明死信交换机
channel.exchangeDeclare(DEAD_EXCHANGE_NAME, "direct");
// 创建死信队列
channel.queueDeclare(DEAD_QUEUE_NAME1, true, false, false, null);
// 将死信队列绑定到死信交换机并授予routingKey
channel.queueBind(DEAD_QUEUE_NAME1, DEAD_EXCHANGE_NAME, "routing-key-dead-queue-one");
// 创建死信队列2
channel.queueDeclare(DEAD_QUEUE_NAME2, true, false, false, null);
// 将死信队列绑定到死信交换机并授予routingKey 2
channel.queueBind(DEAD_QUEUE_NAME2, DEAD_EXCHANGE_NAME, "routing-key-dead-queue-two");
// 开启死信队列的消费监听 1
DeliverCallback deliverCallback1 = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
// 拒绝消息
channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, false);
// 打印接收消息
System.out.println("[1] Dead receive message:" + message + "-key:" + delivery.getEnvelope().getRoutingKey());
};
channel.basicConsume(DEAD_QUEUE_NAME1, false, deliverCallback1, consumerTag -> {
});
// 开启死信队列的消费监听 2
DeliverCallback deliverCallback2 = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
// 拒绝消息
channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, false);
// 打印接收消息
System.out.println("[2] Dead receive message:" + message + "-key:" + delivery.getEnvelope().getRoutingKey());
};
channel.basicConsume(DEAD_QUEUE_NAME2, false, deliverCallback2, consumerTag -> {
});
// 生产消息
Scanner scanner = new Scanner(System.in);
while (scanner.hasNext()) {
String userInput = scanner.nextLine();
String[] s = userInput.split(" ");
if (s.length < 1) {
continue;
}
String message = s[0];
String routingKey = s[1];
channel.basicPublish(WORK_EXCHANGE_NAME, routingKey, null, message.getBytes(StandardCharsets.UTF_8));
System.out.println("[x] Send message:" + message + " with routingKey:" + routingKey);
}
} catch (Exception e) {
log.info("创建链接失败!E:" + e);
}
}
}- 6.消费者代码
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76package com.bi.spring.demomq;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.DeliverCallback;
import lombok.extern.slf4j.Slf4j;
import java.nio.charset.StandardCharsets;
import java.util.HashMap;
import java.util.Map;
public class DlxDirectConsumer {
private static final String DEAD_EXCHANGE_NAME = "dlx-direct-exchange";
public static final String WORK_QUEUE_NAME1 = "work-dlx-one-queue";
public static final String WORK_QUEUE_NAME2 = "work-dlx-two-queue";
private static final String WORK_EXCHANGE_NAME = "work-direct-exchange";
public static void main(String[] args) {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("localhost");
try {
// 从工厂中创建连接
Connection connection = factory.newConnection();
// 获取通信通道
Channel channel = connection.createChannel();
// 声明工作交换机
channel.exchangeDeclare(WORK_EXCHANGE_NAME, "direct");
// 指定死信队列参数
Map<String, Object> parameter = new HashMap<>();
// 绑定交换机
parameter.put("x-dead-letter-exchange", DEAD_EXCHANGE_NAME);
// 指定死信要转发到那个队列
parameter.put("x-dead-letter-routing-key", "routing-key-dead-queue-one");
// 创建队列1 以参数形式多绑定一个死信交换机
channel.queueDeclare(WORK_QUEUE_NAME1, true, false, false, parameter);
channel.queueBind(WORK_QUEUE_NAME1, WORK_EXCHANGE_NAME, "work-queue-one");
// 同样也给2进行绑定
Map<String, Object> parameter2 = new HashMap<>();
parameter2.put("x-dead-letter-exchange", DEAD_EXCHANGE_NAME);
parameter2.put("x-dead-letter-routing-key", "routing-key-dead-queue-two");
// 创建队列2 以参数形式多绑定一个死信交换机
channel.queueDeclare(WORK_QUEUE_NAME2, true, false, false, parameter);
channel.queueBind(WORK_QUEUE_NAME2, WORK_EXCHANGE_NAME, "work-queue-two");
System.out.println("Consumer waiting for producer to send message");
// 开启工作队列的消费监听 1
DeliverCallback deliverCallback1 = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
// 拒绝消息
channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, false);
// 打印接收消息
System.out.println("[1] Dead receive message:" + message + "-key:" + delivery.getEnvelope().getRoutingKey());
};
channel.basicConsume(WORK_QUEUE_NAME1, false, deliverCallback1, consumerTag -> {
});
// 开启工作队列的消费监听 2
DeliverCallback deliverCallback2 = (consumerTag, delivery) -> {
String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
// 拒绝消息
channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, false);
// 打印接收消息
System.out.println("[2] Dead receive message:" + message + "-key:" + delivery.getEnvelope().getRoutingKey());
};
channel.basicConsume(WORK_QUEUE_NAME2, false, deliverCallback2, consumerTag -> {
});
} catch (Exception e) {
log.info("创建链接失败!E:" + e);
}
}
}
rabbit 重点知识
消息队列的概念、模型、应用场景
交换机的类别、路由绑定关系
消息可靠性
消息确认机制(ack、nack、reject)
消息持久化(durable)
消息过期机制
死信队列延迟队列 类似于 死信队列
顺序消费、消费幂等性(做了解)
可扩展性
集群
故障的回复机制
镜像运维监控告警(做了解)
RabbitMQ 项目实战
项目中如何使用 RabbitMQ?
- 使用官方的客户端
优点:兼容性好,换语言成本低,比较灵活
缺点:太灵活,要自己去做一些限制或者是一些事情。比如要自己维护管理连接等等。- 使用封装好的客户端,比如 Spring Boot RabbitMQ Starter
优点:简单易用,直接配置直接用,更方便地去管理连接
缺点:封装的太好了,没基础知识很难去使用或者是配置基础实战
引入依赖
引入时候一定要选择与自己当前使用的 Spring Boot 版本一致
1
2
3
4
5
6<!-- https://mvnrepository.com/artifact/org.springframework.boot/spring-boot-starter-amqp -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
<version>2.7.2</version>
</dependency>在 yml 中进行配置
1
2
3
4
5
6spring:
rabbmit:
host: localhost
port: 5672
username: guest
passowrd: guest创建交换机和队列
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38package com.bi.spring.bizmq;
import com.bi.spring.config.RabbitMQConfig;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component;
import org.springframework.stereotype.Service;
/**
* 用于创建测试程序用到的交换机和队列(只用在程序启动前执行一次)
*/
public class MqInitMain {
private static final String EXCHANGE_NAME = "bi-exchange";
public static void main(String[] args) {
try {
// 创建连接
ConnectionFactory factory = new ConnectionFactory();
// 设置连接属性
factory.setHost("localhost");
Connection connection = factory.newConnection();
Channel channel = connection.createChannel();
channel.exchangeDeclare(EXCHANGE_NAME, "direct");
// 创建队列,随机分配一个队列名称
String queueName = "test_queue";
channel.queueDeclare(queueName, true, false, false, null);
channel.queueBind(queueName, EXCHANGE_NAME, "bi-queue-key");
} catch (Exception e) {
log.info("初始化失败:" + e);
}
}
}- 生产者代码
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17package com.bi.spring.bizmq;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
public class MyMessageProducer {
private RabbitTemplate rabbitTemplate;
public void sendMessage(String exchange, String routingKey, String message) {
rabbitTemplate.convertAndSend(exchange, routingKey, message);
}
}- 消费者代码
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30package com.bi.spring.bizmq;
import com.rabbitmq.client.Channel;
import lombok.SneakyThrows;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.amqp.support.AmqpHeaders;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.stereotype.Component;
public class MyMessageConsumer {
/**
* 指定程序监听的消息队列和确认机制
*
* @param message
* @param channel
* @param deliveryTag
*/
public void receiveMessage(String message, Channel channel, long deliveryTag) {
log.info("receiveMessage message = {}", message);
// 确认收到消息
channel.basicAck(deliveryTag, false);
}
}- 单元测试
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20package com.bi.spring.bizmq;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
class MyMessageConsumerTest {
private MyMessageProducer myMessageProducer;
void receiveMessage() {
myMessageProducer.sendMessage("bi-exchange", "bi-queue-key", "消息进来了");
}
}- 测试结果

智能 BI 改造
按照上一期的做法,是把用户提交的任务放入本地线程池中,在线程池内排队,但是程序如果中断了,任务就会丢失,就丢了,并且还不会去修改状态。
实现步骤
创建交换机和队列 程序启动前要先启动一次初始化
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35package com.bi.spring.bizmq;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
/**
* 初始化MQ
*/
public class BiInitMain {
public static void main(String[] args) {
try {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("xxxx");
factory.setPort(xxxx);
factory.setUsername("xxxx");
factory.setPassword("xxxx");
Connection connection = factory.newConnection();
Channel channel = connection.createChannel();
// 声明交换机 类型为直接交换机
channel.exchangeDeclare(BiMqConstant.BI_EXCHANGE_NAME, "direct");
// 声明队列
channel.queueDeclare(BiMqConstant.BI_QUEUE_NAME, true, false, false, null);
// 绑定队列到交换机
channel.queueBind(BiMqConstant.BI_QUEUE_NAME, BiMqConstant.BI_EXCHANGE_NAME, BiMqConstant.BI_ROUTING_KEY);
// 限制任务数
channel.basicQos(1);
} catch (Exception e) {
System.out.println("MQ初始化失败!");
}
}
}

重构 async 接口为新的 API
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
private BiMessageProducer biMessageProducer;
/**
* AI异步(消息队列)对话分析
*
* @param multipartFile
* @param genChartByAiRequest
* @param request
* @return
*/
public BaseResponse<String> genChartByAIUseAsyncMQ( MultipartFile multipartFile,
GenChartByAiRequest genChartByAiRequest, HttpServletRequest request) {
// 参数获取
String goal = genChartByAiRequest.getGoal();
String name = genChartByAiRequest.getName();
String chartType = genChartByAiRequest.getChartType();
// 参数校验
ThrowUtils.throwIf(StringUtils.isEmpty(goal), ErrorCode.PARAMS_ERROR, "目标不能为空");
ThrowUtils.throwIf(StringUtils.isEmpty(chartType), ErrorCode.PARAMS_ERROR, "类型不能为空");
ThrowUtils.throwIf(StringUtils.isEmpty(name), ErrorCode.PARAMS_ERROR, "图表名称不能为空");
// 文件校验
long size = multipartFile.getSize();
String originalFilename = multipartFile.getOriginalFilename();
ThrowUtils.throwIf(size > MAX_FILE_SIZE, ErrorCode.PARAMS_ERROR, "文件过大不得超过1MB");
ThrowUtils.throwIf(!FILE_NAME_LIST.contains(FileUtil.getSuffix(originalFilename)), ErrorCode.PARAMS_ERROR, "不支持此文件");
// 判断是否登录
User loginUser = userService.getLoginUser(request);
ThrowUtils.throwIf(Objects.isNull(loginUser), ErrorCode.NOT_LOGIN_ERROR);
// 对用户进行限流
redisLimiterManager.doRateLimit("genChartByAI_" + loginUser.getId());
// 用户消息拼接
String dataStr = ExcelUtils.excelToCsv(multipartFile);
// 先保存任务到数据库中
Chart chart = new Chart();
chart.setName(name);
chart.setGoal(goal);
chart.setChartData(dataStr);
chart.setChartType(chartType);
chart.setStatus("wait");
chart.setExecMessage("任务等待中");
chart.setUserId(loginUser.getId());
boolean firstChartSave = chartService.save(chart);
ThrowUtils.throwIf(!firstChartSave, ErrorCode.PARAMS_ERROR, "任务提交失败");
// 任务如果保存成功 将ID步入消息队列中
biMessageProducer.sendMessage(String.valueOf(chart.getId()));
return ResultUtils.success("任务提交成功!");
}将线程池中的代码进行重构
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74/**
* 监听消息消费
*
* @param message
* @param channel
* @param deliverTag
*/
public void receiveMessage(String message, Channel channel, long deliverTag) {
try {
log.info("receiveMessage: {}", message);
if (StringUtils.isBlank(message)) {
// 参数 消息的tag 是否要全部拒绝 false为只拒绝当前 是否要重入队列 false为不入
channel.basicNack(deliverTag, false, false);
throw new BusinessException(ErrorCode.SYSTEM_ERROR);
}
Long id = Long.parseLong(message);
Chart chart = chartService.getById(id);
if (Objects.isNull(chart)) {
channel.basicNack(deliverTag, false, false);
throw new BusinessException(ErrorCode.SYSTEM_ERROR);
}
// 将任务转入任务队列中分配执行
CompletableFuture.runAsync(() -> {
Chart updateFirstChart = new Chart();
updateFirstChart.setId(chart.getId());
updateFirstChart.setStatus("running");
updateFirstChart.setExecMessage("任务正在执行中");
boolean isFirstUpdate = chartService.updateById(updateFirstChart);
if (!isFirstUpdate) {
try {
channel.basicNack(deliverTag, false, false);
} catch (IOException e) {
log.error("消息拒绝失败!");
}
handleChartUpdateError(chart.getId());
return;
}
// AI接口服务
String content = aiManager.doChat(buildUserMessage(chart));
String[] splits = content.split("【【【【【");
if (splits.length != 3) {
try {
channel.basicNack(deliverTag, false, false);
} catch (IOException e) {
log.error("消息拒绝失败!");
}
throw new RuntimeException("AI 生成错误");
}
Chart updateChartData = new Chart();
updateChartData.setId(chart.getId());
updateChartData.setStatus("succeed");
updateFirstChart.setExecMessage("任务已经完成");
updateChartData.setGenResult(splits[2].trim());
updateChartData.setGenChart(splits[1].trim());
boolean isUpdateChartData = chartService.updateById(updateChartData);
if (!isUpdateChartData) {
try {
channel.basicNack(deliverTag, false, false);
} catch (IOException e) {
log.error("消息拒绝失败!");
}
handleChartUpdateError(chart.getId());
return;
}
}, threadPoolExecutor);
// 成功执行后确认消息
channel.basicAck(deliverTag, false);
} catch (IOException e) {
log.error("消息消费失败!");
throw new BusinessException(ErrorCode.SYSTEM_ERROR);
}
}代码中仍然存在缺陷,被拒绝的消息要进行处理,并要修改任务状态
修改初始化 MQ 代码
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49package com.bi.spring.bizmq;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import java.util.HashMap;
import java.util.Map;
/**
* 初始化MQ
*/
public class BiInitMain {
public static void main(String[] args) {
try {
ConnectionFactory factory = new ConnectionFactory();
factory.setHost("124.222.153.56");
factory.setPort(5672);
factory.setUsername("admin");
factory.setPassword("admin");
Connection connection = factory.newConnection();
Channel channel = connection.createChannel();
// 声明死信交换机 类型为直接交换机
channel.exchangeDeclare(BiMqConstant.BI_DEAD_EXCHANGE_NAME, "direct");
// 声明死信队列
channel.queueDeclare(BiMqConstant.BI_DEAD_QUEUE_NAME, true, false, false, null);
// 声明交换机 类型为直接交换机
channel.exchangeDeclare(BiMqConstant.BI_EXCHANGE_NAME, "direct");
// 指定死信队列的参数
Map<String, Object> parameter = new HashMap<>();
parameter.put("x-dead-letter-exchange", BiMqConstant.BI_DEAD_EXCHANGE_NAME);
parameter.put("x-dead-letter-routing-key", BiMqConstant.BI_DEAD_ROUTING_KEY);
// 指定过期时间
parameter.put("x-message-ttl", 3000L);
// 声明队列
channel.queueDeclare(BiMqConstant.BI_QUEUE_NAME, true, false, false, parameter);
// 绑定队列到交换机
channel.queueBind(BiMqConstant.BI_QUEUE_NAME, BiMqConstant.BI_EXCHANGE_NAME, BiMqConstant.BI_ROUTING_KEY);
// 绑定死信交换机和死信队列的routingKey
channel.queueBind(BiMqConstant.BI_DEAD_QUEUE_NAME, BiMqConstant.BI_DEAD_EXCHANGE_NAME, BiMqConstant.BI_DEAD_ROUTING_KEY);
// 限制任务数
channel.basicQos(1);
} catch (Exception e) {
System.out.println("MQ初始化失败!");
}
}
}添加死信队列消费监听
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32/**
* 处理因为某种情况没修改状态
*
* @param message
* @param channel
* @param deliverTage
*/
public void receiveDeadMessage(String message, Channel channel, long deliverTage) {
try {
log.info("receiveMessage dead message:{}", message);
if (StringUtils.isBlank(message)) {
channel.basicAck(deliverTage, false);
} else {
Chart chart = chartService.getById(Long.parseLong(message));
if (!Objects.isNull(chart)) {
chart.setStatus("failed");
chart.setExecMessage("更新图表状态失败!");
boolean save = chartService.save(chart);
if (!save) {
// 再次失败 就判定为系统异常
log.error("更新图表失败状态失败" + chart.getId() + "," + "更新图表状态失败!");
throw new BusinessException(ErrorCode.SYSTEM_ERROR);
}
}
}
channel.basicAck(deliverTage, false);
} catch (IOException e) {
log.error("死信消息消费失败!");
throw new BusinessException(ErrorCode.SYSTEM_ERROR);
}
}
验证
验证发现,如果程序中断了,并且消息未被手动确认或者拒绝,也就是没有 ack 或者 nack 无任何反应
那么这条消息就会重新处于 ready 状态,系统恢复时消费者依然会监听到这条消息并重新消费。
性能优化
改造之后
写了一个专门来接受消息的程序,处理任务
如果程序中断了,消息未被确认,就会重发
消息集中的发到消息队列,可以部署多个后端,从同一个地方取出消息做处理,实现了分布式的负载均衡
今日后端收获
后端项目到此就完成了
- 完成了 由 同步化->异步化->异步化(消息队列)的优化过程
- 学习了如何使用本地线程池来进行异步化接口,从而增强了接口的响应速度,熟悉了线程池参数的意义
- 学习了使用消息中间件 RbbitMQ 来作为消息队列,从原本的本地线程池导致数据丢失,生成状态由于程序的中断而导致丢失,卡状态无法反馈的问题。
- 此项目中收获最大的就是如何对业务进行解耦拆分,把核心处理业务和其他三方服务或者是耗时长的服务抽离,用户只需要提交任务,后端慢慢的去处理
- 初步的使用消息中间件接入项目中,从听理论到结尾、模糊到分析清楚,花费了大量时间。因此此项目对于我收获非常的巨大。
- 特别感谢鱼总能够分享技术,带领学习。
今日前端开发
更新后端开发文档 执行 openAPI

创建新的页面 修改请求 API
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22const onFinish = async (values: any) => {
if (submitting) return;
setSubmitting(true);
const params = {
...values,
file: undefined,
};
try {
const res = await genChartByAIUseAsyncMQUsingPOST(
params,
{},
values?.file[0]?.originFileObj
);
if (res?.data) {
message.success(res.data);
form.resetFields();
}
} catch (e: any) {
message.error("分析失败!");
}
setSubmitting(false);
};提交后 DEBUG 观察消息是否进入后端

到我的图表中观察状态


今日前端收获
由于前端不是侧重点,最后一期主要是对于后端的开发和优化
页面架构都是前期搭建好的直接复制过来进行修改即可。

