
一、什么是訊息中間件
兩個系統或兩個客戶端之間進行訊息傳送,利用高效可靠的訊息傳遞機制進行平臺無關的資料交流,并基于資料通信來進行分布式系統的集成,通過提供訊息傳遞和訊息排隊模型,它可以在分布式環境下擴展行程間的通信,
訊息中間件,總結起來作用有三個:異步化提升性能、降低耦合度、流量削峰,

系統A發送訊息給中間件后,自己的作業已經完成了,不用再去管系統B什么時候完成操作,而系統B拉去訊息后,執行自己的操作也不用告訴系統A執行結果,所以整個的通信程序是異步呼叫的,
二、訊息中間件的應用場景
2.1 異步通信
有些業務不想也不需要立即處理訊息,訊息佇列提供了異步處理機制,允許用戶把一個訊息放入佇列,但并不立即處理它,想向佇列中放入多少訊息就放多少,然后在需要的時候再去處理它們,

2.2 緩沖
在任何重要的系統中,都會有需要不同的處理時間的元素,訊息佇列通過一個緩沖層來幫助任務最高效率的執行,該緩沖有助于控制和優化資料流經過系統的速度,以調節系統回應時間,
2.3 解耦
降低工程間的強依賴程度,針對異構系統進行適配,在專案啟動之初來預測將來專案會碰到什么需求,是極其困難的,通過訊息系統在處理程序中間插入了一個隱含的、基于資料的介面層,兩邊的處理程序都要實作這一介面,當應用發生變化時,可以獨立的擴展或修改兩邊的處理程序,只要確保它們遵守同樣的介面約束,
2.4 冗余
有些情況下,處理資料的程序會失敗,除非資料被持久化,否則將造成丟失,訊息佇列把資料進行持久化直到它們已經被完全處理,通過這一方式規避了資料丟失風險,許多訊息佇列所采用的”插入-獲取-洗掉”范式中,在把一個訊息從佇列中洗掉之前,需要你的處理系統明確的指出該訊息已經被處理完畢,從而確保你的資料被安全的保存直到你使用完畢,
2.5 擴展性
因為訊息佇列解耦了你的處理程序,所以增大訊息入隊和處理的頻率是很容易的,只要另外增加處理程序即可,不需要改變代碼、不需要調節引數,便于分布式擴容,
2.6 可恢復性
系統的一部分組件失效時,不會影響到整個系統,訊息佇列降低了行程間的耦合度,所以即使一個處理訊息的行程掛掉,加入佇列中的訊息仍然可以在系統恢復后被處理,
2.7 順序保證
在大多使用場景下,資料處理的順序都很重要,大部分訊息佇列本來就是排序的,并且能保證資料會按照特定的順序來處理,
2.8 過載保護
在訪問量劇增的情況下,應用仍然需要繼續發揮作用,但是這樣的突發流量無法提取預知;如果以為了能處理這類瞬間峰值訪問為標準來投入資源隨時待命無疑是巨大的浪費,使用訊息佇列能夠使關鍵組件頂住突發的訪問壓力,而不會因為突發的超負荷的請求而完全崩潰,
2.9 資料流處理
分布式系統產生的海量資料流,如:業務日志、監控資料、用戶行為等,針對這些資料流進行實時或批量采集匯總,然后進行大資料分析是當前互聯網的必備技術,通過訊息佇列完成此類資料收集是最好的選擇,
三、常用訊息佇列(ActiveMQ、RabbitMQ、RocketMQ、Kafka)比較
| 特性MQ | ActiveMQ | RabbitMQ | RocketMQ | Kafka |
|---|---|---|---|---|
| 生產者消費者模式 | 支持 | 支持 | 支持 | 支持 |
| 發布訂閱模式 | 支持 | 支持 | 支持 | 支持 |
| 請求回應模式 | 支持 | 支持 | 不支持 | 不支持 |
| Api完備性 | 高 | 高 | 高 | 高 |
| 多語言支持 | 支持 | 支持 | java | 支持 |
| 單機吞吐量 | 萬級 | 萬級 | 萬級 | 十萬級 |
| 訊息延遲 | 無 | 微秒級 | 毫秒級 | 毫秒級 |
| 可用性 | 高(主從) | 高(主從) | 非常高(分布式) | 非常高(分布式) |
| 訊息丟失 | 低 | 低 | 理論上不會丟失 | 理論上不會丟失 |
| 檔案的完備性 | 高 | 高 | 高 | 高 |
| 提供快速入門 | 有 | 有 | 有 | 有 |
| 社區活躍度 | 高 | 高 | 有 | 高 |
| 商業支持 | 無 | 無 | 商業云 | 商業云 |
四、訊息中間件的角色
Queue: 佇列存盤,常用與點對點訊息模型 ,默認只能由唯一的一個消費者處理,一旦處理訊息洗掉,
Topic: 主題存盤,用于訂閱/發布訊息模型,主題中的訊息,會發送給所有的消費者同時處理,只有在訊息可以重復處 理的業務場景中可使用,Queue/Topic都是 Destination 的子介面
ConnectionFactory: 連接工廠,客戶用來創建連接的物件,例如ActiveMQ提供的ActiveMQConnectionFactory
Connection: JMS Connection封裝了客戶與JMS提供者之間的一個虛擬的連接,
Destination: 訊息的目的地,目的地是客戶用來指定它生產的訊息的目標和它消費的訊息的來源的物件,JMS1.0.2規范中定義了兩種訊息傳遞域:點對點(PTP)訊息傳遞域和發布/訂閱訊息傳遞域,
點對點訊息傳遞域的特點如下:
- 每個訊息只能有一個消費者,
- 訊息的生產者和消費者之間沒有時間上的相關性,無論消費者在生產者發送訊息的時候是否處于運行狀態,它都可以提取訊息,
發布/訂閱訊息傳遞域的特點如下:
- 每個訊息可以有多個消費者,
- 生產者和消費者之間有時間上的相關性,
- 訂閱一個主題的消費者只能消費自它訂閱之后發布的訊息,JMS規范允許客戶創建持久訂閱,這在一定程度上放松了時間上的相關性要求 ,持久訂閱允許消費者消費它在未處于激活狀態時發送的訊息,
在點對點訊息傳遞域中,目的地被成為佇列(queue);在發布/訂閱訊息傳遞域中,目的地被成為主題(topic),
五、JMS的訊息格式
JMS訊息由以下三部分組成的:
-
訊息頭:
每個訊息頭欄位都有相應的getter和setter方法,
-
訊息屬性:
如果需要除訊息頭欄位以外的值,那么可以使用訊息屬性,
-
訊息體:
JMS定義的訊息型別有TextMessage、MapMessage、BytesMessage、StreamMessage和ObjectMessage,
訊息型別:
| 屬性 | 型別 |
|---|---|
| TextMessage | 文本訊息 |
| MapMessage | k/v |
| BytesMessage | 位元組流 |
| StreamMessage | java原始的資料流 |
| ObjectMessage | 序列化的java物件 |
六、訊息可靠性機制
只有在被確認之后,才認為已經被成功地消費了,訊息的成功消費通常包含三個階段 :客戶接收訊息、客戶處理訊息和訊息被確認,在事務性會話中,當一個事務被提交的時候,確認自動發生,在非事務性會話中,訊息何時被確認取決于創建會話時的應答模式(acknowledgement mode),該引數有以下三個可選值:
Session.AUTO_ACKNOWLEDGE:當客戶成功的從receive方法回傳的時候,或者從MessageListener.onMessage方法成功回傳的時候,會話自動確認客戶收到的訊息,Session.CLIENT_ACKNOWLEDGE:客戶通過訊息的acknowledge方法確認訊息,需要注意的是,在這種模式中,確認是在會話層上進行:確認一個被消費的訊息將自動確認所有已被會話消費的訊息,例如,如果一個訊息消費者消費了10個訊息,然后確認第5個訊息,那么所有10個訊息都被確認,Session.DUPS_ACKNOWLEDGE:該選擇只是會話遲鈍的確認訊息的提交,如果JMS Provider失敗,那么可能會導致一些重復的訊息,如果是重復的訊息,那么JMS Provider必須把訊息頭的JMSRedelivered欄位設定為true,
6.1 優先級
可以使用訊息優先級來指示JMS Provider首先提交緊急的訊息,優先級分10個級別,從0(最低)到9(最高),如果不指定優先級,默認級別是4,需要注意的是,JMS Provider并不一定保證按照優先級的順序提交訊息,
6.2 訊息過期
可以設定訊息在一定時間后過期,默認是永不過期,
6.3 臨時目的地
可以通過會話上的createTemporaryQueue方法和createTemporaryTopic方法來創建臨時目的地,它們的存在時間只限于創建它們的連接所保持的時間,只有創建該臨時目的地的連接上的訊息消費者才能夠從臨時目的地中提取訊息,
七、什么是ActiveMQ
ActiveMQ是一種開源的基于JMS(Java Message Servie)規范的一種訊息中間件的實作,ActiveMQ的設計目標是提供標準的,面向訊息的,能夠跨越多語言和多系統的應用集成訊息通信中間件,
官網地址:http://activemq.apache.org/
7.1 存盤方式
1. KahaDB存盤: KahaDB是默認的持久化策略,所有訊息順序添加到一個日志檔案中,同時另外有一個索引檔案記錄指向這些日志的存盤地址,還有一個事務日志用于訊息回復操作,是一個專門針對訊息持久化的解決方案,它對典型的訊息使用模式進行了優化
特性:
1、日志形式存盤訊息;
2、訊息索引以 B-Tree 結構存盤,可以快速更新;
3、 完全支持 JMS 事務;
4、支持多種恢復機制kahadb 可以限制每個資料檔案的大小,不代表總計資料容量,
2. AMQ 方式: 只適用于 5.3 版本之前, AMQ 也是一個檔案型資料庫,訊息資訊最終是存盤在檔案中,記憶體中也會有快取資料,
3. JDBC存盤 : 使用JDBC持久化方式,資料庫默認會創建3個表,每個表的作用如下:
activemq_msgs:queue和topic的訊息都存在這個表中
activemq_acks:存盤持久訂閱的資訊和最后一個持久訂閱接收的訊息ID
activemq_lock:跟kahadb的lock檔案類似,確保資料庫在某一時刻只有一個broker在訪問
4. LevelDB存盤 : LevelDB持久化性能高于KahaDB,但是在ActiveMQ官網對LevelDB的表述:LevelDB官方建議使用以及不再支持,推薦使用的是KahaDB
5.Memory 訊息存盤: 顧名思義,基于記憶體的訊息存盤,就是訊息存盤在記憶體中,persistent=”false”,表示不設定持 久化存盤,直接存盤到記憶體中,在broker標簽處設定,
7.2 協議
協議官網API:http://activemq.apache.org/configuring-version-5-transports.html
-
Transmission Control Protocol (TCP):
- 這是默認的Broker配置,TCP的Client監聽埠是61616,
- 在網路傳輸資料前,必須要序列化資料,訊息是通過一個叫wire protocol的來序列化成位元組流,默認情況下,ActiveMQ把wire protocol叫做OpenWire,它的目的是促使網路上的效率和資料快速互動,
- TCP連接的URI形式:tcp://hostname:port?key=value&key=value
- TCP傳輸的優點:
(1)TCP協議傳輸可靠性高,穩定性強
(2)高效性:位元組流方式傳遞,效率很高
(3)有效性、可用性:應用廣泛,支持任何平臺
-
New I/O API Protocol(NIO)
-
NIO協議和TCP協議類似,但NIO更側重于底層的訪問操作,它允許開發人員對同一資源可有更多的client呼叫和服務端有更多的負載,
-
適合使用NIO協議的場景:
(1)可能有大量的Client去鏈接到Broker上一般情況下,大量的Client去鏈接Broker是被作業系統的執行緒數所限制的,因此,NIO的實作比TCP需要更少的執行緒去運行,所以建議使用NIO協議
(2)可能對于Broker有一個很遲鈍的網路傳輸NIO比TCP提供更好的性能 -
NIO連接的URI形式:nio://hostname:port?key=value
-
Transport Connector配置示例:
-
<transportConnectors>
<transportConnector
name="tcp"
uri="tcp://localhost:61616?trace=true" />
<transportConnector
name="nio"
uri="nio://localhost:61618?trace=true" />
</transportConnectors>
- User Datagram Protocol(UDP)
1:UDP和TCP的區別
(1)TCP是一個原始流的傳遞協議,意味著資料包是有保證的,換句話說,資料包是不會被復制和丟失的,UDP,另一方面,它是不會保證資料包的傳遞的
(2)TCP也是一個穩定可靠的資料包傳遞協議,意味著資料在傳遞的程序中不會被丟失,這樣確保了在發送和接收之間能夠可靠的傳遞,相反,UDP僅僅是一個鏈接協議,所以它沒有可靠性之說
2:從上面可以得出:TCP是被用在穩定可靠的場景中使用的;UDP通常用在快速資料傳遞和不怕資料丟失的場景中,還有ActiveMQ通過防火墻時,只能用UDP
3:UDP連接的URI形式:udp://hostname:port?key=value
4:Transport Connector配置示例:
<transportConnectors>
<transportConnector
name="udp"
uri="udp://localhost:61618?trace=true" />
</transportConnectors>
- Secure Sockets Layer Protocol (SSL)
1:連接的URI形式:ssl://hostname:port?key=value
2:Transport Connector配置示例:
<transportConnectors>
<transportConnector name="ssl" uri="ssl://localhost:61617?trace=true"/>
</transportConnectors>
八、案例(Hello World)
這里以windows為案例演示
下載地址:http://activemq.apache.org/components/classic/download/
8.1 安裝啟動
解壓后直接執行
bin/win64/activemq.bat

8.2 web控制臺
http://localhost:8161/
賬號密碼:admin/admin


8.3 web控制臺
修改 ActiveMQ 組態檔 activemq/conf/jetty.xml
jettyport節點: 組態檔修改完畢,保存并重新啟動 ActiveMQ 服務
<bean id="jettyPort" class="org.apache.activemq.web.WebConsolePort" init-method="start">
<!-- the default port number for the web console -->
<property name="host" value="127.0.0.1"/>
<property name="port" value="8161"/>
</bean>
8.4 開發
1. jar引入:
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-activemq</artifactId>
</dependency>
2. Sender :
import org.apache.activemq.ActiveMQConnectionFactory;
import javax.jms.*;
/**
* @program: activemq_01
* @ClassName Sender
* @description: 訊息發送
* @author: muxiaonong
* @create: 2020-10-02 13:01
* @Version 1.0
**/
public class Sender {
public static void main(String[] args) throws Exception{
// 1. 獲取連接工廠
ActiveMQConnectionFactory factory = new ActiveMQConnectionFactory(
ActiveMQConnectionFactory.DEFAULT_USER,
ActiveMQConnectionFactory.DEFAULT_PASSWORD,
"tcp://localhost:61616"
);
// 2. 獲取一個向activeMq的連接
Connection connection = factory.createConnection();
// 3. 獲取session
Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
// 4.找目的地,獲取destination,消費端,也會從這個目的地取訊息
Queue queue = session.createQueue("user");
// 5.1 訊息創建者
MessageProducer producer = session.createProducer(queue);
// consumer --> 消費者
// producer --> 創建者
// 5.2. 創建訊息
for (int i = 0; i < 100; i++) {
TextMessage textMessage = session.createTextMessage("hi:"+i);
// 5.3 向目的地寫入訊息
producer.send(textMessage);
Thread.sleep(1000);
}
// 6.關閉連接
connection.close();
System.out.println("結束,,,,,");
}
}
3. Receiver :
import org.apache.activemq.ActiveMQConnectionFactory;
import javax.jms.*;
/**
* @program: activemq_01
* @ClassName Receiver
* @description: 訊息接收
* @author: muxiaonong
* @create: 2020-10-02 13:01
* @Version 1.0
**/
public class Receiver {
public static void main(String[] args) throws Exception{
// 1. 獲取連接工廠
ActiveMQConnectionFactory factory = new ActiveMQConnectionFactory(
ActiveMQConnectionFactory.DEFAULT_USER,
ActiveMQConnectionFactory.DEFAULT_PASSWORD,
"tcp://localhost:61616"
);
// 2. 獲取一個向activeMq的連接
Connection connection = factory.createConnection();
connection.start();
// 3. 獲取session
Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
// 4.找目的地,獲取destination,消費端,也會從這個目的地取訊息
Destination queue = session.createQueue("user");
// 5 獲取訊息
MessageConsumer consumer = session.createConsumer(queue);
while(true){
TextMessage message = (TextMessage)consumer.receive();
System.out.println("message:"+message.getText());
}
}
}
測驗結果:
message:hi:38
message:hi:39
message:hi:40
message:hi:41
message:hi:42
message:hi:43
message:hi:44
message:hi:45
web后臺顯示有一個消費者處于連接狀態,且已消費了68個message,而該條佇列已沒有message待消費了

九、總結
今天的MQ入門教程系列就這里了,感興趣的小伙伴可以試試,遇到了什么問題,或者有疑問的,都可以在下方留言,小農看見了會第一時間回復大家,MQ作為一個訊息中間件,不管是面試還是作業中都會經常用到,所以是很有必要去了解和學習的一個技術點,今天的分享就到這里了,謝謝各位小伙伴的觀看,我們下篇文章見,大家加油!
CSDN認證博客專家
Java后端
大資料
資料庫
轉載請註明出處,本文鏈接:https://www.uj5u.com/qukuanlian/169552.html
標籤:區塊鏈
下一篇:這次生日我想寫點什么

