3W字!带你玩转「消息队列」
1. 消息隊列解決了什么問題
消息中間件是目前比較流行的一個中間件,其中RabbitMQ更是占有一定的市場份額,主要用來做異步處理、應用解耦、流量削峰、日志處理等等方面。
1. 異步處理
一個用戶登陸網(wǎng)址注冊,然后系統(tǒng)發(fā)短信跟郵件告知注冊成功,一般有三種解決方法。
串行到依次執(zhí)行,問題是用戶注冊后就可以使用了,沒必要等驗證碼跟郵件。
注冊成功后,郵件跟驗證碼用并行等方式執(zhí)行,問題是郵件跟驗證碼是非重要的任務,系統(tǒng)注冊還要等這倆完成么?
基于異步MQ的處理,用戶注冊成功后直接把信息異步發(fā)送到MQ中,然后郵件系統(tǒng)跟驗證碼系統(tǒng)主動去拉取數(shù)據(jù)。
2. 應用解耦
比如我們有一個訂單系統(tǒng),還要一個庫存系統(tǒng),用戶下訂單了就要調(diào)用下庫存系統(tǒng)來處理,直接調(diào)用到話庫存系統(tǒng)出現(xiàn)問題咋辦呢?
3. 流量削峰
舉辦一個 秒殺活動,如何較好到設計?服務層直接接受瞬間搞密度訪問絕對不可以起碼要加入一個MQ。
4. 日志處理
用戶通過WebUI訪問發(fā)送請求到時候后端如何接受跟處理呢一般?
2. RabbitMQ 安裝跟配置
官網(wǎng):https://www.rabbitmq.com/download.html
開發(fā)語言: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。創(chuàng)建一個virtual hosts,一般以/ 開頭。對用戶進行授權(quán),點擊/vhost_mmr,至于WebUI多點點即可了解。
3. 實戰(zhàn)
RabbitMQ 官網(wǎng)支持任務模式:https://www.rabbitmq.com/getstarted.htm
l創(chuàng)建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 消息的生產(chǎn)者 中間: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();//?創(chuàng)建隊列聲明AMQP.Queue.DeclareOk?declareOk?=?channel.queueDeclare(QUEUE_NAME,?false,?false,?false,?null);String?msg?=?"hello?Simple";//?exchange,隊列,參數(shù),消息字節(jié)體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?{//?創(chuàng)建連接Connection?connection?=?ConnectionUtils.getConnection();//?創(chuàng)建頻道Channel?channel?=?connection.createChannel();//?隊列聲明??隊列名,是否持久化,是否獨占模式,無消息后是否自動刪除,消息攜帶參數(shù)channel.queueDeclare(Send.QUEUE_NAME,false,false,false,null);//?定義消費者DefaultConsumer?defaultConsumer?=?new?DefaultConsumer(channel)?{@Override??//?事件模型,消息來了會觸發(fā)該函數(shù)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);}};//?監(jiān)聽隊列channel.basicConsume(Send.QUEUE_NAME,true,defaultConsumer);}//?老方法?消費者 MQ 在3。4以下?用次方法,private?static?void?oldApi()?throws?IOException,?TimeoutException,?InterruptedException?{//?創(chuàng)建連接Connection?connection?=?ConnectionUtils.getConnection();//?創(chuàng)建頻道Channel?channel?=?connection.createChannel();//?定義隊列消費者QueueingConsumer?consumer?=?new?QueueingConsumer(channel);//監(jiān)聽隊列channel.basicConsume(Send.QUEUE_NAME,?true,?consumer);while?(true)?{//?發(fā)貨體QueueingConsumer.Delivery?delivery?=?consumer.nextDelivery();byte[]?body?=?delivery.getBody();String?s?=?new?String(body);System.out.println("---Recv:"?+?s);}} }右上角有可以設置頁面刷新頻率,然后可以在UI界面直接手動消費掉,如下圖:簡單隊列的不足:耦合性過高,生產(chǎn)者一一對應消費者,如果有多個消費者想消費隊列中信息就無法實現(xiàn)了。
2. WorkQueue 工作隊列
Simple隊列中只能一一對應的生產(chǎn)消費,實際開發(fā)中生產(chǎn)者發(fā)消息很簡單,而消費者要跟業(yè)務結(jié)合,消費者接受到消息后要處理從而會耗時。「可能會出現(xiàn)隊列中出現(xiàn)消息積壓」。所以如果多個消費者可以加速消費。
1. round robin 輪詢分發(fā)
代碼編程一個生產(chǎn)者兩個消費者:
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?//?事件觸發(fā)機制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?//?事件觸發(fā)機制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);} }現(xiàn)象:消費者1 跟消費者2 處理的數(shù)據(jù)量完全一樣的個數(shù):消費者1:處理偶數(shù) 消費者2:處理奇數(shù) 這種方式叫輪詢分發(fā)(round-robin)結(jié)果就是不管兩個消費者誰忙,「數(shù)據(jù)總是你一個我一個」,MQ 給兩個消費發(fā)數(shù)據(jù)的時候是不知道消費者性能的,默認就是雨露均沾。此時 autoAck = true。
2. 公平分發(fā) fair dipatch
如果要實現(xiàn)公平分發(fā),要讓消費者消費完畢一條數(shù)據(jù)后就告知MQ,再讓MQ發(fā)數(shù)據(jù)即可。自動應答要關閉!
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);//?每個消費者發(fā)送確認消息之前,消息隊列不發(fā)送下一個消息到消費者,一次只發(fā)送一個消息//?從而限制一次性發(fā)送給消費者到消息不得超過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);//?保證一次只分發(fā)一個channel.basicQos(1);//定義消費者DefaultConsumer?consumer?=?new?DefaultConsumer(channel)?{@Override?//?事件觸發(fā)機制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");//?手動回執(zhí)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);//?保證一次只分發(fā)一個channel.basicQos(1);//定義消費者DefaultConsumer?consumer?=?new?DefaultConsumer(channel)?{@Override?//?事件觸發(fā)機制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");//?手動回執(zhí)channel.basicAck(envelope.getDeliveryTag(),false);}}};//?自動應答boolean?autoAck?=?false;channel.basicConsume(Send.QUEUE_NAME,?autoAck,?consumer);} }結(jié)果:實現(xiàn)了公平分發(fā),消費者2 是消費者1消費數(shù)量的2倍。
3. publish/subscribe 發(fā)布訂閱模式
類似公眾號的訂閱跟發(fā)布,無需指定routingKey:
解讀:
一個生產(chǎn)者多個消費者
每一個消費者都有一個自己的隊列
生產(chǎn)者沒有把消息直接發(fā)送到隊列而是發(fā)送到了交換機轉(zhuǎn)化器(exchange)。
每一個隊列都要綁定到交換機上。
生產(chǎn)者發(fā)送的消息經(jīng)過交換機到達隊列,從而實現(xiàn)一個消息被多個消費者消費。
生產(chǎn)者:
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");//?分發(fā)=?fanout//?發(fā)送消息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);//?綁定隊列到交換機轉(zhuǎn)發(fā)器channel.queueBind(QUEUE_NAME,EXCHANGE_NAME,""?);//?保證一次只分發(fā)一個channel.basicQos(1);//定義消費者DefaultConsumer?consumer?=?new?DefaultConsumer(channel)?{@Override?//?事件觸發(fā)機制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");//?手動回執(zhí)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);//?綁定隊列到交換機轉(zhuǎn)發(fā)器channel.queueBind(QUEUE_NAME,EXCHANGE_NAME,""?);//?保證一次只分發(fā)一個channel.basicQos(1);//定義消費者DefaultConsumer?consumer?=?new?DefaultConsumer(channel)?{@Override?//?事件觸發(fā)機制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");//?手動回執(zhí)channel.basicAck(envelope.getDeliveryTag(),false);}}};//?自動應答boolean?autoAck?=?false;channel.basicConsume(QUEUE_NAME,?autoAck,?consumer);} }「同時還可以自己手動的添加一個隊列監(jiān)控到該exchange」
4. routing 路由選擇 通配符模式
Exchange(交換機,轉(zhuǎn)發(fā)器):「一方面接受生產(chǎn)者消息,另一方面是向隊列推送消息」。匿名轉(zhuǎn)發(fā)用 "" ?表示,比如前面到簡單隊列跟WorkQueue。fanout:不處理路由鍵。「不需要指定routingKey」,我們只需要把隊列綁定到交換機, 「消息就會被發(fā)送到所有到隊列中」。direct:處理路由鍵,「需要指定routingKey」,此時生產(chǎn)者發(fā)送數(shù)據(jù)到時候會指定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?//?事件觸發(fā)機制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");//?手動回執(zhí)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?//?事件觸發(fā)機制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");//?手動回執(zhí)channel.basicAck(envelope.getDeliveryTag(),?false);}}};//?自動應答boolean?autoAck?=?false;channel.basicConsume(QUEUE_NAME,?autoAck,?consumer);} }WebUI:缺點:路由key必須要明確,無法實現(xiàn)規(guī)則性模糊匹配。
5. Topics 主題
將路由鍵跟某個模式匹配,# 表示匹配 >=1個字符, *表示匹配一個。生產(chǎn)者會帶routingKey,但是消費者的MQ會帶模糊routingKey。商品:發(fā)布、刪除、修改、查詢。
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?//?事件觸發(fā)機制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");//?手動回執(zhí)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?//?事件觸發(fā)機制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");//?手動回執(zhí)channel.basicAck(envelope.getDeliveryTag(),?false);}}};//?自動應答boolean?autoAck?=?false;channel.basicConsume(QUEUE_NAME,?autoAck,?consumer);} }6. MQ的持久化跟非持久化
因為消息在內(nèi)存中,如果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是不可以的!因為我們已經(jīng)定義過的test_work_queue,這個queue已聲明為未持久化的。結(jié)論:MQ 不允許修改一個已經(jīng)存在的隊列參數(shù)。
7. 消費者端手動跟自動確認消息
//?自動應答boolean?autoAck?=?false;channel.basicConsume(Send.QUEUE_NAME,?autoAck,?consumer);當MQ發(fā)送數(shù)據(jù)個消費者后,消費者要對收到對信息應答給MQ。
如果autoAck = true 表示「自動確認模式」,一旦MQ把消息分發(fā)給消費者就會把消息從內(nèi)存中刪除。如果消費者收到消息但是還沒有消費完而MQ中數(shù)據(jù)已刪除則會導致丟失了正在處理對消息。
如果autoAck = false表示「手動確認模式」,如果有個消費者掛了,MQ因為沒有收到回執(zhí)信息可以把該信息再發(fā)送給其他對消費者。
MQ支持消息應答(Message acknowledgement),消費者發(fā)送一個消息應答告訴MQ這個消息已經(jīng)被消費了,MQ才從內(nèi)存中刪除。消息應答模式「默認為 false」。
8. RabbitMQ生產(chǎn)者端消息確認機制(事務 + confirm)
在RabbitMQ中我們可以通過持久化來解決MQ服務器異常的數(shù)據(jù)丟失問題,但是「生產(chǎn)者如何確保數(shù)據(jù)發(fā)送到MQ了」?默認情況下生產(chǎn)者也是不知道的。如何解決 呢?
1. AMQP事務
第一種方式AMQP實現(xiàn)了事務機制,類似mysql的事務機制。txSelect:用戶將當前channel設置為transition模式。txCommit:用于提交事務。txRollback:用于回滾事務。
以上都是對生產(chǎn)者對操作。
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模式。
「生產(chǎn)者端confirm實現(xiàn)原理」生產(chǎn)者將信道設置為confirm模式,一旦信道進入了confirm模式,所以該信道上發(fā)布的信息都會被派一個唯一的ID(從1開始),一旦消息被投遞到所有的匹配隊列后,Broker就回發(fā)送一個確認給生產(chǎn)者(包含消息唯一ID),這就使得生產(chǎn)者知道消息已經(jīng)正確到達目的隊列了,如果消息跟隊列是可持久化的,那么確認消息會在消息寫入到磁盤后才發(fā)出。broker回傳給生產(chǎn)者到確認消息中deliver-tag域包含了確認消息到序列號,此外broker也可以設置basic.ack的multiple域,表示這個序列號之前所以信息都已經(jīng)得到處理。
Confirm模式最大的好處在于是異步的。第一條消息發(fā)送后不用一直等待回復后才發(fā)第二條消息。
開啟confirm模式:channel.confimSelect()編程模式:
1. 普通的發(fā)送一個消息后就 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("消息發(fā)送失敗");}?else?{System.out.println("消息發(fā)送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. 批量的發(fā)一批數(shù)據(jù) 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";//?批量發(fā)送for?(int?i?=?0;?i?<?10;?i++)?{channel.basicPublish("",?QUEUE_NAME,?null,?msg.getBytes());}//?確認if?(!channel.waitForConfirms())?{System.out.println("消息發(fā)送失敗");}?else?{System.out.println("消息發(fā)送OK");}channel.close();connection.close();} } --- 接受信息跟上面一樣3. 異步confirm模式,提供一個回調(diào)方法。
Channel對象提供的ConfirmListener()回調(diào)方法只包含deliveryTag(包含當前發(fā)出消息序號),我們需要自己為每一個Channel維護一個unconfirm的消息序號集合,每publish一條數(shù)據(jù),集合中元素加1,每回調(diào)一次handleAck方法,unconfirm集合刪掉響應的一條(multiple=false)或多條(multiple=true)記錄,從運行效率來看,unconfirm集合最好采用有序集合SortedSet存儲結(jié)構(gòu)。
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);//生產(chǎn)者調(diào)用confirmSelectchannel.confirmSelect();//?存放未確認消息final?SortedSet<Long>?confirmSet?=?Collections.synchronizedSortedSet(new?TreeSet<Long>());//?添加監(jiān)聽通道channel.addConfirmListener(new?ConfirmListener()?{//?回執(zhí)有問題的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,如果生產(chǎn)者等routingkey?就會觸發(fā)這個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);}} }總結(jié):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. 生產(chǎn)者
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屬性中的值得隊列,如果沒有則創(chuàng)建,如果有則返回* 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屬性中的值得隊列,如果沒有則創(chuàng)建,如果有則返回*/ @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屬性中的值得隊列,如果沒有則創(chuàng)建,如果有則返回*/ @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...........");} }詳細安裝跟代碼看參考下載:
總結(jié)
如果需要指定模式一般是在消費者端設置,靈活性調(diào)節(jié)。
| Simple(簡單模式少用) | 指定 | 不指定 | 不指定 | 不指定 | 指定 | 不指定 |
| WorkQueue(多個消費者少用) | 指定 | 不指定 | 不指定 | 不指定 | 指定 | 不指定 |
| fanout(publish/subscribe模式) | 不指定 | 指定 | 不指定 | 指定 | 指定 | 不指定 |
| direct(路由模式) | 不指定 | 指定 | 指定 | 指定 | 指定 | 消費者routingKey精確指定多個 |
| topic(主題模糊匹配) | 不指定 | 指定 | 指定 | 指定 | 指定 | 消費者routingKey可以進行模糊匹配 |
用好MySQL的21個好習慣!
2020-11-25
這么簡單的三目運算符,竟然這么多坑?
2020-11-24
5種SpringBoot熱部署方式,你用哪種?
2020-11-23
關注我,每天陪你進步一點點!
總結(jié)
以上是生活随笔為你收集整理的3W字!带你玩转「消息队列」的全部內(nèi)容,希望文章能夠幫你解決所遇到的問題。
- 上一篇: 面试官 | 线程间是如何通信的?
- 下一篇: 面试官 | 说一下什么是代理模式?