日韩性视频-久久久蜜桃-www中文字幕-在线中文字幕av-亚洲欧美一区二区三区四区-撸久久-香蕉视频一区-久久无码精品丰满人妻-国产高潮av-激情福利社-日韩av网址大全-国产精品久久999-日本五十路在线-性欧美在线-久久99精品波多结衣一区-男女午夜免费视频-黑人极品ⅴideos精品欧美棵-人人妻人人澡人人爽精品欧美一区-日韩一区在线看-欧美a级在线免费观看

歡迎訪問 生活随笔!

生活随笔

當前位置: 首頁 > 编程资源 > 编程问答 >内容正文

编程问答

3W字!带你玩转「消息队列」

發布時間:2025/3/11 编程问答 14 豆豆
生活随笔 收集整理的這篇文章主要介紹了 3W字!带你玩转「消息队列」 小編覺得挺不錯的,現在分享給大家,幫大家做個參考.

1. 消息隊列解決了什么問題

消息中間件是目前比較流行的一個中間件,其中RabbitMQ更是占有一定的市場份額,主要用來做異步處理、應用解耦、流量削峰、日志處理等等方面。

1. 異步處理

一個用戶登陸網址注冊,然后系統發短信跟郵件告知注冊成功,一般有三種解決方法。

  • 串行到依次執行,問題是用戶注冊后就可以使用了,沒必要等驗證碼跟郵件。

  • 注冊成功后,郵件跟驗證碼用并行等方式執行,問題是郵件跟驗證碼是非重要的任務,系統注冊還要等這倆完成么?

  • 基于異步MQ的處理,用戶注冊成功后直接把信息異步發送到MQ中,然后郵件系統跟驗證碼系統主動去拉取數據。

  • 2. 應用解耦

    比如我們有一個訂單系統,還要一個庫存系統,用戶下訂單了就要調用下庫存系統來處理,直接調用到話庫存系統出現問題咋辦呢?

    3. 流量削峰

    舉辦一個 秒殺活動,如何較好到設計?服務層直接接受瞬間搞密度訪問絕對不可以起碼要加入一個MQ。

    4. 日志處理

    用戶通過WebUI訪問發送請求到時候后端如何接受跟處理呢一般?

    2. RabbitMQ 安裝跟配置

    官網:https://www.rabbitmq.com/download.html

    開發語言:https://www.erlang.org/

    正式到安裝跟允許需要Erlang跟RabbitMQ倆版本之間相互兼容!我這里圖省事直接用Docker 拉取鏡像了。下載:開啟:管理頁面 默認賬號:guest ?默認密碼:guest 。Docker啟動時候可以指定賬號密碼對外端口以及

    docker?run?-d?--hostname?my-rabbit?--name?rabbit?-e?RABBITMQ_DEFAULT_USER=admin?-e?RABBITMQ_DEFAULT_PASS=admin?-p?15672:15672?-p?5672:5672?-p?25672:25672?-p?61613:61613?-p?1883:1883?rabbitmq:management?

    啟動:用戶添加:vitrual hosts 相當于mysql中的DB。創建一個virtual hosts,一般以/ 開頭。對用戶進行授權,點擊/vhost_mmr,至于WebUI多點點即可了解。

    3. 實戰

    RabbitMQ 官網支持任務模式:https://www.rabbitmq.com/getstarted.htm

    l創建Maven項目導入必要依賴:

    ????<dependencies><dependency><groupId>com.rabbitmq</groupId><artifactId>amqp-client</artifactId><version>4.0.2</version></dependency><dependency><groupId>org.slf4j</groupId><artifactId>slf4j-api</artifactId><version>1.7.10</version></dependency><dependency><groupId>org.slf4j</groupId><artifactId>slf4j-log4j12</artifactId><version>1.7.5</version></dependency><dependency><groupId>log4j</groupId><artifactId>log4j</artifactId><version>1.2.17</version></dependency><dependency><groupId>junit</groupId><artifactId>junit</artifactId><version>4.11</version></dependency></dependencies>

    0. 獲取MQ連接

    package?com.sowhat.mq.util;import?com.rabbitmq.client.Connection; import?com.rabbitmq.client.ConnectionFactory;import?java.io.IOException; import?java.util.concurrent.TimeoutException;public?class?ConnectionUtils?{/***?連接器*?@return*?@throws?IOException*?@throws?TimeoutException*/public?static?Connection?getConnection()?throws?IOException,?TimeoutException?{ConnectionFactory?factory?=?new?ConnectionFactory();factory.setHost("127.0.0.1");factory.setPort(5672);factory.setVirtualHost("/vhost_mmr");factory.setUsername("user_mmr");factory.setPassword("sowhat");Connection?connection?=?factory.newConnection();return?connection;} }

    1. 簡單隊列

    P:Producer 消息的生產者 中間:Queue消息隊列 C:Consumer 消息的消費者

    package?com.sowhat.mq.simple;import?com.rabbitmq.client.AMQP; import?com.rabbitmq.client.Channel; import?com.rabbitmq.client.Connection; import?com.sowhat.mq.util.ConnectionUtils;import?java.io.IOException; import?java.util.concurrent.TimeoutException;public?class?Send?{public?static?final?String?QUEUE_NAME?=?"test_simple_queue";public?static?void?main(String[]?args)?throws?IOException,?TimeoutException?{//?獲取一個連接Connection?connection?=?ConnectionUtils.getConnection();//?從連接獲取一個通道Channel?channel?=?connection.createChannel();//?創建隊列聲明AMQP.Queue.DeclareOk?declareOk?=?channel.queueDeclare(QUEUE_NAME,?false,?false,?false,?null);String?msg?=?"hello?Simple";//?exchange,隊列,參數,消息字節體channel.basicPublish("",?QUEUE_NAME,?null,?msg.getBytes());System.out.println("--send?msg:"?+?msg);channel.close();connection.close();} } --- package?com.sowhat.mq.simple;import?com.rabbitmq.client.*; import?com.sowhat.mq.util.ConnectionUtils;import?java.io.IOException; import?java.util.concurrent.TimeoutException;/***?消費者獲取消息*/ public?class?Recv?{public?static?void?main(String[]?args)?throws?IOException,?TimeoutException,?InterruptedException?{newApi();oldApi();}private?static?void?newApi()?throws?IOException,?TimeoutException?{//?創建連接Connection?connection?=?ConnectionUtils.getConnection();//?創建頻道Channel?channel?=?connection.createChannel();//?隊列聲明??隊列名,是否持久化,是否獨占模式,無消息后是否自動刪除,消息攜帶參數channel.queueDeclare(Send.QUEUE_NAME,false,false,false,null);//?定義消費者DefaultConsumer?defaultConsumer?=?new?DefaultConsumer(channel)?{@Override??//?事件模型,消息來了會觸發該函數public?void?handleDelivery(String?consumerTag,?Envelope?envelope,?AMQP.BasicProperties?properties,?byte[]?body)?throws?IOException?{String?s?=?new?String(body,?"utf-8");System.out.println("---new?api?recv:"?+?s);}};//?監聽隊列channel.basicConsume(Send.QUEUE_NAME,true,defaultConsumer);}//?老方法?消費者 MQ 在3。4以下?用次方法,private?static?void?oldApi()?throws?IOException,?TimeoutException,?InterruptedException?{//?創建連接Connection?connection?=?ConnectionUtils.getConnection();//?創建頻道Channel?channel?=?connection.createChannel();//?定義隊列消費者QueueingConsumer?consumer?=?new?QueueingConsumer(channel);//監聽隊列channel.basicConsume(Send.QUEUE_NAME,?true,?consumer);while?(true)?{//?發貨體QueueingConsumer.Delivery?delivery?=?consumer.nextDelivery();byte[]?body?=?delivery.getBody();String?s?=?new?String(body);System.out.println("---Recv:"?+?s);}} }

    右上角有可以設置頁面刷新頻率,然后可以在UI界面直接手動消費掉,如下圖:簡單隊列的不足:耦合性過高,生產者一一對應消費者,如果有多個消費者想消費隊列中信息就無法實現了。

    2. WorkQueue 工作隊列

    Simple隊列中只能一一對應的生產消費,實際開發中生產者發消息很簡單,而消費者要跟業務結合,消費者接受到消息后要處理從而會耗時。「可能會出現隊列中出現消息積壓」。所以如果多個消費者可以加速消費。

    1. round robin 輪詢分發

    代碼編程一個生產者兩個消費者:

    package?com.sowhat.mq.work;import?com.rabbitmq.client.AMQP; import?com.rabbitmq.client.Channel; import?com.rabbitmq.client.Connection; import?com.sowhat.mq.util.ConnectionUtils;import?java.io.IOException; import?java.util.concurrent.TimeoutException;public?class?Send?{public?static?final?String??QUEUE_NAME?=?"test_work_queue";public?static?void?main(String[]?args)?throws?IOException,?TimeoutException,?InterruptedException?{//?獲取連接Connection?connection?=?ConnectionUtils.getConnection();//?獲取?channelChannel?channel?=?connection.createChannel();//?聲明隊列AMQP.Queue.DeclareOk?declareOk?=?channel.queueDeclare(QUEUE_NAME,?false,?false,?false,?null);for?(int?i?=?0;?i?<50?;?i++)?{String?msg?=?"hello-"?+?i;System.out.println("WQ?send?"?+?msg);channel.basicPublish("",QUEUE_NAME,null,msg.getBytes());Thread.sleep(i*20);}channel.close();connection.close();} }--- package?com.sowhat.mq.work;import?com.rabbitmq.client.*; import?com.sowhat.mq.util.ConnectionUtils;import?java.io.IOException; import?java.util.concurrent.TimeoutException;public?class?Recv1?{public?static?void?main(String[]?args)?throws?IOException,?TimeoutException?{//?獲取連接Connection?connection?=?ConnectionUtils.getConnection();//?獲取通道Channel?channel?=?connection.createChannel();//?聲明隊列channel.queueDeclare(Send.QUEUE_NAME,?false,?false,?false,?null);//定義消費者DefaultConsumer?consumer?=?new?DefaultConsumer(channel)?{@Override?//?事件觸發機制public?void?handleDelivery(String?consumerTag,?Envelope?envelope,?AMQP.BasicProperties?properties,?byte[]?body)?throws?IOException?{String?s?=?new?String(body,?"utf-8");System.out.println("【1】:"?+?s);try?{Thread.sleep(2000);}?catch?(InterruptedException?e)?{e.printStackTrace();}?finally?{System.out.println("【1】?done");}}};boolean?autoAck?=?true;channel.basicConsume(Send.QUEUE_NAME,?autoAck,?consumer);} } --- package?com.sowhat.mq.work;import?com.rabbitmq.client.*; import?com.sowhat.mq.util.ConnectionUtils;import?java.io.IOException; import?java.util.concurrent.TimeoutException;public?class?Recv2?{public?static?void?main(String[]?args)?throws?IOException,?TimeoutException?{//?獲取連接Connection?connection?=?ConnectionUtils.getConnection();//?獲取通道Channel?channel?=?connection.createChannel();//?聲明隊列channel.queueDeclare(Send.QUEUE_NAME,?false,?false,?false,?null);//定義消費者DefaultConsumer?consumer?=?new?DefaultConsumer(channel)?{@Override?//?事件觸發機制public?void?handleDelivery(String?consumerTag,?Envelope?envelope,?AMQP.BasicProperties?properties,?byte[]?body)?throws?IOException?{String?s?=?new?String(body,?"utf-8");System.out.println("【2】:"?+?s);try?{Thread.sleep(1000?);}?catch?(InterruptedException?e)?{e.printStackTrace();}?finally?{System.out.println("【2】?done");}}};boolean?autoAck?=?true;channel.basicConsume(Send.QUEUE_NAME,?autoAck,?consumer);} }

    現象:消費者1 跟消費者2 處理的數據量完全一樣的個數:消費者1:處理偶數 消費者2:處理奇數 這種方式叫輪詢分發(round-robin)結果就是不管兩個消費者誰忙,「數據總是你一個我一個」,MQ 給兩個消費發數據的時候是不知道消費者性能的,默認就是雨露均沾。此時 autoAck = true。

    2. 公平分發 fair dipatch

    如果要實現公平分發,要讓消費者消費完畢一條數據后就告知MQ,再讓MQ發數據即可。自動應答要關閉!

    package?com.sowhat.mq.work;import?com.rabbitmq.client.AMQP; import?com.rabbitmq.client.Channel; import?com.rabbitmq.client.Connection; import?com.sowhat.mq.util.ConnectionUtils;import?java.io.IOException; import?java.util.concurrent.TimeoutException;public?class?Send?{public?static?final?String??QUEUE_NAME?=?"test_work_queue";public?static?void?main(String[]?args)?throws?IOException,?TimeoutException,?InterruptedException?{//?獲取連接Connection?connection?=?ConnectionUtils.getConnection();//?獲取?channelChannel?channel?=?connection.createChannel();//?s聲明隊列AMQP.Queue.DeclareOk?declareOk?=?channel.queueDeclare(QUEUE_NAME,?false,?false,?false,?null);//?每個消費者發送確認消息之前,消息隊列不發送下一個消息到消費者,一次只發送一個消息//?從而限制一次性發送給消費者到消息不得超過1個。int?perfetchCount?=?1;channel.basicQos(perfetchCount);for?(int?i?=?0;?i?<50?;?i++)?{String?msg?=?"hello-"?+?i;System.out.println("WQ?send?"?+?msg);channel.basicPublish("",QUEUE_NAME,null,msg.getBytes());Thread.sleep(i*20);}channel.close();connection.close();} } --- package?com.sowhat.mq.work;import?com.rabbitmq.client.*; import?com.sowhat.mq.util.ConnectionUtils;import?java.io.IOException; import?java.util.concurrent.TimeoutException;public?class?Recv1?{public?static?void?main(String[]?args)?throws?IOException,?TimeoutException?{//?獲取連接Connection?connection?=?ConnectionUtils.getConnection();//?獲取通道final?Channel?channel?=?connection.createChannel();//?聲明隊列channel.queueDeclare(Send.QUEUE_NAME,?false,?false,?false,?null);//?保證一次只分發一個channel.basicQos(1);//定義消費者DefaultConsumer?consumer?=?new?DefaultConsumer(channel)?{@Override?//?事件觸發機制public?void?handleDelivery(String?consumerTag,?Envelope?envelope,?AMQP.BasicProperties?properties,?byte[]?body)?throws?IOException?{String?s?=?new?String(body,?"utf-8");System.out.println("【1】:"?+?s);try?{Thread.sleep(2000);}?catch?(InterruptedException?e)?{e.printStackTrace();}?finally?{System.out.println("【1】?done");//?手動回執channel.basicAck(envelope.getDeliveryTag(),false);}}};//?自動應答boolean?autoAck?=?false;channel.basicConsume(Send.QUEUE_NAME,?autoAck,?consumer);} } --- package?com.sowhat.mq.work;import?com.rabbitmq.client.*; import?com.sowhat.mq.util.ConnectionUtils;import?java.io.IOException; import?java.util.concurrent.TimeoutException;public?class?Recv2?{public?static?void?main(String[]?args)?throws?IOException,?TimeoutException?{//?獲取連接Connection?connection?=?ConnectionUtils.getConnection();//?獲取通道final?Channel?channel?=?connection.createChannel();//?聲明隊列channel.queueDeclare(Send.QUEUE_NAME,?false,?false,?false,?null);//?保證一次只分發一個channel.basicQos(1);//定義消費者DefaultConsumer?consumer?=?new?DefaultConsumer(channel)?{@Override?//?事件觸發機制public?void?handleDelivery(String?consumerTag,?Envelope?envelope,?AMQP.BasicProperties?properties,?byte[]?body)?throws?IOException?{String?s?=?new?String(body,?"utf-8");System.out.println("【2】:"?+?s);try?{Thread.sleep(1000);}?catch?(InterruptedException?e)?{e.printStackTrace();}?finally?{System.out.println("【2】?done");//?手動回執channel.basicAck(envelope.getDeliveryTag(),false);}}};//?自動應答boolean?autoAck?=?false;channel.basicConsume(Send.QUEUE_NAME,?autoAck,?consumer);} }

    結果:實現了公平分發,消費者2 是消費者1消費數量的2倍。

    3. publish/subscribe 發布訂閱模式

    類似公眾號的訂閱跟發布,無需指定routingKey:

    解讀:

  • 一個生產者多個消費者

  • 每一個消費者都有一個自己的隊列

  • 生產者沒有把消息直接發送到隊列而是發送到了交換機轉化器(exchange)。

  • 每一個隊列都要綁定到交換機上。

  • 生產者發送的消息經過交換機到達隊列,從而實現一個消息被多個消費者消費。

  • 生產者:

    package?com.sowhat.mq.ps;import?com.rabbitmq.client.Channel; import?com.rabbitmq.client.Connection; import?com.sowhat.mq.util.ConnectionUtils;import?java.io.IOException; import?java.util.concurrent.TimeoutException;public?class?Send?{public?static?final?String?EXCHANGE_NAME?=?"test_exchange_fanout";public?static?void?main(String[]?args)?throws?IOException,?TimeoutException?{Connection?connection?=?ConnectionUtils.getConnection();Channel?channel?=?connection.createChannel();//聲明交換機channel.exchangeDeclare(EXCHANGE_NAME,"fanout");//?分發=?fanout//?發送消息String?msg?=?"hello?ps?";channel.basicPublish(EXCHANGE_NAME,"",null,msg.getBytes());System.out.println("Send:"?+?msg);channel.close();connection.close();} }

    消息哪兒去了?丟失了,在RabbitMQ中只有隊列有存儲能力,「因為這個時候隊列還沒有綁定到交換機 所以消息丟失了」。消費者:

    package?com.sowhat.mq.ps;import?com.rabbitmq.client.*; import?com.sowhat.mq.util.ConnectionUtils;import?java.io.IOException; import?java.util.concurrent.TimeoutException;public?class?Recv1?{public?static?final?String??QUEUE_NAME?=?"test_queue_fanout_email";public?static?final?String?EXCHANGE_NAME?=?"test_exchange_fanout";public?static?void?main(String[]?args)?throws?IOException,?TimeoutException?{Connection?connection?=?ConnectionUtils.getConnection();final?Channel?channel?=?connection.createChannel();//?隊列聲明channel.queueDeclare(QUEUE_NAME,false,false,false,null);//?綁定隊列到交換機轉發器channel.queueBind(QUEUE_NAME,EXCHANGE_NAME,""?);//?保證一次只分發一個channel.basicQos(1);//定義消費者DefaultConsumer?consumer?=?new?DefaultConsumer(channel)?{@Override?//?事件觸發機制public?void?handleDelivery(String?consumerTag,?Envelope?envelope,?AMQP.BasicProperties?properties,?byte[]?body)?throws?IOException?{String?s?=?new?String(body,?"utf-8");System.out.println("【1】:"?+?s);try?{Thread.sleep(2000);}?catch?(InterruptedException?e)?{e.printStackTrace();}?finally?{System.out.println("【1】?done");//?手動回執channel.basicAck(envelope.getDeliveryTag(),false);}}};//?自動應答boolean?autoAck?=?false;channel.basicConsume(QUEUE_NAME,?autoAck,?consumer);} } --- package?com.sowhat.mq.ps;import?com.rabbitmq.client.*; import?com.sowhat.mq.util.ConnectionUtils;import?java.io.IOException; import?java.util.concurrent.TimeoutException;public?class?Recv2?{public?static?final?String??QUEUE_NAME?=?"test_queue_fanout_sms";public?static?final?String?EXCHANGE_NAME?=?"test_exchange_fanout";public?static?void?main(String[]?args)?throws?IOException,?TimeoutException?{Connection?connection?=?ConnectionUtils.getConnection();final?Channel?channel?=?connection.createChannel();//?隊列聲明channel.queueDeclare(QUEUE_NAME,false,false,false,null);//?綁定隊列到交換機轉發器channel.queueBind(QUEUE_NAME,EXCHANGE_NAME,""?);//?保證一次只分發一個channel.basicQos(1);//定義消費者DefaultConsumer?consumer?=?new?DefaultConsumer(channel)?{@Override?//?事件觸發機制public?void?handleDelivery(String?consumerTag,?Envelope?envelope,?AMQP.BasicProperties?properties,?byte[]?body)?throws?IOException?{String?s?=?new?String(body,?"utf-8");System.out.println("【2】:"?+?s);try?{Thread.sleep(1000);}?catch?(InterruptedException?e)?{e.printStackTrace();}?finally?{System.out.println("【2】?done");//?手動回執channel.basicAck(envelope.getDeliveryTag(),false);}}};//?自動應答boolean?autoAck?=?false;channel.basicConsume(QUEUE_NAME,?autoAck,?consumer);} }

    「同時還可以自己手動的添加一個隊列監控到該exchange」

    4. routing 路由選擇 通配符模式

    Exchange(交換機,轉發器):「一方面接受生產者消息,另一方面是向隊列推送消息」。匿名轉發用 "" ?表示,比如前面到簡單隊列跟WorkQueue。fanout:不處理路由鍵。「不需要指定routingKey」,我們只需要把隊列綁定到交換機, 「消息就會被發送到所有到隊列中」。direct:處理路由鍵,「需要指定routingKey」,此時生產者發送數據到時候會指定key,任務隊列也會指定key,只有key一樣消息才會被傳送到隊列中。如下圖

    package?com.sowhat.mq.routing;import?com.rabbitmq.client.Channel; import?com.rabbitmq.client.Connection; import?com.sowhat.mq.util.ConnectionUtils;import?java.io.IOException; import?java.util.concurrent.TimeoutException;public?class?Send?{public?static?final?String??EXCHANGE_NAME?=?"test_exchange_direct";public?static?void?main(String[]?args)?throws?IOException,?TimeoutException?{Connection?connection?=?ConnectionUtils.getConnection();Channel?channel?=?connection.createChannel();//?exchangechannel.exchangeDeclare(EXCHANGE_NAME,"direct");String?msg?=?"hello?info!";//?可以指定類型String?routingKey?=?"info";channel.basicPublish(EXCHANGE_NAME,routingKey,null,msg.getBytes());System.out.println("Send?:?"?+?msg);channel.close();connection.close();} } --- package?com.sowhat.mq.routing;import?com.rabbitmq.client.*; import?com.sowhat.mq.util.ConnectionUtils;import?java.io.IOException; import?java.util.concurrent.TimeoutException;public?class?Recv1?{public?static?final?String??EXCHANGE_NAME?=?"test_exchange_direct";public?static?final?String?QUEUE_NAME?=?"test_queue_direct_1";public?static?void?main(String[]?args)?throws?IOException,?TimeoutException?{Connection?connection?=?ConnectionUtils.getConnection();final?Channel?channel?=?connection.createChannel();channel.queueDeclare(QUEUE_NAME,false,false,false,null);channel.basicQos(1);channel.queueBind(QUEUE_NAME,EXCHANGE_NAME,"error");//定義消費者DefaultConsumer?consumer?=?new?DefaultConsumer(channel)?{@Override?//?事件觸發機制public?void?handleDelivery(String?consumerTag,?Envelope?envelope,?AMQP.BasicProperties?properties,?byte[]?body)?throws?IOException?{String?s?=?new?String(body,?"utf-8");System.out.println("【1】:"?+?s);try?{Thread.sleep(2000);}?catch?(InterruptedException?e)?{e.printStackTrace();}?finally?{System.out.println("【1】?done");//?手動回執channel.basicAck(envelope.getDeliveryTag(),false);}}};//?自動應答boolean?autoAck?=?false;channel.basicConsume(QUEUE_NAME,?autoAck,?consumer);} } --- package?com.sowhat.mq.routing;import?com.rabbitmq.client.*; import?com.sowhat.mq.util.ConnectionUtils;import?java.io.IOException; import?java.util.concurrent.TimeoutException;public?class?Recv2?{public?static?final?String?EXCHANGE_NAME?=?"test_exchange_direct";public?static?final?String?QUEUE_NAME?=?"test_queue_direct_2";public?static?void?main(String[]?args)?throws?IOException,?TimeoutException?{Connection?connection?=?ConnectionUtils.getConnection();final?Channel?channel?=?connection.createChannel();channel.queueDeclare(QUEUE_NAME,?false,?false,?false,?null);channel.basicQos(1);//?綁定種類似?Keychannel.queueBind(QUEUE_NAME,?EXCHANGE_NAME,?"error");channel.queueBind(QUEUE_NAME,?EXCHANGE_NAME,?"info");channel.queueBind(QUEUE_NAME,?EXCHANGE_NAME,?"warning");//定義消費者DefaultConsumer?consumer?=?new?DefaultConsumer(channel)?{@Override?//?事件觸發機制public?void?handleDelivery(String?consumerTag,?Envelope?envelope,?AMQP.BasicProperties?properties,?byte[]?body)?throws?IOException?{String?s?=?new?String(body,?"utf-8");System.out.println("【2】:"?+?s);try?{Thread.sleep(1000);}?catch?(InterruptedException?e)?{e.printStackTrace();}?finally?{System.out.println("【2】?done");//?手動回執channel.basicAck(envelope.getDeliveryTag(),?false);}}};//?自動應答boolean?autoAck?=?false;channel.basicConsume(QUEUE_NAME,?autoAck,?consumer);} }

    WebUI:缺點:路由key必須要明確,無法實現規則性模糊匹配。

    5. Topics 主題

    將路由鍵跟某個模式匹配,# 表示匹配 >=1個字符, *表示匹配一個。生產者會帶routingKey,但是消費者的MQ會帶模糊routingKey。商品:發布、刪除、修改、查詢。

    package?com.sowhat.mq.topic;import?com.rabbitmq.client.Channel; import?com.rabbitmq.client.Connection; import?com.sowhat.mq.util.ConnectionUtils;import?java.io.IOException; import?java.util.concurrent.TimeoutException;public?class?Send?{public?static?final?String?EXCHANGE_NAME?=?"test_exchange_topic";public?static?void?main(String[]?args)?throws?IOException,?TimeoutException?{Connection?connection?=?ConnectionUtils.getConnection();Channel?channel?=?connection.createChannel();//?exchangechannel.exchangeDeclare(EXCHANGE_NAME,?"topic");String?msg?=?"商品!";//?可以指定類型String?routingKey?=?"goods.find";channel.basicPublish(EXCHANGE_NAME,?routingKey,?null,?msg.getBytes());System.out.println("Send?:?"?+?msg);channel.close();connection.close();} } --- package?com.sowhat.mq.topic;import?com.rabbitmq.client.*; import?com.sowhat.mq.util.ConnectionUtils;import?java.io.IOException; import?java.util.concurrent.TimeoutException;public?class?Recv1?{public?static?final?String??EXCHANGE_NAME?=?"test_exchange_topic";public?static?final?String?QUEUE_NAME?=?"test_queue_topic_1";public?static?void?main(String[]?args)?throws?IOException,?TimeoutException?{Connection?connection?=?ConnectionUtils.getConnection();final?Channel?channel?=?connection.createChannel();channel.queueDeclare(QUEUE_NAME,false,false,false,null);channel.basicQos(1);channel.queueBind(QUEUE_NAME,EXCHANGE_NAME,"goods.add");//定義消費者DefaultConsumer?consumer?=?new?DefaultConsumer(channel)?{@Override?//?事件觸發機制public?void?handleDelivery(String?consumerTag,?Envelope?envelope,?AMQP.BasicProperties?properties,?byte[]?body)?throws?IOException?{String?s?=?new?String(body,?"utf-8");System.out.println("【1】:"?+?s);try?{Thread.sleep(2000);}?catch?(InterruptedException?e)?{e.printStackTrace();}?finally?{System.out.println("【1】?done");//?手動回執channel.basicAck(envelope.getDeliveryTag(),false);}}};//?自動應答boolean?autoAck?=?false;channel.basicConsume(QUEUE_NAME,?autoAck,?consumer);} } --- package?com.sowhat.mq.topic;import?com.rabbitmq.client.*; import?com.sowhat.mq.util.ConnectionUtils;import?java.io.IOException; import?java.util.concurrent.TimeoutException;public?class?Recv2?{public?static?final?String?EXCHANGE_NAME?=?"test_exchange_topic";public?static?final?String?QUEUE_NAME?=?"test_queue_topic_2";public?static?void?main(String[]?args)?throws?IOException,?TimeoutException?{Connection?connection?=?ConnectionUtils.getConnection();final?Channel?channel?=?connection.createChannel();channel.queueDeclare(QUEUE_NAME,?false,?false,?false,?null);channel.basicQos(1);//?此乃重點channel.queueBind(QUEUE_NAME,?EXCHANGE_NAME,?"goods.#");//定義消費者DefaultConsumer?consumer?=?new?DefaultConsumer(channel)?{@Override?//?事件觸發機制public?void?handleDelivery(String?consumerTag,?Envelope?envelope,?AMQP.BasicProperties?properties,?byte[]?body)?throws?IOException?{String?s?=?new?String(body,?"utf-8");System.out.println("【2】:"?+?s);try?{Thread.sleep(1000);}?catch?(InterruptedException?e)?{e.printStackTrace();}?finally?{System.out.println("【2】?done");//?手動回執channel.basicAck(envelope.getDeliveryTag(),?false);}}};//?自動應答boolean?autoAck?=?false;channel.basicConsume(QUEUE_NAME,?autoAck,?consumer);} }

    6. MQ的持久化跟非持久化

    因為消息在內存中,如果MQ掛了那么消息也丟失了,所以應該考慮MQ的持久化。MQ是支持持久化的,

    //?聲明隊列 channel.queueDeclare(Send.QUEUE_NAME,?false,?false,?false,?null);/***?Declare?a?queue*?@see?com.rabbitmq.client.AMQP.Queue.Declare*?@see?com.rabbitmq.client.AMQP.Queue.DeclareOk*?@param?queue?the?name?of?the?queue*?@param?durable?true?if?we?are?declaring?a?durable?queue?(the?queue?will?survive?a?server?restart)*?@param?exclusive?true?if?we?are?declaring?an?exclusive?queue?(restricted?to?this?connection)*?@param?autoDelete?true?if?we?are?declaring?an?autodelete?queue?(server?will?delete?it?when?no?longer?in?use)*?@param?arguments?other?properties?(construction?arguments)?for?the?queue*?@return?a?declaration-confirm?method?to?indicate?the?queue?was?successfully?declared*?@throws?java.io.IOException?if?an?error?is?encountered*/Queue.DeclareOk?queueDeclare(String?queue,?boolean?durable,?boolean?exclusive,?boolean?autoDelete,Map<String,?Object>?arguments)?throws?IOException;

    boolean durable就是表明是否可以持久化,如果我們將程序中的durable = false改為true是不可以的!因為我們已經定義過的test_work_queue,這個queue已聲明為未持久化的。結論:MQ 不允許修改一個已經存在的隊列參數。

    7. 消費者端手動跟自動確認消息

    //?自動應答boolean?autoAck?=?false;channel.basicConsume(Send.QUEUE_NAME,?autoAck,?consumer);

    當MQ發送數據個消費者后,消費者要對收到對信息應答給MQ。

    如果autoAck = true 表示「自動確認模式」,一旦MQ把消息分發給消費者就會把消息從內存中刪除。如果消費者收到消息但是還沒有消費完而MQ中數據已刪除則會導致丟失了正在處理對消息。

    如果autoAck = false表示「手動確認模式」,如果有個消費者掛了,MQ因為沒有收到回執信息可以把該信息再發送給其他對消費者。

    MQ支持消息應答(Message acknowledgement),消費者發送一個消息應答告訴MQ這個消息已經被消費了,MQ才從內存中刪除。消息應答模式「默認為 false」

    8. RabbitMQ生產者端消息確認機制(事務 + confirm)

    在RabbitMQ中我們可以通過持久化來解決MQ服務器異常的數據丟失問題,但是「生產者如何確保數據發送到MQ了」?默認情況下生產者也是不知道的。如何解決 呢?

    1. AMQP事務

    第一種方式AMQP實現了事務機制,類似mysql的事務機制。txSelect:用戶將當前channel設置為transition模式。txCommit:用于提交事務。txRollback:用于回滾事務。

    以上都是對生產者對操作。

    package?com.sowhat.mq.tx;import?com.rabbitmq.client.Channel; import?com.rabbitmq.client.Connection; import?com.sowhat.mq.util.ConnectionUtils;import?java.io.IOException; import?java.util.concurrent.TimeoutException;public?class?TxSend?{public?static?final?String?QUEUE_NAME?=?"test_queue_tx";public?static?void?main(String[]?args)?throws?IOException,?TimeoutException?{Connection?connection?=?ConnectionUtils.getConnection();Channel?channel?=?connection.createChannel();channel.queueDeclare(QUEUE_NAME,?false,?false,?false,?null);String?msg?=?"hello?tx?message";try?{//開啟事務模式channel.txSelect();channel.basicPublish("",?QUEUE_NAME,?null,?msg.getBytes());int?x?=?1?/?0;//?提交事務channel.txCommit();}?catch?(IOException?e)?{//?回滾channel.txRollback();System.out.println("send?message?rollback");}?finally?{channel.close();connection.close();}} } --- package?com.sowhat.mq.tx;import?com.rabbitmq.client.*; import?com.sowhat.mq.util.ConnectionUtils;import?java.io.IOException; import?java.util.concurrent.TimeoutException;public?class?TxRecv?{public?static?final?String?QUEUE_NAME?=?"test_queue_tx";public?static?void?main(String[]?args)?throws?IOException,?TimeoutException?{Connection?connection?=?ConnectionUtils.getConnection();Channel?channel?=?connection.createChannel();channel.queueDeclare(QUEUE_NAME,?false,?false,?false,?null);String?s?=?channel.basicConsume(QUEUE_NAME,?true,?new?DefaultConsumer(channel)?{@Overridepublic?void?handleDelivery(String?consumerTag,?Envelope?envelope,?AMQP.BasicProperties?properties,?byte[]?body)?throws?IOException?{System.out.println("recv[tx]?msg:"?+?new?String(body,?"utf-8"));}});channel.close();connection.close();} }

    缺點就是大量對請求嘗試然后失敗然后回滾,會降低MQ的吞吐量。

    2. Confirm模式。

    「生產者端confirm實現原理」生產者將信道設置為confirm模式,一旦信道進入了confirm模式,所以該信道上發布的信息都會被派一個唯一的ID(從1開始),一旦消息被投遞到所有的匹配隊列后,Broker就回發送一個確認給生產者(包含消息唯一ID),這就使得生產者知道消息已經正確到達目的隊列了,如果消息跟隊列是可持久化的,那么確認消息會在消息寫入到磁盤后才發出。broker回傳給生產者到確認消息中deliver-tag域包含了確認消息到序列號,此外broker也可以設置basic.ack的multiple域,表示這個序列號之前所以信息都已經得到處理。

    Confirm模式最大的好處在于是異步的。第一條消息發送后不用一直等待回復后才發第二條消息。

    開啟confirm模式:channel.confimSelect()編程模式:

    1. 普通的發送一個消息后就 waitForConfirms()
    package?com.sowhat.confirm;import?com.rabbitmq.client.Channel; import?com.rabbitmq.client.Connection; import?com.sowhat.mq.util.ConnectionUtils;import?java.io.IOException; import?java.util.concurrent.TimeoutException;public?class?Send1?{public?static?final?String?QUEUE_NAME?=?"test_queue_confirm1";public?static?void?main(String[]?args)?throws?IOException,?TimeoutException,?InterruptedException?{Connection?connection?=?ConnectionUtils.getConnection();Channel?channel?=?connection.createChannel();channel.queueDeclare(QUEUE_NAME,?false,?false,?false,?null);//?將channel模式設置為 confirm模式,注意設置這個不能設置為事務模式。channel.confirmSelect();String?msg?=?"hello?confirm?message";channel.basicPublish("",?QUEUE_NAME,?null,?msg.getBytes());if?(!channel.waitForConfirms())?{System.out.println("消息發送失敗");}?else?{System.out.println("消息發送OK");}channel.close();connection.close();} } --- package?com.sowhat.confirm;import?com.rabbitmq.client.*; import?com.sowhat.mq.util.ConnectionUtils;import?java.io.IOException; import?java.util.concurrent.TimeoutException;public?class?Recv?{public?static?final?String?QUEUE_NAME?=?"test_queue_confirm1";public?static?void?main(String[]?args)?throws?IOException,?TimeoutException?{Connection?connection?=?ConnectionUtils.getConnection();Channel?channel?=?connection.createChannel();channel.queueDeclare(QUEUE_NAME,?false,?false,?false,?null);String?s?=?channel.basicConsume(QUEUE_NAME,?true,?new?DefaultConsumer(channel)?{@Overridepublic?void?handleDelivery(String?consumerTag,?Envelope?envelope,?AMQP.BasicProperties?properties,?byte[]?body)?throws?IOException?{System.out.println("recv[tx]?msg:"?+?new?String(body,?"utf-8"));}});} }
    2. 批量的發一批數據 waitForConfirms()
    package?com.sowhat.confirm;import?com.rabbitmq.client.Channel; import?com.rabbitmq.client.Connection; import?com.sowhat.mq.util.ConnectionUtils;import?java.io.IOException; import?java.util.concurrent.TimeoutException;public?class?Send2?{public?static?final?String?QUEUE_NAME?=?"test_queue_confirm1";public?static?void?main(String[]?args)?throws?IOException,?TimeoutException,?InterruptedException?{Connection?connection?=?ConnectionUtils.getConnection();Channel?channel?=?connection.createChannel();channel.queueDeclare(QUEUE_NAME,?false,?false,?false,?null);//?將channel模式設置為 confirm模式,注意設置這個不能設置為事務模式。channel.confirmSelect();String?msg?=?"hello?confirm?message";//?批量發送for?(int?i?=?0;?i?<?10;?i++)?{channel.basicPublish("",?QUEUE_NAME,?null,?msg.getBytes());}//?確認if?(!channel.waitForConfirms())?{System.out.println("消息發送失敗");}?else?{System.out.println("消息發送OK");}channel.close();connection.close();} } --- 接受信息跟上面一樣
    3. 異步confirm模式,提供一個回調方法。

    Channel對象提供的ConfirmListener()回調方法只包含deliveryTag(包含當前發出消息序號),我們需要自己為每一個Channel維護一個unconfirm的消息序號集合,每publish一條數據,集合中元素加1,每回調一次handleAck方法,unconfirm集合刪掉響應的一條(multiple=false)或多條(multiple=true)記錄,從運行效率來看,unconfirm集合最好采用有序集合SortedSet存儲結構。

    package?com.sowhat.mq.confirm;import?com.rabbitmq.client.*; import?com.sowhat.mq.util.ConnectionUtils;import?java.io.IOException; import?java.util.Collections; import?java.util.SortedSet; import?java.util.TreeSet; import?java.util.concurrent.TimeoutException;public?class?Send3?{public?static?final?String?QUEUE_NAME?=?"test_queue_confirm3";public?static?void?main(String[]?args)?throws?IOException,?TimeoutException,?InterruptedException?{Connection?connection?=?ConnectionUtils.getConnection();Channel?channel?=?connection.createChannel();channel.queueDeclare(QUEUE_NAME,?false,?false,?false,?null);//生產者調用confirmSelectchannel.confirmSelect();//?存放未確認消息final?SortedSet<Long>?confirmSet?=?Collections.synchronizedSortedSet(new?TreeSet<Long>());//?添加監聽通道channel.addConfirmListener(new?ConfirmListener()?{//?回執有問題的public?void?handleAck(long?deliveryTag,?boolean?multiple)?throws?IOException?{if?(multiple)?{System.out.println("--handleNack---multiple");confirmSet.headSet(deliveryTag?+?1).clear();}?else?{System.out.println("--handleNack--?multiple?false");confirmSet.remove(deliveryTag);}}//?沒有問題的handleAckpublic?void?handleNack(long?deliveryTag,?boolean?multiple)?throws?IOException?{if?(multiple)?{System.out.println("--handleAck---multiple");confirmSet.headSet(deliveryTag?+?1).clear();}?else?{System.out.println("--handleAck--multiple?false");confirmSet.remove(deliveryTag);}}});//?一般情況下是先開啟?消費者,指定好?exchange跟routingkey,如果生產者等routingkey?就會觸發這個return?方法channel.addReturnListener(new?ReturnListener()?{public?void?handleReturn(int?replyCode,?String?replyText,?String?exchange,?String?routingKey,?AMQP.BasicProperties?properties,?byte[]?body)?throws?IOException?{System.out.println("----?handle?return----");System.out.println("replyCode:"?+?replyCode?);System.out.println("replyText:"?+replyText?);System.out.println("exchange:"?+?exchange);System.out.println("routingKey:"?+?routingKey);System.out.println("properties:"?+?properties);System.out.println("body:"?+?new?String(body));}});String?msgStr?=?"sssss";while(true){long?nextPublishSeqNo?=?channel.getNextPublishSeqNo();channel.basicPublish("",QUEUE_NAME,null,msgStr.getBytes());confirmSet.add(nextPublishSeqNo);Thread.sleep(1000);}} }

    總結:AMQP模式相對來說沒Confirm模式性能好些,推薦使用后者。

    9. RabbitMQ延遲隊列 跟死信

    淘寶訂單付款,驗證碼等限時類型服務。

    ????????Map<String,Object>?headers?=??new?HashMap<String,Object>();headers.put("my1","111");headers.put("my2","222");AMQP.BasicProperties?build?=?new?AMQP.BasicProperties().builder().deliveryMode(2).contentEncoding("utf-8").expiration("10000").headers(headers).build();

    死信的處理:

    10. SpringBoot Tpoic Demo

    需求圖:新建SpringBoot 項目添加如下依賴:

    ???????<dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-amqp</artifactId></dependency>
    1. 生產者

    application.yml

    spring:rabbitmq:host:?127.0.0.1username:?adminpassword:?admin

    測試用例:

    package?com.sowhat.mqpublisher;import?org.junit.jupiter.api.Test; import?org.springframework.amqp.core.AmqpTemplate; import?org.springframework.beans.factory.annotation.Autowired; import?org.springframework.boot.test.context.SpringBootTest;@SpringBootTest class?MqpublisherApplicationTests?{@Autowiredprivate?AmqpTemplate?amqpTemplate;@Testvoid?userInfo()?{/***?exchange,routingKey,message*/this.amqpTemplate.convertAndSend("log.topic","user.log.error","Users...");} }
    2. 消費者

    application.xml

    spring:rabbitmq:host:?127.0.0.1username:?adminpassword:?admin#?自定義配置 mq:config:exchange_name:?log.topic#?配置隊列名稱queue_name:info:?log.infoerror:?log.errorlogs:?log.logs

    三個不同的消費者:

    package?com.sowhat.mqconsumer.service;import?org.springframework.amqp.core.ExchangeTypes; import?org.springframework.amqp.rabbit.annotation.Exchange; import?org.springframework.amqp.rabbit.annotation.Queue; import?org.springframework.amqp.rabbit.annotation.QueueBinding; import?org.springframework.amqp.rabbit.annotation.RabbitListener; import?org.springframework.stereotype.Service;/***?@QueueBinding?value屬性:用于綁定一個隊列。@Queue去查找一個名字為value屬性中的值得隊列,如果沒有則創建,如果有則返回* type = ExchangeTypes.TOPIC 指定交換器類型。默認的direct交換器*/ @Service public?class?ErrorReceiverService?{/***?把一個方法跟一個隊列進行綁定,收到消息后綁定給msg*/@RabbitListener(bindings?=?@QueueBinding(value?=?@Queue(value?=?"${mq.config.queue_name.error}"),exchange?=?@Exchange(value?=?"${mq.config.exchange_name}",?type?=?ExchangeTypes.TOPIC),key?=?"*.log.error"))public?void?process(String?msg)?{System.out.println(msg?+?"?Logs...........");} } --- package?com.sowhat.mqconsumer.service;import?org.springframework.amqp.core.ExchangeTypes; import?org.springframework.amqp.rabbit.annotation.Exchange; import?org.springframework.amqp.rabbit.annotation.Queue; import?org.springframework.amqp.rabbit.annotation.QueueBinding; import?org.springframework.amqp.rabbit.annotation.RabbitListener; import?org.springframework.stereotype.Service;/***?@QueueBinding?value屬性:用于綁定一個隊列。*?@Queue去查找一個名字為value屬性中的值得隊列,如果沒有則創建,如果有則返回*/ @Service public?class?InfoReceiverService?{/***?添加一個能夠處理消息的方法*/@RabbitListener(bindings?=?@QueueBinding(value?=?@Queue(value?="${mq.config.queue_name.info}"),exchange?=?@Exchange(value?=?"${mq.config.exchange_name}",type?=?ExchangeTypes.TOPIC),key?=?"*.log.info"))public?void?process(String?msg){System.out.println(msg+"?Info...........");} } -- package?com.sowhat.mqconsumer.service;import?org.springframework.amqp.core.ExchangeTypes; import?org.springframework.amqp.rabbit.annotation.Exchange; import?org.springframework.amqp.rabbit.annotation.Queue; import?org.springframework.amqp.rabbit.annotation.QueueBinding; import?org.springframework.amqp.rabbit.annotation.RabbitListener; import?org.springframework.stereotype.Service;/***?@QueueBinding?value屬性:用于綁定一個隊列。*?@Queue去查找一個名字為value屬性中的值得隊列,如果沒有則創建,如果有則返回*/ @Service public?class?LogsReceiverService?{/***?添加一個能夠處理消息的方法*/@RabbitListener(bindings?=?@QueueBinding(value?=?@Queue(value?="${mq.config.queue_name.logs}"),exchange?=?@Exchange(value?=?"${mq.config.exchange_name}",type?=?ExchangeTypes.TOPIC),key?=?"*.log.*"))public?void?process(String?msg){System.out.println(msg+"?Error...........");} }

    詳細安裝跟代碼看參考下載:

    總結

    如果需要指定模式一般是在消費者端設置,靈活性調節。

    模式生產者Queue生產者exchange生產者routingKey消費者exchange消費者queueroutingKey
    Simple(簡單模式少用)指定不指定不指定不指定指定不指定
    WorkQueue(多個消費者少用)指定不指定不指定不指定指定不指定
    fanout(publish/subscribe模式)不指定指定不指定指定指定不指定
    direct(路由模式)不指定指定指定指定指定消費者routingKey精確指定多個
    topic(主題模糊匹配)不指定指定指定指定指定消費者routingKey可以進行模糊匹配

    往期推薦

    用好MySQL的21個好習慣!

    2020-11-25

    這么簡單的三目運算符,竟然這么多坑?

    2020-11-24

    5種SpringBoot熱部署方式,你用哪種?

    2020-11-23

    關注我,每天陪你進步一點點!

    總結

    以上是生活随笔為你收集整理的3W字!带你玩转「消息队列」的全部內容,希望文章能夠幫你解決所遇到的問題。

    如果覺得生活随笔網站內容還不錯,歡迎將生活随笔推薦給好友。