小小MQ,知識點竟然這么多???
- 一、MQ的基本概念
- 1.MQ概述
- 二、MQ的優勢
- 1.應用解耦
- 2.異步提速
- 3.削峰填谷
- 三、MQ的劣勢
- 系統可用性降低
- 系統復雜度提高
- 四、常見的MQ產品
- 五、RabbitMQ 介紹
- 1.RabbitMQ 簡介
- 1.1 Producer(生產者) 和 Consumer(消費者)
- 1.2 Exchange(交換器)
- 1.3 Queue(訊息佇列)
- 1.4 Broker(訊息中間件的服務節點)
- 1.5 Exchange Types(交換器型別)
- ① fanout
- ② direct
- ③ topic
- ④ headers(不推薦)
- 2.總結
- 3.Windows本地環境安裝RabbitMQ
- 3.1 下載Erlang
- 3.2 安裝RabbitMQ(3.9.5)
- 3.3啟動
- 3.4安裝管理插件
- 5.訪問后臺
- 4.從一個簡單的例子認識RabbitMQ
- 1.示例
- 2.重復消費
- 解決方案
- 3.訊息丟失
- 4.訊息積壓
- 5.MQ高可用
- 六、原始碼分析
- 1.回顧設計模式
- 1.1 工廠方法模式
- 1.2 抽象工廠模式
- 1.3 建造者模式
- 2.理解RabbitMQ與客戶端的資料互動
- 3.帶著問題去看代碼
- 3.1 訊息是如何入隊的?
- 3.2 訊息是如何被消費的?
- 3.3 客戶端是如何建立連接的?
- 3.4 channel是什么?
- 3.5 RabbitMQ是如何實作異步的?
- 3.6 RabbitMQ如何確保訊息不會丟失?Ack
- 3.7 為什么消費者在啟動后會自動消費佇列中的訊息?
- 3.8 消費者接收到消費資訊,停止運行消費者后,會繼續執行消費行為?
- 4.核心思想
一、MQ的基本概念
1.MQ概述
MQ全稱Message Queue(訊息佇列),是在訊息的傳輸程序中保存訊息的容器,多用于分布式系統之間進行通信,
常見的服務通信:

加入MQ后:

二、MQ的優勢
1.應用解耦
服務與服務之間不再約定協議而對接介面,而是通過生產者/消費者的模式讓中間件的MQ來對接兩邊的資料通信實作解耦合(擴展性更強),
常見的服務通信:
加入MQ后:

2.異步提速
常見的服務通信:
需要阻塞成功獲取到回應狀態在寫入資料到資料庫,阻塞同步往下執行,
加入MQ后:

3.削峰填谷
常見的服務通信:當一瞬間有5000個請求給到服務器時,服務器最大能處理1000請求,受不了了當場去世,

加入MQ后:
先把請求丟到佇列中等待處理,系統每秒從mq中拉取1000個請求進行處理(剛好卡在能處理的1000個,真實壓榨案例),這樣就變成了從1秒鐘處理5000個請求的高峰期,拆分成了 5秒鐘,每秒鐘處理1000請求的平緩處理器,哦不,是滿載處理器,

根據下面的圖所示,確實高峰期被削掉了

三、MQ的劣勢
系統可用性降低
系統引入的外部依賴越多,系統穩定性越差,一旦MQ宕機,就會對業務產生影響,如何保證MQ的高可用?
系統復雜度提高
MQ的加入大大增加了系統的復雜度,以前系統間是同步的遠程呼叫,現在是通過MQ進行的異步呼叫,如何保證訊息不背丟失等等,
四、常見的MQ產品
MQ是一種抽象的概念,衍生了各種基于其思想的實作,如常見的:RabbitMQ,RocketMQ,Kafka,

五、RabbitMQ 介紹
建議直接看總結,下面的知識可以慢慢品(以下參考Guide哥的)
1.RabbitMQ 簡介
RabbitMQ 是采用 Erlang 語言實作 AMQP(Advanced Message Queuing Protocol,高級訊息佇列協議)的訊息中間件,它最初起源于金融系統,用于在分布式系統中存盤轉發訊息,
RabbitMQ 發展到今天,被越來越多的人認可,這和它在易用性、擴展性、可靠性和高可用性等方面的卓著表現是分不開的,RabbitMQ 的具體特點可以概括為以下幾點:
- 可靠性: RabbitMQ使用一些機制來保證訊息的可靠性,如持久化、傳輸確認及發布確認等,
- 靈活的路由: 在訊息進入佇列之前,通過交換器來路由訊息,對于典型的路由功能,RabbitMQ 己經提供了一些內置的交換器來實作,針對更復雜的路由功能,可以將多個交換器系結在一起,也可以通過插件機制來實作自己的交換器,這個后面會在我們將 RabbitMQ 核心概念的時候詳細介紹到,
- 擴展性: 多個RabbitMQ節點可以組成一個集群,也可以根據實際業務情況動態地擴展集群中節點,
- 高可用性: 佇列可以在集群中的機器上設定鏡像,使得在部分節點出現問題的情況下佇列仍然可用,
- 支持多種協議: RabbitMQ 除了原生支持 AMQP 協議,還支持 STOMP、MQTT 等多種訊息中間件協議,
- 多語言客戶端: RabbitMQ幾乎支持所有常用語言,比如 Java、Python、Ruby、PHP、C#、JavaScript等,
- 易用的管理界面: RabbitMQ提供了一個易用的用戶界面,使得用戶可以監控和管理訊息、集群中的節點等,在安裝 RabbitMQ 的時候會介紹到,安裝好 RabbitMQ 就自帶管理界面,
- 插件機制: RabbitMQ 提供了許多插件,以實作從多方面進行擴展,當然也可以撰寫自己的插件,感覺這個有點類似 Dubbo 的 SPI機制,
RabbitMQ 整體上是一個生產者與消費者模型,主要負責接收、存盤和轉發訊息,可以把訊息傳遞的程序想象成:當你將一個包裹送到郵局,郵局會暫存并最終將郵件通過郵遞員送到收件人的手上,RabbitMQ就好比由郵局、郵箱和郵遞員組成的一個系統,從計算機術語層面來說,RabbitMQ 模型更像是一種交換機模型,
下面再來看看圖—— RabbitMQ 的整體模型架構,

1.1 Producer(生產者) 和 Consumer(消費者)
- Producer(生產者) :生產訊息的一方(郵件投遞者)
- Consumer(消費者) :消費訊息的一方(郵件收件人)
訊息一般由 2 部分組成:訊息頭(或者說是標簽 Label)和 訊息體,訊息體也可以稱為 payLoad ,訊息體是不透明的,而訊息頭則由一系列的可選屬性組成,這些屬性包括 routing-key(路由鍵)、priority(相對于其他訊息的優先權)、delivery-mode(指出該訊息可能需要持久性存盤)等,生產者把訊息交由 RabbitMQ 后,RabbitMQ 會根據訊息頭把訊息發送給感興趣的 Consumer(消費者),
1.2 Exchange(交換器)
在 RabbitMQ 中,訊息并不是直接被投遞到 Queue(訊息佇列) 中的,中間還必須經過 Exchange(交換器) 這一層,Exchange(交換器) 會把我們的訊息分配到對應的 Queue(訊息佇列) 中,
Exchange(交換器) 用來接收生產者發送的訊息并將這些訊息路由給服務器中的佇列中,如果路由不到,或許會回傳給 Producer(生產者) ,或許會被直接丟棄掉 ,這里可以將RabbitMQ中的交換器看作一個簡單的物體,
RabbitMQ 的 Exchange(交換器) 有4種型別,不同的型別對應著不同的路由策略:direct(默認),fanout, topic, 和 headers,不同型別的Exchange轉發訊息的策略有所區別,這個會在介紹 Exchange Types(交換器型別) 的時候介紹到,
Exchange(交換器) 示意圖如下:

生產者將訊息發給交換器的時候,一般會指定一個 RoutingKey(路由鍵),用來指定這個訊息的路由規則,而這個 RoutingKey 需要與交換器型別和系結鍵(BindingKey)聯合使用才能最終生效,
RabbitMQ 中通過 Binding(系結) 將 Exchange(交換器) 與 Queue(訊息佇列) 關聯起來,在系結的時候一般會指定一個 BindingKey(系結建) ,這樣 RabbitMQ 就知道如何正確將訊息路由到佇列了,如下圖所示,一個系結就是基于路由鍵將交換器和訊息佇列連接起來的路由規則,所以可以將交換器理解成一個由系結構成的路由表,Exchange 和 Queue 的系結可以是多對多的關系,
Binding(系結) 示意圖:

生產者將訊息發送給交換器時,需要一個RoutingKey,當 BindingKey 和 RoutingKey 相匹配時,訊息會被路由到對應的佇列中,在系結多個佇列到同一個交換器的時候,這些系結允許使用相同的 BindingKey,BindingKey 并不是在所有的情況下都生效,它依賴于交換器型別,比如fanout型別的交換器就會無視,而是將訊息路由到所有系結到該交換器的佇列中,
1.3 Queue(訊息佇列)
Queue(訊息佇列) 用來保存訊息直到發送給消費者,它是訊息的容器,也是訊息的終點,一個訊息可投入一個或多個佇列,訊息一直在佇列里面,等待消費者連接到這個佇列將其取走,
RabbitMQ 中訊息只能存盤在 佇列 中,這一點和 Kafka 這種訊息中間件相反,Kafka 將訊息存盤在 topic(主題) 這個邏輯層面,而相對應的佇列邏輯只是topic實際存盤檔案中的位移標識, RabbitMQ 的生產者生產訊息并最終投遞到佇列中,消費者可以從佇列中獲取訊息并消費,
多個消費者可以訂閱同一個佇列,這時佇列中的訊息會被平均分攤(Round-Robin,即輪詢)給多個消費者進行處理,而不是每個消費者都收到所有的訊息并處理,這樣避免的訊息被重復消費,
RabbitMQ 不支持佇列層面的廣播消費,如果有廣播消費的需求,需要在其上進行二次開發,這樣會很麻煩,不建議這樣做,
1.4 Broker(訊息中間件的服務節點)
對于 RabbitMQ 來說,一個 RabbitMQ Broker 可以簡單地看作一個 RabbitMQ 服務節點,或者RabbitMQ服務實體,大多數情況下也可以將一個 RabbitMQ Broker 看作一臺 RabbitMQ 服務器,
下圖展示了生產者將訊息存入 RabbitMQ Broker,以及消費者從Broker中消費資料的整個流程,

這樣圖1中的一些關于 RabbitMQ 的基本概念我們就介紹完畢了,下面再來介紹一下 Exchange Types(交換器型別) ,
1.5 Exchange Types(交換器型別)
RabbitMQ 常用的 Exchange Type 有 fanout、direct、topic、headers 這四種(AMQP規范里還提到兩種 Exchange Type,分別為 system 與 自定義,這里不予以描述),
① fanout
fanout 型別的Exchange路由規則非常簡單,它會把所有發送到該Exchange的訊息路由到所有與它系結的Queue中,不需要做任何判斷操作,所以 fanout 型別是所有的交換機型別里面速度最快的,fanout 型別常用來廣播訊息,
② direct
direct 型別的Exchange路由規則也很簡單,它會把訊息路由到那些 Bindingkey 與 RoutingKey 完全匹配的 Queue 中,

以上圖為例,如果發送訊息的時候設定路由鍵為“warning”,那么訊息會路由到 Queue1 和 Queue2,如果在發送訊息的時候設定路由鍵為"Info”或者"debug”,訊息只會路由到Queue2,如果以其他的路由鍵發送訊息,則訊息不會路由到這兩個佇列中,
direct 型別常用在處理有優先級的任務,根據任務的優先級把訊息發送到對應的佇列,這樣可以指派更多的資源去處理高優先級的佇列,
③ topic
前面講到direct型別的交換器路由規則是完全匹配 BindingKey 和 RoutingKey ,但是這種嚴格的匹配方式在很多情況下不能滿足實際業務的需求,topic型別的交換器在匹配規則上進行了擴展,它與 direct 型別的交換器相似,也是將訊息路由到 BindingKey 和 RoutingKey 相匹配的佇列中,但這里的匹配規則有些不同,它約定:
- RoutingKey 為一個點號“.”分隔的字串(被點號“.”分隔開的每一段獨立的字串稱為一個單詞),如 “com.rabbitmq.client”、“java.util.concurrent”、“com.hidden.client”;
- BindingKey 和 RoutingKey 一樣也是點號“.”分隔的字串;
- BindingKey 中可以存在兩種特殊字串“*”和“#”,用于做模糊匹配,其中“*”用于匹配一個單詞,“#”用于匹配多個單詞(可以是零個),

以上圖為例: - 路由鍵為 “com.rabbitmq.client” 的訊息會同時路由到 Queuel 和 Queue2;
- 路由鍵為 “com.hidden.client” 的訊息只會路由到 Queue2 中;
- 路由鍵為 “com.hidden.demo” 的訊息只會路由到 Queue2 中;
- 路由鍵為 “java.rabbitmq.demo” 的訊息只會路由到Queuel中;
- 路由鍵為 “java.util.concurrent” 的訊息將會被丟棄或者回傳給生產者(需要設定 mandatory 引數),因為它沒有匹配任何路由鍵,
④ headers(不推薦)
headers 型別的交換器不依賴于路由鍵的匹配規則來路由訊息,而是根據發送的訊息內容中的 headers 屬性進行匹配,在系結佇列和交換器時制定一組鍵值對,當發送訊息到交換器時,RabbitMQ會獲取到該訊息的 headers(也是一個鍵值對的形式)'對比其中的鍵值對是否完全匹配佇列和交換器系結時指定的鍵值對,如果完全匹配則訊息會路由到該佇列,否則不會路由到該佇列,headers 型別的交換器性能會很差,而且也不實用,基本上不會看到它的存在,
2.總結
用大白話說,RabbitMQ整個架構就是:客戶端、交換機、佇列,這三個角色因為交換機不同的模式(直連交換機、扇形交換機、主體交換機、首部交換機)以及不同的組裝形成了RabbitMQ使用的各種模式:簡單模式()、作業佇列模式、發布/訂閱模式、路由模式、通配符模式等,其最核心的莫過于佇列,提供者及對應入隊,消費后對應出隊,
3.Windows本地環境安裝RabbitMQ
3.1 下載Erlang
RabbitMQ是基于Erlang環境開發的,先下載個Erlang(24),https://erlang.org/download/otp_win64_24.0.exe,下載直接一鍵點安裝
3.2 安裝RabbitMQ(3.9.5)
https://github.com/rabbitmq/rabbitmq-server/releases/download/v3.9.5/rabbitmq-server-3.9.5.exe,·1下載直接一鍵點安裝
3.3啟動
這個應該不用我說按啥啟動了吧

啟動成功,如下

3.4安裝管理插件
右鍵快捷方式打開檔案夾所在目錄

執行:
rabbitmq-plugins enable rabbitmq_management
5.訪問后臺
后臺管理網址:
http://localhost:15672/
安裝后的默認初始
賬號:guest
密碼:guset
然后就得到了黑化版的RabbitMQ?

4.從一個簡單的例子認識RabbitMQ
1.示例
依賴:
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
<version>2.4.1</version>
</dependency>
生產者:
package boot.spring.test;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import java.io.IOException;
import java.util.concurrent.TimeoutException;
/**
* @description:
* @author:lx
* @date: 2021/09/04 下午 3:12
* @Copyright: lx
*/
public class Provider {
/**
* 宣告的佇列名
*/
private final static String QUEUE_NAME = "test_queue";
public static void main(String[] args) throws IOException, TimeoutException {
ConnectionFactory connectionFactory = new ConnectionFactory();
connectionFactory.setHost("127.0.0.1");
// 默認埠號
connectionFactory.setPort(5672);
connectionFactory.setUsername("guest");
connectionFactory.setPassword("guest");
connectionFactory.setVirtualHost("/");
// 獲取TCP長連接
Connection conn = connectionFactory.newConnection();
// 創建通信“通道”,相當于TCP中的虛擬連接
Channel channel = conn.createChannel();
// 開啟RabbitMQ事務,當沒有接收到MQ反饋時拋出例外并回滾
channel.txSelect();
// 創建佇列,宣告并創建一個佇列,如果佇列已存在,則使用這個佇列
// 第一個引數:佇列名稱
// 第二個引數:是否持久化,false對應不持久化資料,MQ停掉資料就會丟失
// 第三個引數:是否佇列私有化,false則代表所有消費者都可以訪問,true代表只有一次則擁有它的消費者才能一直使用,其它消費者不讓訪問
// 第四個引數:是否自動洗掉,false代表連接停掉后不自動洗掉這個佇列
// 其它額外引數:null
channel.queueDeclare(QUEUE_NAME, true, false, false, null);
String message = "hello world!";
try {
// 第一個引數:交換機,這里時簡單demo版本,沒有用到交換機
// 第二個引數:佇列名稱
// 第三個引數:額外的設定屬性
// 第四個引數:要傳遞的訊息位元組陣列
channel.basicPublish("", QUEUE_NAME, null, message.getBytes());
} catch (Exception e) {
// 發生例外的回滾
channel.txRollback();
}
// 正常流程提交事務
channel.txCommit();
channel.close();
conn.close();
System.out.println("發送成功!");
}
}
消費者:
package boot.spring.test;
import com.rabbitmq.client.*;
import java.io.IOException;
import java.util.concurrent.TimeoutException;
/**
* @description:
* @author:lx
* @date: 2021/09/04 下午 3:33
* @Copyright: lx
*/
public class Consumer {
private final static String QUEUE_NAME = "test_queue";
public static void main(String[] args) throws IOException, TimeoutException {
ConnectionFactory connectionFactory = new ConnectionFactory();
connectionFactory.setHost("127.0.0.1");
connectionFactory.setPort(5672);
connectionFactory.setUsername("guest");
connectionFactory.setPassword("guest");
connectionFactory.setVirtualHost("/");
// 獲取TCP長連接
Connection conn = connectionFactory.newConnection();
// 創建通信“通道”,相當于TCP中的虛擬連接
Channel channel = conn.createChannel();
// 創建佇列,宣告并創建一個佇列,如果佇列已存在,則使用這個佇列
// 第一個引數:佇列名稱
// 第二個引數:是否持久化,false對應不持久化資料,MQ停掉資料就會丟失
// 第三個引數:是否佇列私有化,false則代表所有消費者都可以訪問,true代表只有一次則擁有它的消費者才能一直使用,其它消費者不讓訪問
// 第四個引數:是否自動洗掉,false代表連接停掉后不自動洗掉這個佇列
// 其它額外引數:null
channel.queueDeclare(QUEUE_NAME, true, false, false, null);
// 創建一個訊息消費者
// 第一個引數:佇列名稱
// 第二個引數:第二個引數表示是否自動確認收到訊息,false代表手動編程確認訊息
// 第三個引數:傳入DefaultConsumer的實作類
channel.basicConsume(QUEUE_NAME, false, new Receiver(channel));
}
}
/**
* @description:
* @author:lx
* @date: 2021/09/04 下午 3:37
* @Copyright: lx
*/
class Receiver extends DefaultConsumer {
private Channel channel;
/**
* 重寫建構式,channel通道物件需要從外層傳入,在handleDelivery中要用到
*
* @param channel
*/
public Receiver(Channel channel) {
super(channel);
this.channel = channel;
}
@Override
public void handleDelivery(String consumerTag,
Envelope envelope,
AMQP.BasicProperties properties,
byte[] body)
throws IOException {
String message = new String(body);
System.out.println("消費者接收到的訊息:" + message);
System.out.println("訊息的TagID:" + envelope.getDeliveryTag());
// int i = 1 / 0;
// false只簽收當前的訊息,設定為true的時候代表簽收該消費者所有未簽收的訊息
channel.basicAck(envelope.getDeliveryTag(), false);
}
}
2.重復消費
啟動生產者,來到管理頁面,可以看到訊息已經入隊,進入準備被消費的狀態,

Debug啟動消費者后,在管理頁面可以看到,隊內訊息狀態由準備-轉變到了待回應狀態,total總數還是存在的,此時訊息已被接收到,但未被回應,

由于斷點導致MQ超時未收到回應,狀態回滾到ready,訊息仍然在隊中,但事實上消費者已經消費過一次了,這里引出一個問題-重復消費,

還是很多場景導致你發生例外回滾的情況還有很多,比如:
程式發送例外導致,順帶一提,程式例外會導致消費者程式崩潰,MQ也會一直阻塞在等待回應的階段
消費者行程突然GG

使用代理將捕獲轉發,模擬丟包的情況
解決方案
既然重復消費這種情況是難以避免的,那么我們如何去處理這種情況呢?
消費端處理訊息的業務邏輯保持冪等性,
冪等性,通俗點說,就一個資料,或者一個請求,給你重復來多次,你得確保對應的資料是不會改變的,不能出錯,
1.你拿到這個訊息做資料庫的insert操作,那就容易了,給這個訊息做一個唯一主鍵,那么就算出現重復消費的情況,就會導致主鍵沖突,避免資料庫出現臟資料,
2.你拿到這個訊息做redis的set的操作,那就容易了,不用解決,因為你無論set幾次結果都是一樣的,set操作本來就算冪等操作,
3.準備一個第三方介質,來做消費記錄,以redis為例,給訊息分配一個全域id,只要消費過該訊息,將<id,message>以K-V形式寫入redis,那消費者開始消費前,先去redis中查詢有沒消費記錄即可,
3.訊息丟失
訊息在網路傳輸中丟失,MQ宕機丟失訊息
4.訊息積壓
程式例外會導致消費者程式崩潰,MQ也會一直阻塞在等待回應的階段,導致訊息一直堆積
5.MQ高可用
https://blog.csdn.net/yygEwing/article/details/116329666?utm_source=app&app_version=4.14.1
六、原始碼分析
從簡單的原始碼分析更深入的認知RabbitMQ

Broker: 接收和分發訊息的應用,RabbitMQ Server就是Message Broker
Virtual Host: 處于多租戶和安全因素設計的,把AMQP的基本組件劃分到一個虛擬的分組中,類似于網路中的namespace概念,當多個不同的用戶使用同一個RabbitMQserver提供的服務時,可以劃分出多個vhost,每個用戶在自己的vhost創建exchange/queue等
Connection: publisher/consumer和broker之間的TCP連接
Channel: 如何每一次訪問RabbitMQ都建立一個Connection,在訊息量大的時候建立TCP Connection的開銷將是巨大的,效率也低,channel是在connection內部建立的邏輯連接,如果應用程式支持多執行緒,通常每個thread創建單獨的channel進行通訊,AMQPmethod包含了channelId幫助客戶端和message broker識別channel,所以channel之間是完全隔離的,Channel作為輕量級的Connection極大減少了作業系統建立TCP connection的開銷,
1.回顧設計模式
1.1 工廠方法模式
將多段代碼的共性行為抽象到介面中去定義,具體的實作由子類實作父類后去定義,最后,通過一個工廠類去根據傳參來選擇回傳對應的實體化物件,
關鍵詞:工廠類一般帶有Factory
1.2 抽象工廠模式
抽象工廠的本質是其它工廠類的抽象類,也就是將其他工廠類中的共性行為提取到了抽象工廠類
AbstractXXX
1.3 建造者模式
日常生活中,裝修房子會根據不同的場景、品牌、型號、價格等等組合形成了各式各樣的裝修風格(套餐A:現代簡約,套餐B:輕奢田園,套餐C:歐式豪華)
一些基本物料不會變,而其組合經常變化的時候,就可以選擇這樣的構建者模式來構建代碼,
Builder
2.理解RabbitMQ與客戶端的資料互動
3.帶著問題去看代碼
3.1 訊息是如何入隊的?
3.2 訊息是如何被消費的?
3.3 客戶端是如何建立連接的?
3.4 channel是什么?
3.5 RabbitMQ是如何實作異步的?
3.6 RabbitMQ如何確保訊息不會丟失?Ack
3.7 為什么消費者在啟動后會自動消費佇列中的訊息?
3.8 消費者接收到消費資訊,停止運行消費者后,會繼續執行消費行為?
4.核心思想
RabbitMQ實作MQ的核心思想是?其特點是?其中用到了什么設計模式?最有印象的是?

轉載請註明出處,本文鏈接:https://www.uj5u.com/qita/298080.html
標籤:其他
上一篇:【elasticsearch系列】windows安裝IK分詞器插件
下一篇:2020年技能樹
