RabbitMQ隊列的選擇及應(yīng)用場景
前言
本篇是Rabbit MQ高級特性的學(xué)習(xí)筆記,記錄RabbitMQ的Classic,Quorum,Stream隊列,懶隊列的特性和運用場景
一、Rabbit MQ隊列的選擇
Rabbit MQ默認提供了三種隊列:
- Classic:經(jīng)典隊列,在單機模式下是最常用的。
- Quorum:仲裁隊列,通常用于集群環(huán)境下,保證集群的高可用性。
- Stream:流式隊列,是自3.9.0版本開始引入的新特性,這種隊列類型的消息是持久化到磁盤并且具備分布式備份的。
1.1、Classic
經(jīng)典隊列是Rabbit MQ默認的隊列,與Spring Boot整合時,使用new Queue,創(chuàng)建的隊列就是普通隊列:

經(jīng)典隊列在創(chuàng)建時除了指定隊列的名稱,還有額外的四個選項:

對應(yīng)的是頁面上的:

durable代表了隊列是否進行持久化,如果開啟持久化,則會將隊列保存到磁盤上,將來Rabbit MQ重啟后,可以從磁盤恢復(fù),否則隊列只存在于內(nèi)存中,服務(wù)重啟后會自動刪除。這里的durable只是設(shè)置了隊列的持久化。消息和交換機,同樣可以設(shè)置持久化。如果消息設(shè)置了持久化,而隊列未設(shè)置持久化,重啟之后,消息依舊會丟失。所以如果要保證消息的零丟失,消息和隊列都需要設(shè)置持久化。
// 隊列聲明
channel.queueDeclare("safe_queue", true, false, false, null); // Durability=true
// 消息發(fā)布
AMQP.BasicProperties props = new AMQP.BasicProperties.Builder()
.deliveryMode(2) // 持久化消息
.build(); exclusive代表了隊列是否排他。當(dāng)屬性為true時,表示僅允許聲明了該隊列的連接進行訪問,其他連接無法訪問該隊列。通常用于單消費者模式,確保隊列只能被特定消費者訪問。
autoDelete代表了隊列是否自動刪除,如果設(shè)置為true,表示該隊列在最后一個消費者斷開連接之后,進行刪除操作。通常用于臨時隊列的場景。
arguments可以指定更多的參數(shù),比如指定死信交換機和路由鍵,消息過期時間,最大長度,是否為懶隊列等。
1.2、Quorum
Quorum是針對鏡像隊列的一種優(yōu)化,目前已經(jīng)取代了鏡像隊列,作為Rabbit MQ集群部署保證高可用性的解決方案。傳統(tǒng)的鏡像隊列,是將消息副本存儲在一組節(jié)點上,以提高可用性和可靠性。鏡像隊列將隊列中的消息復(fù)制到一個或多個其他節(jié)點上,并使這些節(jié)點上的隊列保持同步。當(dāng)一個節(jié)點失敗時,其他節(jié)點上的隊列不受影響,因為它們上面都有消息的備份。
鏡像隊列使用主從模式,所有消息寫入和讀取均通過主節(jié)點,并異步復(fù)制到鏡像節(jié)點。主節(jié)點故障時需重新選舉,期間隊列不可用。而仲裁隊列基于Raft分布式共識算法,所有節(jié)點組成仲裁組。消息需被多數(shù)節(jié)點持久化后才確認成功,Leader故障時自動觸發(fā)選舉。
相比較于傳統(tǒng)的主從模式,避免了發(fā)生網(wǎng)絡(luò)分區(qū)時的腦裂問題(基于Raft分布式共識算法避免)。

和普通隊列的區(qū)別:

相比較于普通隊列,仲裁隊列增加了一個對于有毒消息的處理。什么是有毒消息?首先,消費者從隊列中獲取到了元素,隊列會將該元素刪除,但是消費者消費失敗了,會給隊列nack,并且可以設(shè)置消息重新入隊。這樣可能存在因為業(yè)務(wù)代碼的問題,某條消息一直處理不成功的問題。仲裁隊列會記錄消息的重新投遞次數(shù),判斷是否超過了設(shè)置的閾值,如果超過了就直接丟棄,或者放入死信隊列人工處理。
如果需要聲明一個仲裁隊列,只需要加入?yún)?shù):
@Configuration
public class QuorumConfig {
@Bean
public Queue quorumQueue() {
Map<String,Object> params = new HashMap<>();
params.put("x-queue-type","quorum");
return new Queue(MyConstants.QUEUE_QUORUM,true,false,false,params);
}
}仲裁隊列適用于集群環(huán)境下,隊列長期存在,并且對于消息可靠性要求高,允許犧牲一部分性能(因為raft算法,消息需被多數(shù)節(jié)點持久化后才確認成功)的場景。
1.3、Stream
在傳統(tǒng)的隊列模型中,同一條消息只能被一個消費者消費(一個隊列如果有多個消費者,是工作分發(fā)的機制。消息1->消費者1,消息2->消費者2,消息3->消費者1,不能兩個消費者讀同一條消息。),并且消息是閱后即焚的(消費者接收到消息后,隊列中的該消息就刪除,如果消費者拒絕簽收并且設(shè)置了重新入隊,再把消息重新放入隊列中),無法重復(fù)從隊列中獲取相同的消息。并且在當(dāng)隊列中積累的消息過多時,性能下降會非常明顯。
Stream隊列正是解決了以上的這些問題。Stream隊列的核心是用aof文件的形式存儲隊列,將消息以aof的方式追加到文件中。允許用戶在日志的任何一個連接點開始重新讀取數(shù)據(jù)。(需要用戶自己記錄偏移量)
聲明一個stream隊列:
@Configuration
public class StreamConfig {
@Bean
public Queue streamQueue() {
Map<String,Object> params = new HashMap<>();
params.put("x-queue-type","stream");
params.put("x-max-length-bytes", 20_000_000_000L); // 指定隊列的大小
params.put("x-stream-max-segment-size-bytes", 100_000_000); // 文件分片存儲,每一片的大小
//必須設(shè)置持久化為true,同時獨占和自動刪除模式為false
return new Queue(MyConstants.QUEUE_STREAM,true,false,false,params);
}
}聲明消費者:
public void stremReceiver(Channel channel,String message){
try {
channel.basicQos(100);
Consumer myconsumer = new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag, Envelope envelope,
AMQP.BasicProperties properties, byte[] body)
throws IOException {
System.out.println("========================");
String routingKey = envelope.getRoutingKey();
System.out.println("routingKey >"+routingKey);
String contentType = properties.getContentType();
System.out.println("contentType >"+contentType);
long deliveryTag = envelope.getDeliveryTag();
System.out.println("deliveryTag >"+deliveryTag);
System.out.println("content:"+new String(body,"UTF-8"));
// (process the message components here ...)
//消息處理完后,進行答復(fù)。答復(fù)過的消息,服務(wù)器就不會再次轉(zhuǎn)發(fā)。
//沒有答復(fù)過的消息,服務(wù)器會一直不停轉(zhuǎn)發(fā)。
channel.basicAck(deliveryTag, false);
}
};
Map<String,Object> consumeParam = new HashMap<>();
//first: 從日志隊列中第一個可消費的消息開始消費
//last: 消費消息日志中最后一個消息
//next: 相當(dāng)于不指定offset,消費不到消息。
//Offset: 一個數(shù)字型的偏
//Timestamp:一個代表時間的Data類型變量,表示從這個時間點開始消費。
//例如 一個小時前 Date timestamp = new Date(System.currentTimeMillis() - 60 * 60 * 1_000)
consumeParam.put("x-stream-offset","last");
channel.basicConsume(MyConstants.QUEUE_STREAM, false,consumeParam, myconsumer);
} catch (IOException e) {
e.printStackTrace();
}
System.out.println("quorumReceiver received message : "+ message);
}
1.4、懶隊列
Rabbit MQ對于常規(guī)隊列的處理是,將消息優(yōu)先存在于內(nèi)存中,在合適的時機再持久化到磁盤上,而懶隊列則相反,懶隊列會盡可能早的將消息內(nèi)容保存到磁盤當(dāng)中,并且只有在用戶請求到時,才臨時從磁盤加載到內(nèi)存當(dāng)中。懶隊列的設(shè)計也是為了應(yīng)對消息堆積問題的。
聲明懶隊列的方式,只需要加入?yún)?shù),相應(yīng)的,當(dāng)一個隊列被聲明為懶隊列,那即使隊列被設(shè)定為不持久化,消息依然會寫入到硬盤中。
Map<String, Object> args = new HashMap<String, Object>();
args.put("x-queue-mode", "lazy");懶隊列適合消息量大且長期有堆積的隊列,可以減少內(nèi)存使用,加快消費速度。
到此這篇關(guān)于RabbitMQ隊列的選擇的文章就介紹到這了,更多相關(guān)RabbitMQ隊列內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
Java對接樂橙攝像頭詳細步驟(綁定設(shè)備/直播/控制)
大華樂橙SDK(LechangeSDK)是一套由大華科技推出的智能安防領(lǐng)域?qū)S密浖_發(fā)工具包,下面這篇文章主要介紹了Java對接樂橙攝像頭(綁定設(shè)備/直播/控制)的相關(guān)資料,文中通過代碼介紹的非常詳細,需要的朋友可以參考下2025-12-12
java實現(xiàn)解析二進制文件的方法(字符串、圖片)
本篇文章主要介紹了java實現(xiàn)解析二進制文件的方法(字符串、圖片),小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧2017-02-02
java微信企業(yè)號開發(fā)之開發(fā)模式的開啟
這篇文章主要為大家詳細介紹了java微信企業(yè)號開發(fā)之開發(fā)模式的開啟方法,感興趣的小伙伴們可以參考一下2016-06-06
解讀String字符串導(dǎo)致的JVM內(nèi)存泄漏問題
這篇文章主要介紹了解讀String字符串導(dǎo)致的JVM內(nèi)存泄漏問題,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教2023-07-07
詳解springboot設(shè)置cors跨域請求的兩種方式
這篇文章主要介紹了詳解springboot設(shè)置cors跨域請求的兩種方式,小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧2018-11-11

