目錄
一些基本概念:
訊息佇列三大功能:
MQ的四大核心概念:
其他的一些:
基礎代碼:
生產者部分:
消費者部分:
作業佇列:
訊息應答:
自動應答:
手動應答:
訊息自動重新入隊:
代碼部分:
RabbitMQ持久化:
佇列的持久化:
訊息持久化:
發布確認:
不公平發布:
預取值:
發布確認原理:
單個發布確認:
批量發布確認:
一些基本概念:
-
訊息佇列三大功能:
- 流量消峰:超過極限之后,后續的訪問人員需要等待;可以避免宕機,但是需要更多的時間;
- 應用解耦:可以使系統之間解耦,一個系統呼叫別的系統的時候不會因為被呼叫的系統故障而一起發生故障;這樣在呼叫的時候,會通過佇列去訪問別的系統,那么需要呼叫的系統只需要將這樣的一個請求交給了佇列則就完成了他的操作,后續的出錯不會影響它,
- 異步處理:可以使得模塊之間呼叫的時候,呼叫的模塊不再需要等待被呼叫的模塊執行結束,而是被呼叫的模塊執行結束以后由佇列去通知呼叫的模塊;
-
MQ的四大核心概念:
- 生產者;
- 交換機:交換機和佇列是系結的關系,一個交換機可以系結多個佇列;
- 佇列:交換機和佇列是MQ的重要組成部分;佇列和消費者一一對應;
- 消費者;
-
其他的一些:
- 生產者和交換機、佇列(broker快取代理)之間通過connection中的chennel(信道)連接,
基礎代碼:
-
生產者部分:
/**
* 生產者
*/
public class Producer {
public static final String Queue_Name="hello";
public static void main(String[] args) throws IOException, TimeoutException {
//創建連接工廠
ConnectionFactory connectionFactory = new ConnectionFactory();
//工廠ip 連接mq
connectionFactory.setHost("172.20.10.6");
connectionFactory.setUsername("ljw");
connectionFactory.setPassword("666666");
Connection connection = connectionFactory.newConnection();
//連接需要通過信道發送訊息
Channel channel = connection.createChannel();
//通過信道獲取佇列
//1.佇列名、佇列中的訊息是否需要持久化---默認存記憶體,持久化后存磁盤;
//2.是否進行訊息共享,是否可以被多個消費者共享,true:不共享,false:共享;
//3.是否自動洗掉
channel.queueDeclare(Queue_Name,false,false,false,null);
String message="hello world";
//1.交換機
//2.路由的key
channel.basicPublish("",Queue_Name,null,message.getBytes(StandardCharsets.UTF_8));
System.out.println("訊息發送完畢");
}
}
-
消費者部分:
/**
* 消費者接收訊息
*/
public class Consumer {
public static final String Queue_Name = "hello";
//接收訊息
public static void main(String[] args) throws IOException, TimeoutException {
//創建連接工廠
ConnectionFactory connectionFactory = new ConnectionFactory();
//工廠ip 連接mq
connectionFactory.setHost("172.20.10.6");
connectionFactory.setUsername("ljw");
connectionFactory.setPassword("666666");
Connection connection = connectionFactory.newConnection();
//連接需要通過信道接受訊息
Channel channel = connection.createChannel();
//宣告 接收訊息
DeliverCallback deliverCallback = (var1, var2) -> {
//將byte陣列轉化為String型別
String message = new String(var2.getBody(), StandardCharsets.UTF_8);;
System.out.println(message);
};
//宣告 取消訊息
CancelCallback cancelCallback=(var1)->{
System.out.println("消費被中斷");
};
//1.消費哪個佇列
//2.消費成功以后是否自動應答
//3.消費者收到訊息的回呼
//4.消費者取消消費的回呼
channel.basicConsume(Queue_Name, true,deliverCallback, cancelCallback);
}
}
綜上,可以將獲取信道的程序封裝到工具類里面:
public class RabbitUtils {
public Channel getChannel() throws IOException, TimeoutException {
//創建連接工廠
ConnectionFactory connectionFactory = new ConnectionFactory();
//工廠ip 連接mq
connectionFactory.setHost("172.20.10.6");
connectionFactory.setUsername("ljw");
connectionFactory.setPassword("666666");
Connection connection = connectionFactory.newConnection();
//連接需要通過信道接受訊息
Channel channel = connection.createChannel();
return channel;
}
}
作業佇列:
- 概念:避免立即執行資源密集型任務,將這些任務封裝成訊息在后臺執行,可以由多個執行緒一起處理,那么多個執行緒操作的時候需要保證一條訊息只能被處理一次,各個執行緒輪詢作業,
- 測驗的時候可以設定允許并行操作,模擬多個執行緒---同時跑兩遍;


- 此時需要生產者發送大量訊息:
/** * 生產者 */ public class Producer { public static final String QUEUE_NAME="hello"; public static void main(String[] args) throws IOException, TimeoutException { Channel channel = RabbitUtils.getChannel(); //通過信道獲取佇列 //1.佇列名、佇列中的訊息是否需要持久化---默認存記憶體,持久化后存磁盤; //2.是否進行訊息共享,是否可以被多個消費者共享,true:不共享,false:共享; //3.是否自動洗掉 channel.queueDeclare(QUEUE_NAME,false,false,false,null); String message="hello world"; //1.交換機 //2.路由的key Scanner scanner = new Scanner(System.in); while (scanner.hasNext()){ scanner=new Scanner(System.in); channel.basicPublish("", QUEUE_NAME, null, message.getBytes(StandardCharsets.UTF_8)); System.out.println("訊息發送完畢"); } } }運行會發現,消費者是輪詢消費多條訊息的;
訊息應答:
- 如果沒有訊息應答那么rabbitmq一旦向消費者傳遞了一條資訊,就會立刻將該條訊息標記為洗掉,那么如果一旦有一個消費者掛掉了,那么就會丟失訊息;那么為了保證訊息在發送程序中不丟失,就產生了訊息應答機制,即:消費者在接收到訊息并且處理訊息之后,會告訴rabbitmq處理好了,此時rabbitmq可以將訊息洗掉了,
-
自動應答:
-
訊息發送以后立即被認為已經傳送成功,這種情況下,如果訊息在接收到之前,消費者的channel關閉了,那么訊息就會丟失;另一方面,這種情況沒有對傳遞訊息的數量進行限制,那么可能會導致消費者來不及處理訊息,導致訊息積壓記憶體耗盡,最終使得消費者執行緒被作業系統殺死,所以這種應答方式僅適用于消費者可以高效并以某種速率處理這些訊息的情況下,
-
手動應答:
- Channel.basicAck(用于肯定確認):消費者已經接收到訊息并且成功處理了訊息,rabbitmq可以丟棄該條訊息了,
- Channel.basicNack(用于否認確認):
- Channel.basicReject(用于否認確認):比Channel.basicNack少一個引數--是否批量處理,表示不處理該訊息直接拒絕,rabbitmq可以丟棄該條訊息了,
- 批量處理:如果此時channel中有多條未應答訊息,批量處理為true的話,那么這多條訊息都會收到訊息應答,如果批量處理為false的話,那么只有最新的一條未應答訊息收到訊息應答,批量處理的話有可能導致后續訊息處理失敗,但是rabbitmq已經收到應答從而導致訊息的丟失,所以不推薦使用批量處理,
-
訊息自動重新入隊:
- 如果消費者由于某些原因失去連接(如channel關閉了),導致訊息未發送ack確認,那么rabbitmq將會知道訊息沒有完全處理,就會將訊息重新排隊,此時如果有其他消費者可以處理,那么將會重新分發給其他消費者,這樣就不會導致訊息的丟失,
-
代碼部分:
- 生產者:
public class Producer { //佇列名稱 private static final String QUEUE_NAME = "ack_queue"; public static void main(String[] args) throws IOException, TimeoutException { Channel channel = RabbitUtils.getChannel(); //申明佇列 channel.queueDeclare(QUEUE_NAME, false, false, false, null); Scanner scanner = new Scanner(System.in); String message = "hello world"; while (scanner.hasNext()) { scanner = new Scanner(System.in); channel.basicPublish("", QUEUE_NAME, null, message.getBytes(StandardCharsets.UTF_8)); System.out.println("訊息已發送"); } } } - 消費者:
public class Consumer01 { //佇列名稱 private static final String QUEUE_NAME = "ack_queue"; public static void main(String[] args) throws IOException, TimeoutException { Channel channel = RabbitUtils.getChannel(); System.out.println("C1等待接收訊息處理時間較短"); //這里寫收到訊息后如何消費 DeliverCallback deliverCallback = (var1, var2) -> { try { Thread.sleep(10000L); } catch (InterruptedException e) { e.printStackTrace(); } String message = new String(var2.getBody(), StandardCharsets.UTF_8); System.out.println("接收到的訊息是:"+message); //手動應答 //1.訊息的標記 tag -->envelope是屬性 //2.批量應答 channel.basicAck(var2.getEnvelope().getDeliveryTag(),false); }; //宣告 取消訊息 CancelCallback cancelCallback=(var1)->{ System.out.println("消費被中斷"); }; //關閉自動應答,采用手動應答 boolean autoAck=false; channel.basicConsume(QUEUE_NAME,autoAck,deliverCallback,cancelCallback); } }public class Consumer02 { //佇列名稱 private static final String QUEUE_NAME = "ack_queue"; public static void main(String[] args) throws IOException, TimeoutException { Channel channel = RabbitUtils.getChannel(); System.out.println("C2等待接收訊息處理時間較長"); //這里寫收到訊息后如何消費 DeliverCallback deliverCallback = (var1, var2) -> { try { Thread.sleep(300000L); } catch (InterruptedException e) { e.printStackTrace(); } String message = new String(var2.getBody(), StandardCharsets.UTF_8); System.out.println("接收到的訊息是:"+message); //手動應答 //1.訊息的標記 tag -->envelope是屬性 //2.批量應答 channel.basicAck(var2.getEnvelope().getDeliveryTag(),false); }; //宣告 取消訊息 CancelCallback cancelCallback=(var1)->{ System.out.println("消費被中斷"); }; //關閉自動應答,采用手動應答 boolean autoAck=false; channel.basicConsume(QUEUE_NAME,autoAck,deliverCallback,cancelCallback); } }測驗后發現,如果在訊息處理的程序中某個消費者掛掉了,那么該條訊息會重新入隊并且很快分給其他消費者;并且如果開啟手動應答以后,會在手動應答執行以后才將訊息從佇列里洗掉,
RabbitMQ持久化:
- 以上是保證了傳給消費者進行處理的程序中訊息不丟失,那么如果是rabbitmq服務停掉以后,生產者發來訊息,這個訊息如何保證不丟失呢?那么也就是將佇列和訊息標記為持久化,
-
佇列的持久化:
-
將佇列申明中的durable引數改為true,那么也就是開啟持久化;如果沒有開啟持久化,那么重啟mq以后該佇列就不存在了,但是持久化的仍舊存在,
- 特別注意:如果是已經創建了沒有開啟持久化的佇列,那么需要將其洗掉重新創建持久化的佇列,否則會報錯;
//申明佇列 //第二個引數即durable,true則為開啟持久化 channel.queueDeclare(QUEUE_NAME, true, false, false, null);
(開啟持久化的佇列在feature一列會顯示“D”)
-
訊息持久化:
-
是在生產者發布訊息的時候開啟持久化,也就是在props引數加上開啟持久化屬性(MessageProperties.PERSISTENT_TEXT_PLAIN),
- 特別注意:將訊息標記為持久化并不能完全保證不會丟失訊息,這里會存在訊息剛準備保存在磁盤的時候,但還沒有儲存完就宕機了,訊息在快取的一個間隔點,也就是沒有真正的寫進磁盤,如果需要更強的持久化策略,需要參考發布確認,
//設定開啟訊息持久化(MessageProperties.PERSISTENT_TEXT_PLAIN) channel.basicPublish("",QUEUE_NAME, MessageProperties.PERSISTENT_TEXT_PLAIN, message.getBytes(StandardCharsets.UTF_8));
發布確認:
-
不公平發布:
-
輪詢分發相當于是一種公平發布,在這種情況下,如果一個消費者處理訊息特別快,一個消費者處理訊息特別慢,那么處理訊息快的消費者大部分時間是空閑的,而處理訊息慢的消費者則是一直處于作業狀態,所以這樣的發布方式不是非常合理,那么為了避免這種情況,我們可以在每個消費者端設定引數channel.basicQos(int prefetchCount=1)(默認是0,也就是輪詢分發),從而變為不公平發布,也就相當于能者多勞,也就是一個消費者只能同一時刻只能處理一個訊息,要是目前的訊息還沒處理好,就不會分給他新的訊息;注意:此設定需要在手動應答的時候才生效;
-
預取值:
- 也就是指定每個消費者分到幾條訊息,分配合理可以提高效率;
//consumer01 int prefetchCount=3; channel.basicQos(prefetchCount);//consumer02 int prefetchCount=3; channel.basicQos(prefetchCount);

-
發布確認原理:
- 發布確認就是存在磁盤上完成以后,mq告知生產者,
- 開啟發布確認就可以保證持久化,真的存盤在磁盤上(前提是開啟了佇列、訊息的持久化),
- 開啟發布確認:需要在生產者的channel上呼叫channel.confirmSelect();
//開啟發布確認 channel.confirmSelect();
-
單個發布確認:
- 這是一種同步發布確認的方式,發布訊息后只有被確認了才會發布下一條;如果指定時間內沒有被確認,那么就會拋出例外;缺點:發布速度特別慢,
//單個確認 public static void publishMessageIndividually() throws IOException, TimeoutException, InterruptedException { Channel channel = RabbitUtils.getChannel(); //通過uuid獲取一個隨機的佇列名 String queueName = UUID.randomUUID().toString(); channel.queueDeclare(queueName, true, false, false, null); //開啟發布確認 channel.confirmSelect(); //開始時間 long begin = System.currentTimeMillis(); //批量發訊息 for (int i=0;i<MESSAGE_COUNTS;i++){ String message=i+""; channel.basicPublish("",queueName, MessageProperties.PERSISTENT_TEXT_PLAIN,message.getBytes(StandardCharsets.UTF_8)); //單個訊息就發布確認 boolean flag = channel.waitForConfirms(); if (flag){ System.out.println("發布成功"); } } //結束時間 long end = System.currentTimeMillis(); System.out.println("發布"+MESSAGE_COUNTS+"個單獨發布確認訊息用時:"+(end - begin)+"ms"); } -
單個確認發布耗時如圖:

-
批量發布確認:
- 發布一批訊息然后一起確認,提高了吞吐量;缺點:如果發布出現問題,就不知道是哪個訊息出現的問題,
//批量確認 public static void publishMessageBatch() throws IOException, TimeoutException, InterruptedException { Channel channel = RabbitUtils.getChannel(); //通過uuid獲取一個隨機的佇列名 String queueName = UUID.randomUUID().toString(); channel.queueDeclare(queueName, true, false, false, null); //開啟發布確認 channel.confirmSelect(); //開始時間 long begin = System.currentTimeMillis(); //設定批量確認訊息的數量 int batchCount = 100; //批量發訊息 for (int i = 1; i <= MESSAGE_COUNTS; i++) { String message = i + ""; channel.basicPublish("", queueName, MessageProperties.PERSISTENT_TEXT_PLAIN, message.getBytes(StandardCharsets.UTF_8)); //判斷每到達100條的時候,批量確認一次 if (i % batchCount == 0) { boolean flag = channel.waitForConfirms(); if (flag) { System.out.println("發布成功"); } } } //結束時間 long end = System.currentTimeMillis(); System.out.println("發布" + MESSAGE_COUNTS + "個批量發布確認訊息用時:" + (end - begin) + "ms"); }代碼如圖
-
批量確認發布耗時如圖:

轉載請註明出處,本文鏈接:https://www.uj5u.com/qita/436403.html
標籤:其他
