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...........");} }詳細安裝跟代碼看參考下載:
總結
如果需要指定模式一般是在消費者端設置,靈活性調節。
| Simple(簡單模式少用) | 指定 | 不指定 | 不指定 | 不指定 | 指定 | 不指定 |
| WorkQueue(多個消費者少用) | 指定 | 不指定 | 不指定 | 不指定 | 指定 | 不指定 |
| fanout(publish/subscribe模式) | 不指定 | 指定 | 不指定 | 指定 | 指定 | 不指定 |
| direct(路由模式) | 不指定 | 指定 | 指定 | 指定 | 指定 | 消費者routingKey精確指定多個 |
| topic(主題模糊匹配) | 不指定 | 指定 | 指定 | 指定 | 指定 | 消費者routingKey可以進行模糊匹配 |
用好MySQL的21個好習慣!
2020-11-25
這么簡單的三目運算符,竟然這么多坑?
2020-11-24
5種SpringBoot熱部署方式,你用哪種?
2020-11-23
關注我,每天陪你進步一點點!
總結
以上是生活随笔為你收集整理的3W字!带你玩转「消息队列」的全部內容,希望文章能夠幫你解決所遇到的問題。
- 上一篇: 面试官 | 线程间是如何通信的?
- 下一篇: 面试官 | 说一下什么是代理模式?