Java工作隊(duì)列代碼詳解
我們寫(xiě)了通過(guò)一個(gè)命名的隊(duì)列發(fā)送和接收消息,如果你還不了解請(qǐng)點(diǎn)擊:RabbitMQJava入門(mén)。這篇中我們將會(huì)創(chuàng)建一個(gè)工作隊(duì)列用來(lái)在工作者(consumer)間分發(fā)耗時(shí)任務(wù)。
工作隊(duì)列的主要任務(wù)是:避免立刻執(zhí)行資源密集型任務(wù),然后必須等待其完成。相反地,我們進(jìn)行任務(wù)調(diào)度:我們把任務(wù)封裝為消息發(fā)送給隊(duì)列。工作進(jìn)行在后臺(tái)運(yùn)行并不斷的從隊(duì)列中取出任務(wù)然后執(zhí)行。當(dāng)你運(yùn)行了多個(gè)工作進(jìn)程時(shí),任務(wù)隊(duì)列中的任務(wù)將會(huì)被工作進(jìn)程共享執(zhí)行。
這樣的概念在web應(yīng)用中極其有用,當(dāng)在很短的HTTP請(qǐng)求間需要執(zhí)行復(fù)雜的任務(wù)。
1、準(zhǔn)備
我們使用Thread.sleep來(lái)模擬耗時(shí)的任務(wù)。我們?cè)诎l(fā)送到隊(duì)列的消息的末尾添加一定數(shù)量的點(diǎn),每個(gè)點(diǎn)代表在工作線程中需要耗時(shí)1秒,例如hello…將會(huì)需要等待3秒。
發(fā)送端:
NewTask.java
import java.io.IOException; import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory; public class NewTask { //隊(duì)列名稱 private final static String QUEUE_NAME = "workqueue"; public static void main(String[] args) throws IOException { //創(chuàng)建連接和頻道 ConnectionFactory factory = new ConnectionFactory(); factory.setHost("localhost"); Connection connection = factory.newConnection(); Channel channel = connection.createChannel(); //聲明隊(duì)列 channel.queueDeclare(QUEUE_NAME, false, false, false, null); //發(fā)送10條消息,依次在消息后面附加1-10個(gè)點(diǎn) for (int i = 0; i < 10; i++) { String dots = ""; for (int j = 0; j <= i; j++) { dots += "."; } String message = "helloworld" + dots+dots.length(); channel.basicPublish("", QUEUE_NAME, null, message.getBytes()); System.out.println(" [x] Sent '" + message + "'"); } //關(guān)閉頻道和資源 channel.close(); connection.close(); } }
接收端:
Work.java
import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory; import com.rabbitmq.client.QueueingConsumer; public class Work { //隊(duì)列名稱 private final static String QUEUE_NAME = "workqueue"; public static void main(String[] argv) throws java.io.IOException, java.lang.InterruptedException { //區(qū)分不同工作進(jìn)程的輸出 int hashCode = Work.class.hashCode(); //創(chuàng)建連接和頻道 ConnectionFactory factory = new ConnectionFactory(); factory.setHost("localhost"); Connection connection = factory.newConnection(); Channel channel = connection.createChannel(); //聲明隊(duì)列 channel.queueDeclare(QUEUE_NAME, false, false, false, null); System.out.println(hashCode + " [*] Waiting for messages. To exit press CTRL+C"); QueueingConsumer consumer = new QueueingConsumer(channel); // 指定消費(fèi)隊(duì)列 channel.basicConsume(QUEUE_NAME, true, consumer); while (true) { QueueingConsumer.Delivery delivery = consumer.nextDelivery(); String message = new String(delivery.getBody()); System.out.println(hashCode + " [x] Received '" + message + "'"); doWork(message); System.out.println(hashCode + " [x] Done"); } } /** * 每個(gè)點(diǎn)耗時(shí)1s * @param task * @throws InterruptedException */ private static void doWork(String task) throws InterruptedException { for (char ch : task.toCharArray()) { if (ch == '.') Thread.sleep(1000); } } }
Round-robin 轉(zhuǎn)發(fā)
使用任務(wù)隊(duì)列的好處是能夠很容易的并行工作。如果我們積壓了很多工作,我們僅僅通過(guò)增加更多的工作者就可以解決問(wèn)題,使系統(tǒng)的伸縮性更加容易。
下面我們先運(yùn)行3個(gè)工作者(Work.java)實(shí)例,然后運(yùn)行NewTask.java,3個(gè)工作者實(shí)例都會(huì)得到信息。但是如何分配呢?讓我們來(lái)看輸出結(jié)果:
[x] Sent 'helloworld.1' [x] Sent 'helloworld..2' [x] Sent 'helloworld...3' [x] Sent 'helloworld....4' [x] Sent 'helloworld.....5' [x] Sent 'helloworld......6' [x] Sent 'helloworld.......7' [x] Sent 'helloworld........8' [x] Sent 'helloworld.........9' [x] Sent 'helloworld..........10' 工作者1: 605645 [*] Waiting for messages. To exit press CTRL+C 605645 [x] Received 'helloworld.1' 605645 [x] Done 605645 [x] Received 'helloworld....4' 605645 [x] Done 605645 [x] Received 'helloworld.......7' 605645 [x] Done 605645 [x] Received 'helloworld..........10' 605645 [x] Done 工作者2: 18019860 [*] Waiting for messages. To exit press CTRL+C 18019860 [x] Received 'helloworld..2' 18019860 [x] Done 18019860 [x] Received 'helloworld.....5' 18019860 [x] Done 18019860 [x] Received 'helloworld........8' 18019860 [x] Done 工作者3: 18019860 [*] Waiting for messages. To exit press CTRL+C 18019860 [x] Received 'helloworld...3' 18019860 [x] Done 18019860 [x] Received 'helloworld......6' 18019860 [x] Done 18019860 [x] Received 'helloworld.........9' 18019860 [x] Done
可以看到,默認(rèn)的,RabbitMQ會(huì)一個(gè)一個(gè)的發(fā)送信息給下一個(gè)消費(fèi)者(consumer),而不考慮每個(gè)任務(wù)的時(shí)長(zhǎng)等等,且是一次性分配,并非一個(gè)一個(gè)分配。平均的每個(gè)消費(fèi)者將會(huì)獲得相等數(shù)量的消息。這樣分發(fā)消息的方式叫做round-robin。
2、消息應(yīng)答(messageacknowledgments)
執(zhí)行一個(gè)任務(wù)需要花費(fèi)幾秒鐘。你可能會(huì)擔(dān)心當(dāng)一個(gè)工作者在執(zhí)行任務(wù)時(shí)發(fā)生中斷。我們上面的代碼,一旦RabbItMQ交付了一個(gè)信息給消費(fèi)者,會(huì)馬上從內(nèi)存中移除這個(gè)信息。在這種情況下,如果殺死正在執(zhí)行任務(wù)的某個(gè)工作者,我們會(huì)丟失它正在處理的信息。我們也會(huì)丟失已經(jīng)轉(zhuǎn)發(fā)給這個(gè)工作者且它還未執(zhí)行的消息。
上面的例子,我們首先開(kāi)啟兩個(gè)任務(wù),然后執(zhí)行發(fā)送任務(wù)的代碼(NewTask.java),然后立即關(guān)閉第二個(gè)任務(wù),結(jié)果為:
工作者2: 31054905[*]Waitingformessages.ToexitpressCTRL+C 31054905[x]Received'helloworld..2' 31054905[x]Done 31054905[x]Received'helloworld....4' 工作者1: 18019860[*]Waitingformessages.ToexitpressCTRL+C 18019860[x]Received'helloworld.1' 18019860[x]Done 18019860[x]Received'helloworld...3' 18019860[x]Done 18019860[x]Received'helloworld.....5' 18019860[x]Done 18019860[x]Received'helloworld.......7' 18019860[x]Done 18019860[x]Received'helloworld.........9' 18019860[x]Done
可以看到,第二個(gè)工作者至少丟失了6,8,10號(hào)任務(wù),且4號(hào)任務(wù)未完成。
但是,我們不希望丟失任何任務(wù)(信息)。當(dāng)某個(gè)工作者(接收者)被殺死時(shí),我們希望將任務(wù)傳遞給另一個(gè)工作者。
為了保證消息永遠(yuǎn)不會(huì)丟失,RabbitMQ支持消息應(yīng)答(messageacknowledgments)。消費(fèi)者發(fā)送應(yīng)答給RabbitMQ,告訴它信息已經(jīng)被接收和處理,然后RabbitMQ可以自由的進(jìn)行信息刪除。
如果消費(fèi)者被殺死而沒(méi)有發(fā)送應(yīng)答,RabbitMQ會(huì)認(rèn)為該信息沒(méi)有被完全的處理,然后將會(huì)重新轉(zhuǎn)發(fā)給別的消費(fèi)者。通過(guò)這種方式,你可以確認(rèn)信息不會(huì)被丟失,即使消者偶爾被殺死。
這種機(jī)制并沒(méi)有超時(shí)時(shí)間這么一說(shuō),RabbitMQ只有在消費(fèi)者連接斷開(kāi)是重新轉(zhuǎn)發(fā)此信息。如果消費(fèi)者處理一個(gè)信息需要耗費(fèi)特別特別長(zhǎng)的時(shí)間是允許的。
消息應(yīng)答默認(rèn)是打開(kāi)的。上面的代碼中我們通過(guò)顯示的設(shè)置autoAsk=true關(guān)閉了這種機(jī)制。下面我們修改代碼(Work.java):
boolean ack = false ; //打開(kāi)應(yīng)答機(jī)制 channel.basicConsume(QUEUE_NAME, ack, consumer); //另外需要在每次處理完成一個(gè)消息后,手動(dòng)發(fā)送一次應(yīng)答。 channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
完整修改后的Work.java
import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory; import com.rabbitmq.client.QueueingConsumer; public class Work { //隊(duì)列名稱 private final static String QUEUE_NAME = "workqueue"; public static void main(String[] argv) throws java.io.IOException, java.lang.InterruptedException { //區(qū)分不同工作進(jìn)程的輸出 int hashCode = Work.class.hashCode(); //創(chuàng)建連接和頻道 ConnectionFactory factory = new ConnectionFactory(); factory.setHost("localhost"); Connection connection = factory.newConnection(); Channel channel = connection.createChannel(); //聲明隊(duì)列 channel.queueDeclare(QUEUE_NAME, false, false, false, null); System.out.println(hashCode + " [*] Waiting for messages. To exit press CTRL+C"); QueueingConsumer consumer = new QueueingConsumer(channel); // 指定消費(fèi)隊(duì)列 Boolean ack = false ; //打開(kāi)應(yīng)答機(jī)制 channel.basicConsume(QUEUE_NAME, ack, consumer); while (true) { QueueingConsumer.Delivery delivery = consumer.nextDelivery(); String message = new String(delivery.getBody()); System.out.println(hashCode + " [x] Received '" + message + "'"); doWork(message); System.out.println(hashCode + " [x] Done"); //發(fā)送應(yīng)答 channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false); } } }
測(cè)試:
我們把消息數(shù)量改為5,然后先打開(kāi)兩個(gè)消費(fèi)者(Work.java),然后發(fā)送任務(wù)(NewTask.java),立即關(guān)閉一個(gè)消費(fèi)者,觀察輸出:
[x]Sent'helloworld.1' [x]Sent'helloworld..2' [x]Sent'helloworld...3' [x]Sent'helloworld....4' [x]Sent'helloworld.....5' 工作者2 18019860[*]Waitingformessages.ToexitpressCTRL+C 18019860[x]Received'helloworld..2' 18019860[x]Done 18019860[x]Received'helloworld....4' 工作者1 31054905[*]Waitingformessages.ToexitpressCTRL+C 31054905[x]Received'helloworld.1' 31054905[x]Done 31054905[x]Received'helloworld...3' 31054905[x]Done 31054905[x]Received'helloworld.....5' 31054905[x]Done 31054905[x]Received'helloworld....4' 31054905[x]Done
可以看到工作者2沒(méi)有完成的任務(wù)4,重新轉(zhuǎn)發(fā)給工作者1進(jìn)行完成了。
3、消息持久化(Messagedurability)
我們已經(jīng)學(xué)習(xí)了即使消費(fèi)者被殺死,消息也不會(huì)被丟失。但是如果此時(shí)RabbitMQ服務(wù)被停止,我們的消息仍然會(huì)丟失。
當(dāng)RabbitMQ退出或者異常退出,將會(huì)丟失所有的隊(duì)列和信息,除非你告訴它不要丟失。我們需要做兩件事來(lái)確保信息不會(huì)被丟失:我們需要給所有的隊(duì)列和消息設(shè)置持久化的標(biāo)志。
第一,我們需要確認(rèn)RabbitMQ永遠(yuǎn)不會(huì)丟失我們的隊(duì)列。為了這樣,我們需要聲明它為持久化的。
booleandurable=true;
channel.queueDeclare("task_queue",durable,false,false,null);
注:RabbitMQ不允許使用不同的參數(shù)重新定義一個(gè)隊(duì)列,所以已經(jīng)存在的隊(duì)列,我們無(wú)法修改其屬性。
第二,我們需要標(biāo)識(shí)我們的信息為持久化的。通過(guò)設(shè)置MessageProperties(implementsBasicProperties)值為PERSISTENT_TEXT_PLAIN。
channel.basicPublish("","task_queue",MessageProperties.PERSISTENT_TEXT_PLAIN,message.getBytes());
現(xiàn)在你可以執(zhí)行一個(gè)發(fā)送消息的程序,然后關(guān)閉服務(wù),再重新啟動(dòng)服務(wù),運(yùn)行消費(fèi)者程序做下實(shí)驗(yàn)。
4、公平轉(zhuǎn)發(fā)(Fairdispatch)
或許會(huì)發(fā)現(xiàn),目前的消息轉(zhuǎn)發(fā)機(jī)制(Round-robin)并非是我們想要的。例如,這樣一種情況,對(duì)于兩個(gè)消費(fèi)者,有一系列的任務(wù),奇數(shù)任務(wù)特別耗時(shí),而偶數(shù)任務(wù)卻很輕松,這樣造成一個(gè)消費(fèi)者一直繁忙,另一個(gè)消費(fèi)者卻很快執(zhí)行完任務(wù)后等待。
造成這樣的原因是因?yàn)镽abbitMQ僅僅是當(dāng)消息到達(dá)隊(duì)列進(jìn)行轉(zhuǎn)發(fā)消息。并不在乎有多少任務(wù)消費(fèi)者并未傳遞一個(gè)應(yīng)答給RabbitMQ。僅僅盲目轉(zhuǎn)發(fā)所有的奇數(shù)給一個(gè)消費(fèi)者,偶數(shù)給另一個(gè)消費(fèi)者。
為了解決這樣的問(wèn)題,我們可以使用basicQos方法,傳遞參數(shù)為prefetchCount=1。這樣告訴RabbitMQ不要在同一時(shí)間給一個(gè)消費(fèi)者超過(guò)一條消息。換句話說(shuō),只有在消費(fèi)者空閑的時(shí)候會(huì)發(fā)送下一條信息。
int prefetchCount = 1; channel.basicQos(prefetchCount);
注:如果所有的工作者都處于繁忙狀態(tài),你的隊(duì)列有可能被填充滿。你可能會(huì)觀察隊(duì)列的使用情況,然后增加工作者,或者使用別的什么策略。
測(cè)試:改變發(fā)送消息的代碼,將消息末尾點(diǎn)數(shù)改為6-2個(gè),然后首先開(kāi)啟兩個(gè)工作者,接著發(fā)送消息:
[x] Sent 'helloworld......6' [x] Sent 'helloworld.....5' [x] Sent 'helloworld....4' [x] Sent 'helloworld...3' [x] Sent 'helloworld..2' 工作者1: 18019860 [*] Waiting for messages. To exit press CTRL+C 18019860 [x] Received 'helloworld......6' 18019860 [x] Done 18019860 [x] Received 'helloworld...3' 18019860 [x] Done 工作者2: 31054905 [*] Waiting for messages. To exit press CTRL+C 31054905 [x] Received 'helloworld.....5' 31054905 [x] Done 31054905 [x] Received 'helloworld....4' 31054905 [x] Done 31054905 [x] Received 'helloworld..2' 31054905 [x] Done
可以看出此時(shí)并沒(méi)有按照之前的Round-robin機(jī)制進(jìn)行轉(zhuǎn)發(fā)消息,而是當(dāng)消費(fèi)者不忙時(shí)進(jìn)行轉(zhuǎn)發(fā)。且這種模式下支持動(dòng)態(tài)增加消費(fèi)者,因?yàn)橄⒉](méi)有發(fā)送出去,動(dòng)態(tài)增加了消費(fèi)者馬上投入工作。而默認(rèn)的轉(zhuǎn)發(fā)機(jī)制會(huì)造成,即使動(dòng)態(tài)增加了消費(fèi)者,此時(shí)的消息已經(jīng)分配完畢,無(wú)法立即加入工作,即使有很多未完成的任務(wù)。
5、完整的代碼
NewTask.java
import java.io.IOException; import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory; import com.rabbitmq.client.MessageProperties; public class NewTask { // 隊(duì)列名稱 private final static String QUEUE_NAME = "workqueue_persistence"; public static void main(String[] args) throws IOException { // 創(chuàng)建連接和頻道 ConnectionFactory factory = new ConnectionFactory(); factory.setHost("localhost"); Connection connection = factory.newConnection(); Channel channel = connection.createChannel(); // 聲明隊(duì)列 Boolean durable = true; // 1、設(shè)置隊(duì)列持久化 channel.queueDeclare(QUEUE_NAME, durable, false, false, null); // 發(fā)送10條消息,依次在消息后面附加1-10個(gè)點(diǎn) for (int i = 5; i > 0; i--) { String dots = ""; for (int j = 0; j <= i; j++) { dots += "."; } String message = "helloworld" + dots + dots.length(); // MessageProperties 2、設(shè)置消息持久化 channel.basicPublish("", QUEUE_NAME, MessageProperties.PERSISTENT_TEXT_PLAIN, message.getBytes()); System.out.println(" [x] Sent '" + message + "'"); } // 關(guān)閉頻道和資源 channel.close(); connection.close(); } }
Work.java
import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import com.rabbitmq.client.ConnectionFactory; import com.rabbitmq.client.QueueingConsumer; public class Work { // 隊(duì)列名稱 private final static String QUEUE_NAME = "workqueue_persistence"; public static void main(String[] argv) throws java.io.IOException, java.lang.InterruptedException { // 區(qū)分不同工作進(jìn)程的輸出 int hashCode = Work.class.hashCode(); // 創(chuàng)建連接和頻道 ConnectionFactory factory = new ConnectionFactory(); factory.setHost("localhost"); Connection connection = factory.newConnection(); Channel channel = connection.createChannel(); // 聲明隊(duì)列 Boolean durable = true; channel.queueDeclare(QUEUE_NAME, durable, false, false, null); System.out.println(hashCode + " [*] Waiting for messages. To exit press CTRL+C"); //設(shè)置最大服務(wù)轉(zhuǎn)發(fā)消息數(shù)量 int prefetchCount = 1; channel.basicQos(prefetchCount); QueueingConsumer consumer = new QueueingConsumer(channel); // 指定消費(fèi)隊(duì)列 Boolean ack = false; // 打開(kāi)應(yīng)答機(jī)制 channel.basicConsume(QUEUE_NAME, ack, consumer); while (true) { QueueingConsumer.Delivery delivery = consumer.nextDelivery(); String message = new String(delivery.getBody()); System.out.println(hashCode + " [x] Received '" + message + "'"); doWork(message); System.out.println(hashCode + " [x] Done"); //channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false); channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false); } } /** * 每個(gè)點(diǎn)耗時(shí)1s * * @param task * @throws InterruptedException */ private static void doWork(String task) throws InterruptedException { for (char ch : task.toCharArray()) { if (ch == '.') Thread.sleep(1000); } } }
總結(jié)
以上就是本文關(guān)于Java工作隊(duì)列代碼詳解的全部?jī)?nèi)容,希望對(duì)大家有所幫助。如有不足之處,歡迎留言指出。感謝朋友們對(duì)本站的支持!
相關(guān)文章
Java多線程并發(fā)執(zhí)行demo代碼實(shí)例
這篇文章主要介紹了Java多線程并發(fā)執(zhí)行demo代碼實(shí)例,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下2020-06-06springboot zuul實(shí)現(xiàn)網(wǎng)關(guān)的代碼
這篇文章主要介紹了springboot zuul實(shí)現(xiàn)網(wǎng)關(guān)的代碼,在為服務(wù)架構(gòu)體系里,網(wǎng)關(guān)是非常重要的環(huán)節(jié),他實(shí)現(xiàn)了很多功能,具體哪些功能大家跟隨小編一起通過(guò)本文學(xué)習(xí)吧2018-10-10Intellij idea下使用不同tomcat編譯maven項(xiàng)目的服務(wù)器路徑方法詳解
今天小編就為大家分享一篇關(guān)于Intellij idea下使用不同tomcat編譯maven項(xiàng)目的服務(wù)器路徑方法詳解,小編覺(jué)得內(nèi)容挺不錯(cuò)的,現(xiàn)在分享給大家,具有很好的參考價(jià)值,需要的朋友一起跟隨小編來(lái)看看吧2019-02-02springboot 如何解決cross跨域請(qǐng)求的問(wèn)題
這篇文章主要介紹了springboot 如何解決cross跨域請(qǐng)求的問(wèn)題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2021-10-10springboot?vue接口測(cè)試前后端樹(shù)節(jié)點(diǎn)編輯刪除功能
這篇文章主要為大家介紹了springboot?vue接口測(cè)試前后端樹(shù)節(jié)點(diǎn)編輯刪除功能,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪2022-05-05