前言
我們可以使用SpringCloud框架中Feign組完成微服務之間的遠程呼叫;
但是Feign組件底層基于HTTP協議,HTTP協議的特點是請求同步,而且既需要請求也需要回應,屬于同步遠程呼叫;
微服務架構在同步遠程呼叫的場景下,如果服務提供者一直沒有回應服務消費者,很容易造成服務雪崩;
如果我們通過MQ協議發送異步訊息,就可以實作服務消費者和服務提供者之間的應用解耦,中間的佇列還能起到流量削峰的作用,進而實作高效、可靠的異步遠程呼叫;
那么MQ能完全替代Feign嗎?
2種方式各有優劣,打電話可以立即得到回應,但是你卻不能跟多個人同時通話,發送訊息可以同時與多個人溝通,但是往往回應會有延遲;
同步遠程呼叫的優勢是服務消費者會一直等待服務提供者的回應訊息;
一旦服務提供者的回應訊息回傳,服務消費者會在第一時間得知,在接收服務提供者回應訊息方面,同步遠程呼叫比異步遠程呼叫更加及時;
所以在企業中我們要根據不同的業務需求,靈活選擇同步遠程呼叫和異步遠程呼叫;
一、MQ概念
MQ (Message Queue),中文是訊息佇列,字面來看就是存放訊息的佇列,
它是分布式系統中重要的組件,主要解決應用解耦,流量削峰,異步訊息等問題,
常見的角色有:Producer(生產者)、Consumer(消費者)、Broker(中介),
1.常見的訊息佇列
| RabbitMQ | ActiveMQ | RocketMQ | Kafka | |
|---|---|---|---|---|
| 公司/社區 | Rabbit | Apache | 阿里 | Apache |
| 開發語言 | Erlang | Java | Java | Scala&Java |
| 協議支持 | AMQP,XMPP,SMTP,STOMP | OpenWire,AMQP,STOMP,MQTT | 自定義協議 | 自定義協議 |
| 可用性 | 高 | 一般 | 高 | 高 |
| 單機吞吐量 | 一般 | 差 | 高 | 非常高 |
| 訊息延遲 | 微秒級 | 毫秒級 | 毫秒級 | 毫秒以內 |
| 訊息可靠性 | 高 | 一般 | 高 | 一般 |
2.常見訊息佇列的使用場景
- RocketMQ是阿里巴巴開發的MQ其特點是訊息可靠性較高適用于電商專案;
- Kafka傳輸的是資料流,雖然訊息可靠性較低,但是吞吐量非常高,適用于大資料專案;
- RabbitMQ的各方面性能適中適用于中小型專案;
二、RabbitMQ概念
RabbitMQ是基于AMQP(Advanced Message Queuing Protocol)協議的一款訊息中間件管理系統;
官網地址:http://www.rabbitmq.com/
官方教程:http://www.rabbitmq.com/getstarted.html
1.安裝RabbitMQ
在Centos7虛擬機中使用Docker來安裝RabbitMQ
#1.下載rabbitmq鏡像 docker pull rabbitmq:3.8-management #2.運行容器 docker run -d -p 15672:15672 -p 5672:5672 --name mq -v mq-plugins:/plugins --hostname mq rabbitmq:3.8-management
設定5672埠提供API服務,15672埠提供用戶管理界面;
2.創建主機
為了讓各個用戶可以互不干擾的作業,RabbitMQ添加了虛擬主機(Virtual Hosts)的概念;
虛擬主機其實就是1個獨立的訪問路徑,每1個用戶使用不同路徑,每1個路徑中包含多個佇列、交換機;
虛擬主機和虛擬主機之間互相隔離不會影響;
3.創建用戶

4.賦予用戶主機管理權限
將創建好的zhanggen虛擬主機交給用戶zhanggen管理;
5.兩大基本訊息傳輸模型
訊息佇列的訊息分為2大類傳輸模型:點對點模型、發布 /訂閱模型;
- 點對點: 生產者生產的同1條訊息只能被1個消費者消費;(微信私聊)
- 發布/訂閱:生產者生產的同1條訊息可以被多個消費者消費;(微信群聊)
5.1.為什么有了點對點還需要發布/訂閱訊息傳輸模型?
因為生產者1次生產的1條訊息只能被1個消費者消費;
為了避免訊息重復冗余才需要發布/訂閱訊息傳輸模型,把生產者1次生產的1條訊息,轉發給N個人;
就像你家丟了一條小狗,挨家挨戶地去問(點對點),不如在大喇叭里吆喝廣播1下(發布/訂閱)省時省力;
6.六大訊息傳輸模型
RabbitMQ在2基本傳輸模型的基礎上進行細化,提供了6種訊息模型,但是第6種其實是RPC,并不是MQ,因此不予學習,那么也就剩下5種;
-
1、2(點對點模型)
-
3、4、5(發布/訂閱模型)

三、Java呼叫RabbitMQ
生產者負責想訊息物件的某個佇列中發送訊息,生產者向訊息佇列發送訊息成功之后程式退出;
消費者啟動之后會一直監聽訊息佇列的某個佇列是否有訊息,程式不會退出;
1.依賴引入
<dependencies> <!--AMQP依賴,包含RabbitMQ--> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-amqp</artifactId> </dependency> <!--單元測驗--> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-test</artifactId> </dependency> <dependency> <groupId>org.projectlombok</groupId> <artifactId>lombok</artifactId> </dependency> </dependencies>pom.xml
2.生產者
package com.itheima.test; import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory; import org.junit.Test; import java.io.IOException; import java.util.concurrent.TimeoutException; public class PublisherTest { @Test public void testSendMessage() throws IOException, TimeoutException { // 1.建立連接 ConnectionFactory factory = new ConnectionFactory(); factory.setHost("192.168.56.18"); factory.setPort(5672); factory.setVirtualHost("zhanggen"); factory.setUsername("zhanggen"); factory.setPassword("123.com"); Connection connection = factory.newConnection(); // 2.創建通道Channel Channel channel = connection.createChannel(); // 3.創建佇列 String queueName = "p2p";//佇列名稱 channel.queueDeclare(queueName, false, false, false, null); // 4.發送訊息 String message = "hello, rabbitmq!"; channel.basicPublish("", queueName, null, message.getBytes()); System.out.println("發送訊息成功:【" + message + "】"); // 5.關閉通道和連接 channel.close(); connection.close(); } }PublisherTest.java
3.消費者
package com.itheima.test; import com.rabbitmq.client.*; import java.io.IOException; import java.util.concurrent.TimeoutException; public class ConsumerTest { public static void main(String[] args) throws IOException, TimeoutException { // 1.建立連接 ConnectionFactory factory = new ConnectionFactory(); factory.setHost("192.168.56.18"); factory.setPort(5672); factory.setVirtualHost("zhanggen"); factory.setUsername("zhanggen"); factory.setPassword("123.com"); Connection connection = factory.newConnection(); // 2.創建通道Channel Channel channel = connection.createChannel(); // 3.創建佇列 String queueName = "p2p"; channel.queueDeclare(queueName, false, false, false, null); // 4.訂閱訊息 channel.basicConsume(queueName, true, new DefaultConsumer(channel) { @Override public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException { // 5.處理訊息 String message = new String(body); System.out.println("接收到訊息:【" + message + "】"); } }); System.out.println("消費者等待接收訊息,,,,"); } }ConsumerTest.java
四、SringBoot呼叫RabbitMQ
AMQP:是用于在應用程式之間傳遞業務訊息的開放標準,該協議與語言和平臺無關,更符合微服務中獨立性的要求,
Spring AMQP:基于AMQP協議定義的一套API,提供了模板來發送和接收訊息,模板底層是基于RabbitMQ封裝,
SpringAmqp的官方地址:https://spring.io/projects/spring-amqp
以下我們將借助SpringAmqp實作RabbitMQ的5中訊息傳輸模型;
1.點對點訊息傳輸模型-BasicQueue
簡單訊息模型,1個生產者和1個消費者進行點對點訊息傳輸;

1.1.生產者
1.1.1.application.yml配置
在application.yml中添加MQ配置
#RabbitMQ相關配置 spring: rabbitmq: host: 192.168.56.18 # 主機名 port: 5672 # 埠 virtual-host: zhanggen # 虛擬主機 username: zhanggen # 用戶名 password: 123.com # 密碼
1.1.2.測驗類
撰寫測驗類SpringAmqpTest,并利用RabbitTemplate實作訊息發送
package com.zhanggen.test; import org.junit.Test; import org.junit.runner.RunWith; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.test.context.SpringBootTest; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; @SpringBootTest @RunWith(SpringJUnit4ClassRunner.class) public class SpringAmqpTest { @Autowired private RabbitTemplate rabbitTemplate; //發送簡單訊息 @Test public void testSimpleQueue() { //引數一: 佇列名稱(次佇列需要是提前創建好的) 引數二: 訊息內容 rabbitTemplate.convertAndSend("p2p", "hello,spring amqp!"); } }
1.2.消費者
1.2.1.application.yml配置
#RabbitMQ相關配置
spring:
rabbitmq:
host: 192.168.56.18 # 主機名
port: 5672 # 埠
virtual-host: zhanggen # 虛擬主機
username: zhanggen # 用戶名
password: 123.com # 密碼
1.2.2.測驗類
package com.zhanggen.listener; import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.stereotype.Component; @Component public class SpringRabbitListener { //簡單型別訊息 @RabbitListener(queues = "p2p")//宣告佇列名稱 public void listenSimpleQueueMessage(String msg) { System.out.println("消費者接收到訊息:【" + msg + "】"); } }
2.點對點訊息傳輸模型-WorkQueue
當訊息處理比較耗時的時候,可能生產訊息的速度會遠遠大于訊息的消費速度,長此以往,訊息就會堆積越來越多造成佇列積壓,

此時就可以使用work 模型,多個消費者共同處理訊息處理,訊息的消費速度就能大大提高了,
2.2.生產者
這次我們回圈發送,模擬大量訊息堆積現象,
在publisher服務中的SpringAmqpTest類中添加一個測驗方法:
//批量發送訊息 @Test public void testWorkQueue() { for (int i = 1; i <= 50; i++) { rabbitTemplate.convertAndSend("p2p", "message_" + i); } }
2.3.消費者
要模擬多個消費者系結到同1個佇列,我們在consumer服務的SpringRabbitListener中添加2個新的方法:
//使用下面兩個方法來接收p2p佇列中的訊息 @RabbitListener(queues = "p2p") public void listenWorkQueue1(String msg) throws InterruptedException { System.out.println("消費者1接收到訊息:【" + msg + "】"); Thread.sleep(20); } @RabbitListener(queues = "p2p") public void listenWorkQueue2(String msg) throws InterruptedException { System.err.println("消費者2........接收到訊息:【" + msg + "】"); Thread.sleep(200); }
2.4.測驗
先啟動ConsumerApplication后(消費者),再執行publisher服務中剛剛撰寫的發送測驗方法testWorkQueue(提供者),
可以看到消費者1很快完成了自己的25條訊息,費者2卻在緩慢的處理自己的25條訊息,
也就是說RabbitMQ按訊息的總條數,平均分配給2個消費者,并沒有考慮到消費者的處理能力,這樣顯然是有問題的,

2.5.消費者能者多勞
在spring中有一個簡單的配置,可以解決以上問題,我們修改consumer服務的application.yml檔案,添加配置:
spring:
rabbitmq:
listener:
simple:
prefetch: 1 # 消費者一次處理一條訊息,處理完畢后再從MQ中獲取
2.6.測驗消費者能者多勞

2.7.小結
以上2中訊息傳輸模型都屬于點對點傳輸模型
Work模型的使用:
-
多個消費者系結到同1個佇列,同1條訊息只會被1個消費者處理
-
通過設定prefetch來控制消費者預取的訊息數量
3.發布/訂閱訊息傳輸模型-Fanout
以上介紹了2種點對點訊息傳輸模型,以下介紹3種發布/訂閱訊息傳輸模型;
在下面的發布/訂閱訊息傳輸模型中,需要使用Exchanges轉發訊息到多個Queue,不再局限于使用單個Queue傳輸訊息;
3.1.點對對訊息傳輸模型的缺陷
如果要實作以下功能:
- UserService作訊息生產者向訊息中間發送1條用戶資訊訊息;
- 郵件微服務和短信微服務都是用戶訊息的消費者,它們從訊息中間獲取用戶訊息之后,實作給用戶發郵件+發短信功能;

3.2.引入發布/訂閱訊息傳輸模型
由于點對對訊息傳輸模型的缺陷:生產者生產的1條訊息,只能被1個消費者所消費,導致上述的2個功能只能實作其中的1個;
1條用戶資訊要么被郵件微服務搶到,發了郵件,要么被短信微服務搶到,發了短信;
如果想讓生產者生產的1條訊息,同時被多個消費者同時搶到,即完成發郵件的功能,也要完成發短信的功能;
就需修改當前點對點訊息傳輸模型為發布/訂閱訊息傳輸模型;
發布/訂閱訊息傳輸模型(Fanout)在MQ中可以理解為廣播模式,

3.2.發布/訂閱訊息傳輸模型術語作業流程

3.3.宣告交換機和佇列
我們習慣在消費者一方來創建交換機和佇列;
package com.zhanggen.config; import org.springframework.amqp.core.Binding; import org.springframework.amqp.core.BindingBuilder; import org.springframework.amqp.core.FanoutExchange; import org.springframework.amqp.core.Queue; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; //關于Fanout訊息傳輸模型配置 @Configuration public class FanoutConfiguration { //1.配置1個交換機 @Bean public FanoutExchange fanoutExchange() { return new FanoutExchange("MyFanoutExchange"); } //2.配置2個佇列myQueue1和myQueue2,注意方法名稱就是就是bean的名稱 @Bean public Queue myQueue1() { return new Queue("MyQueue1"); } @Bean public Queue myQueue2() { return new Queue("MyQueue2"); } //3.將2個佇列(myQueue1和myQueue2)系結到交換機(MyFanoutExchange) @Bean public Binding bindQueue1(FanoutExchange fanoutExchange, Queue myQueue1) { return BindingBuilder.bind(myQueue1).to(fanoutExchange); } @Bean public Binding bindQueue2(FanoutExchange fanoutExchange, Queue myQueue2) { return BindingBuilder.bind(myQueue2).to(fanoutExchange); } }
3.4.消費者
//FanOut訊息傳輸模型 @RabbitListener(queues = "MyQueue1") public void listenFanOut1(String msg) throws InterruptedException { System.err.println("FanOut訊息傳輸模型消費者1........接收到訊息:【" + msg + "】"); Thread.sleep(200); } @RabbitListener(queues = "MyQueue2") public void listenFanOut2(String msg) throws InterruptedException { System.err.println("FanOut訊息傳輸模型消費者2........接收到訊息:【" + msg + "】"); Thread.sleep(200); }
3.5.生產者
//測驗FanOut @Test public void testFanOut() { //引數一: 交換機名稱 引數二:暫時沒用 引數三: 訊息內容 rabbitTemplate.convertAndSend("MyFanoutExchange", "", "交換機"); }
3.6.測驗

------------------------------------- ------------------------------------- -------------------------------------

4.發布/訂閱訊息傳輸模型-Direct
在Fanout訊息傳輸模型中新增了交換機和多佇列的概念,實作1條訊息可以被所有訂閱的佇列消費的功能,
但是只要佇列系結到了1個交換機上,佇列就能收到這個交換機廣播的所有訊息;
如果交換機把所有訊息廣播到了所有佇列里,不僅產生廣播風暴,也會有資訊安全隱患,這時就要用到Direct模型;
RoutingKey:生產者向Exchange發送訊息時,一般會指定一個RoutingKey;
BindingKey:當系結Exchange和Queue時,一般會指定一個BindingKey;
BindingKey與RoutingKey相匹配時,訊息將會被路由到對應的Queue中,
4.1.注解宣告消費者
基于@Bean的方式宣告佇列和交換機比較麻煩,Spring還提供了基于注解方式來宣告,
在consumer的SpringRabbitListener中添加兩個消費者,同時基于注解來宣告佇列和交換機:
//測驗Direct模型 @RabbitListener( bindings = @QueueBinding(//系結 exchange = @Exchange(value = "https://www.cnblogs.com/sss4/archive/2022/06/29/direct.exchange", type = ExchangeTypes.DIRECT),//設定交換機的名字和型別 value = https://www.cnblogs.com/sss4/archive/2022/06/29/@Queue("direct.queue1"),//設定對列名字 key = "base"//設定bindingkey ) ) public void listenDirectQueue1(String message) { System.out.println("消費者1接收到了Direct訊息:" + message); } @RabbitListener( bindings = @QueueBinding(//系結 exchange = @Exchange(value = "https://www.cnblogs.com/sss4/archive/2022/06/29/direct.exchange", type = ExchangeTypes.DIRECT),//設定交換機的名字和型別 value = https://www.cnblogs.com/sss4/archive/2022/06/29/@Queue("direct.queue2"),//設定對列名字 key = {"base", "vip"}//設定bindingkey ) ) public void listenDirectQueue2(String message) { System.out.println("消費者2接收到了Direct訊息:" + message); }
4.2.生產者
// 測驗direct @Test public void testSendDirect() throws Exception { //引數一: 交換機名稱 引數二:routingKey 引數三: 訊息內容 rabbitTemplate.convertAndSend("direct.exchange", "vip", "hello,everyone!"); }
5.發布/訂閱訊息傳輸模型-Topic
Topic訊息傳輸模型與Direct相比,新增了RoutingKey和BidingKey可使用通配符的方式進行匹配的功能;
RoutingKey:一般由有1個或多個單詞組成,多個單詞之間以”.”分割,例如:china.news
BindingKey:使用通配符匹配RoutingKey匹配規則如下:
#:代指匹配RoutingKey的0個或多個單詞*:代指匹配RoutingKey的1個單詞

5.1.消費者
//測驗topic模型 @RabbitListener(bindings = @QueueBinding( value = @Queue("topic.queue1"), exchange = @Exchange(value = "https://www.cnblogs.com/sss4/archive/2022/06/29/topic.exchange", type = ExchangeTypes.TOPIC), key = "china.#" )) public void listenTopicQueue1(String msg) { System.out.println("消費者1接收到Topic訊息:【" + msg + "】"); } @RabbitListener(bindings = @QueueBinding( value = @Queue("topic.queue2"), exchange = @Exchange(value = "https://www.cnblogs.com/sss4/archive/2022/06/29/topic.exchange", type = ExchangeTypes.TOPIC), key = "#.weather" )) public void listenTopicQueue2(String msg) { System.out.println("消費者2接收到Topic訊息:【" + msg + "】"); }
5.2.生產者
// 測驗topic @Test public void testTopicExchange() throws Exception { //引數一: 交換機名稱 引數二:routingKey 引數三: 訊息內容 rabbitTemplate.convertAndSend("topic.exchange", "china.news", "喜報!孫悟空大戰孫行者,勝!!!"); }
6.訊息轉換器
默認情況下SpringAMQP采用的序列化方式是JDK序列化,眾所周知,JDK序列化存在下列問題:
-
資料體積過大
-
可讀性差

6.1.配置JSON轉換器
如果我們希望訊息的體積更小、可讀性更高,因此可以使用JSON方式來做序列化和反序列化,
6.1.1.在父工程中引入依賴
<dependency> <groupId>com.fasterxml.jackson.core</groupId> <artifactId>jackson-databind</artifactId> <version>2.9.10</version> </dependency>
6.1.2.在publisher和consumer兩個服務啟動類中添加json轉換器
SpringBoot的啟動類也是1個配置類;
需要注意的是在publisher和consumer配置完了Json轉換器之后,雙方經應該傳輸JSON資料型別的訊息,否則訊息反序列化是就會報錯;
@Bean public MessageConverter jsonMessageConverter(){ return new Jackson2JsonMessageConverter(); }
6.1.3.測驗

參考
轉載請註明出處,本文鏈接:https://www.uj5u.com/ruanti/498990.html
標籤:其他
下一篇:設計模式之概述篇

