最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

在RabbitMQ中實(shí)現(xiàn)Work queues工作隊(duì)列模式

 更新時(shí)間:2021年04月16日 15:00:37   作者:Java_Caiyo  
這篇文章主要介紹了如何在RabbitMQ中實(shí)現(xiàn)Work queues模式,代碼詳細(xì),解釋清晰,可以幫助大家更好理解java,對(duì)這方面感興趣的朋友可以參考下

一、模式說(shuō)明

Work Queues 與入門(mén)程序的簡(jiǎn)單模式相比,多了一個(gè)或一些消費(fèi)端,多個(gè)消費(fèi)端共同消費(fèi)同一個(gè)隊(duì)列中的消息。

應(yīng)用場(chǎng)景 :對(duì)于任務(wù)過(guò)重或任務(wù)較多情況使用工作隊(duì)列可以提高任務(wù)處理的速度。

二、代碼

Work Queues 與入門(mén)程序的 簡(jiǎn)單模式 的代碼是幾乎一樣的:可以完全復(fù)制,并復(fù)制多一個(gè)消費(fèi)者進(jìn)行多個(gè)消費(fèi)者同時(shí)消費(fèi)消息的測(cè)試。

①生產(chǎn)者

package com.itheima.rabbitmq.work; 
import com.itheima.rabbitmq.util.ConnectionUtil; 
import com.rabbitmq.client.Channel; 
import com.rabbitmq.client.Connection; 
import com.rabbitmq.client.ConnectionFactory; 
public class Producer { 
	static final String QUEUE_NAME = "work_queue"; 
	public static void main(String[] args) throws Exception { 
		//創(chuàng)建連接 
		Connection connection = ConnectionUtil.getConnection(); 
		// 創(chuàng)建頻道 
		Channel channel = connection.createChannel(); 
		// 聲明(創(chuàng)建)隊(duì)列 
		/**
		 * 參數(shù)1:隊(duì)列名稱(chēng) 
		 * 參數(shù)2:是否定義持久化隊(duì)列 
		 * 參數(shù)3:是否獨(dú)占本次連接 
		 * 參數(shù)4:是否在不使用的時(shí)候自動(dòng)刪除隊(duì)列 
		 * 參數(shù)5:隊(duì)列其它參數(shù) 
		*/ 
		channel.queueDeclare(QUEUE_NAME, true, false, false, null); 
		for (int i = 1; i <= 30; i++) { 
			// 發(fā)送信息 
			String message = "你好;小兔子!work模式--" + i; 
			/**
			 * 參數(shù)1:交換機(jī)名稱(chēng),如果沒(méi)有指定則使用默認(rèn)Default Exchage 
			 * 參數(shù)2:路由key,簡(jiǎn)單模式可以傳遞隊(duì)列名稱(chēng) 
			 * 參數(shù)3:消息其它屬性 
			 * 參數(shù)4:消息內(nèi)容 
			*/ 
			channel.basicPublish("", QUEUE_NAME, null, message.getBytes()); 
			System.out.println("已發(fā)送消息:" + message); 
		}
		// 關(guān)閉資源 
		channel.close(); connection.close(); 
	} 
}

②消費(fèi)者1

package com.itheima.rabbitmq.work; 
import com.itheima.rabbitmq.util.ConnectionUtil; 
import com.rabbitmq.client.*;
import java.io.IOException; 
public class Consumer1 { 
	public static void main(String[] args) throws Exception { 
		Connection connection = ConnectionUtil.getConnection(); 
		// 創(chuàng)建頻道 
		Channel channel = connection.createChannel(); 
		// 聲明(創(chuàng)建)隊(duì)列 
		/**
		 * 參數(shù)1:隊(duì)列名稱(chēng) 
		 * 參數(shù)2:是否定義持久化隊(duì)列 
		 * 參數(shù)3:是否獨(dú)占本次連接 
		 * 參數(shù)4:是否在不使用的時(shí)候自動(dòng)刪除隊(duì)列 
		 * 參數(shù)5:隊(duì)列其它參數(shù) 
		*/ 
		channel.queueDeclare(Producer.QUEUE_NAME, true, false, false, null); 
		//一次只能接收并處理一個(gè)消息 
		channel.basicQos(1); 
		//創(chuàng)建消費(fèi)者;并設(shè)置消息處理 
		DefaultConsumer consumer = new DefaultConsumer(channel){ 
			@Override 
			/**
			 * consumerTag 消息者標(biāo)簽,在channel.basicConsume時(shí)候可以指定 
			 * envelope 消息包的內(nèi)容,可從中獲取消息id,消息routingkey,交換機(jī),消息和重傳標(biāo)志(收到消息失敗后是否需要重新發(fā)送) 
			 * properties 屬性信息 
			 * body 消息 
			*/ 
			public void handleDelivery(String consumerTag, Envelope envelope, 
					AMQP.BasicProperties properties, byte[] body) throws IOException { 
				try {
					//路由key 
					System.out.println("路由key為:" + envelope.getRoutingKey()); 
					//交換機(jī) 
					System.out.println("交換機(jī)為:" + envelope.getExchange()); 
					//消息id 
					System.out.println("消息id為:" + envelope.getDeliveryTag()); 
					//收到的消息 
					System.out.println("消費(fèi)者1-接收到的消息為:" + new String(body, "utf-8")); 
					Thread.sleep(1000); 
					//確認(rèn)消息 
					channel.basicAck(envelope.getDeliveryTag(), false); 
				} 
				catch (InterruptedException e) { 
					e.printStackTrace(); 
				} 
			} 
		};
		//監(jiān)聽(tīng)消息 
		/**
		 * 參數(shù)1:隊(duì)列名稱(chēng)
		 * 參數(shù)2:是否自動(dòng)確認(rèn),設(shè)置為true為表示消息接收到自動(dòng)向mq回復(fù)接收到了,mq接收到回復(fù)會(huì)刪除消息,設(shè)置為false則需要手動(dòng)確認(rèn) 
		 * 參數(shù)3:消息接收到后回調(diào) 
		*/ 
		channel.basicConsume(Producer.QUEUE_NAME, false, consumer); 
	} 
}

③消費(fèi)者2

package com.itheima.rabbitmq.work; 
import com.itheima.rabbitmq.util.ConnectionUtil; 
import com.rabbitmq.client.*; 
import java.io.IOException; 
public class Consumer2 { 
	public static void main(String[] args) throws Exception { 
		Connection connection = ConnectionUtil.getConnection(); 
		// 創(chuàng)建頻道 
		Channel channel = connection.createChannel(); 
		// 聲明(創(chuàng)建)隊(duì)列 
		/**
		 * 參數(shù)1:隊(duì)列名稱(chēng) 
		 * 參數(shù)2:是否定義持久化隊(duì)列 
		 * 參數(shù)3:是否獨(dú)占本次連接 
		 * 參數(shù)4:是否在不使用的時(shí)候自動(dòng)刪除隊(duì)列 
		 * 參數(shù)5:隊(duì)列其它參數(shù) 
		*/ 
		channel.queueDeclare(Producer.QUEUE_NAME, true, false, false, null); 
		//一次只能接收并處理一個(gè)消息 
		channel.basicQos(1); 
		//創(chuàng)建消費(fèi)者;并設(shè)置消息處理 
		DefaultConsumer consumer = new DefaultConsumer(channel){ 
			@Override 
			/**
			 * consumerTag 消息者標(biāo)簽,在channel.basicConsume時(shí)候可以指定 
			 * envelope 消息包的內(nèi)容,可從中獲取消息id,消息routingkey,交換機(jī),消息和重傳標(biāo)志(收到消息失敗后是否需要重新發(fā)送) 
			 * properties 屬性信息 
			 * body 消息 
			*/ 
			public void handleDelivery(String consumerTag, Envelope envelope, 
					AMQP.BasicProperties properties, byte[] body) throws IOException { 
				try {
					//路由key 
					System.out.println("路由key為:" + envelope.getRoutingKey()); 
					//交換機(jī) 
					System.out.println("交換機(jī)為:" + envelope.getExchange()); 
					//消息id 
					System.out.println("消息id為:" + envelope.getDeliveryTag());
					//收到的消息 
					System.out.println("消費(fèi)者2-接收到的消息為:" + new String(body, "utf-8")); 
					Thread.sleep(1000); 
					//確認(rèn)消息 
					channel.basicAck(envelope.getDeliveryTag(), false); 
				} catch (InterruptedException e) { 
					e.printStackTrace(); 
				} 
			} 
		};
		//監(jiān)聽(tīng)消息 
		/**
		 * 參數(shù)1:隊(duì)列名稱(chēng) 
		 * 參數(shù)2:是否自動(dòng)確認(rèn),設(shè)置為true為表示消息接收到自動(dòng)向mq回復(fù)接收到了,mq接收到回復(fù)會(huì)刪除消息,設(shè)置為false則需要手動(dòng)確認(rèn) 
		 * 參數(shù)3:消息接收到后回調(diào) 
		*/ 
		channel.basicConsume(Producer.QUEUE_NAME, false, consumer); 
	} 
}

三、測(cè)試

啟動(dòng)兩個(gè)消費(fèi)者,然后再啟動(dòng)生產(chǎn)者發(fā)送消息;到IDEA的兩個(gè)消費(fèi)者對(duì)應(yīng)的控制臺(tái)查看是否競(jìng)爭(zhēng)性的接收到消息。

總結(jié)

在一個(gè)隊(duì)列中如果有多個(gè)消費(fèi)者,那么消費(fèi)者之間對(duì)于同一個(gè)消息的關(guān)系是競(jìng)爭(zhēng)的關(guān)系。

到此這篇關(guān)于如何在RabbitMQ中實(shí)現(xiàn)Work queues模式的文章就介紹到這了,希望對(duì)你有所幫助,更多相關(guān)RabbitMQ內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章,希望大家以后多多支持腳本之家!

相關(guān)文章

  • springboot項(xiàng)目或其他項(xiàng)目使用@Test測(cè)試項(xiàng)目接口配置

    springboot項(xiàng)目或其他項(xiàng)目使用@Test測(cè)試項(xiàng)目接口配置

    這篇文章主要介紹了springboot項(xiàng)目或其他項(xiàng)目使用@Test測(cè)試項(xiàng)目接口配置,具有很好的參考價(jià)值,希望對(duì)大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2024-07-07
  • Java聊天室之解決連接超時(shí)問(wèn)題

    Java聊天室之解決連接超時(shí)問(wèn)題

    這篇文章主要為大家詳細(xì)介紹了Java簡(jiǎn)易聊天室之解決連接超時(shí)問(wèn)題的方法,文中的示例代碼講解詳細(xì),具有一定的借鑒價(jià)值,需要的可以了解一下
    2022-10-10
  • hbase訪(fǎng)問(wèn)方式之java api

    hbase訪(fǎng)問(wèn)方式之java api

    這篇文章主要介紹了hbase訪(fǎng)問(wèn)方式之java api,需要的朋友可以參考下
    2017-09-09
  • 解決Feign配置RequestContextHolder.getRequestAttributes()為null的問(wèn)題

    解決Feign配置RequestContextHolder.getRequestAttributes()為null的問(wèn)題

    這篇文章主要介紹了解決Feign配置RequestContextHolder.getRequestAttributes()為null的問(wèn)題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2024-01-01
  • Spring 加載多份配置文件的問(wèn)題及解決方案

    Spring 加載多份配置文件的問(wèn)題及解決方案

    在Spring項(xiàng)目中,有時(shí)候需要加載多份配置文件以簡(jiǎn)化復(fù)雜的配置管理,解決這一問(wèn)題的方法是使用spring.config.import屬性,通過(guò)這種方式,可以在主配置文件中指定額外的配置文件路徑,支持文件、classpath或URL形式的路徑,感興趣的朋友跟隨小編一起看看吧
    2024-10-10
  • java調(diào)用webService接口的代碼實(shí)現(xiàn)

    java調(diào)用webService接口的代碼實(shí)現(xiàn)

    本文主要介紹了java調(diào)用webService接口的代碼實(shí)現(xiàn),文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧
    2023-02-02
  • mybatis?collection和association的區(qū)別解析

    mybatis?collection和association的區(qū)別解析

    這篇文章主要介紹了mybatis?collection解析以及和association的區(qū)別,本文給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2022-07-07
  • 將java程序打包成可執(zhí)行文件的實(shí)現(xiàn)方式

    將java程序打包成可執(zhí)行文件的實(shí)現(xiàn)方式

    本文介紹了將Java程序打包成可執(zhí)行文件的三種方法:手動(dòng)打包(將編譯后的代碼及JRE運(yùn)行環(huán)境一起打包),使用第三方打包工具(如Launch4j)和JDK自帶工具(jpackage),每種方法都有其優(yōu)缺點(diǎn),可根據(jù)實(shí)際需求選擇合適的方式
    2025-02-02
  • 簡(jiǎn)單了解Java類(lèi)成員初始化順序

    簡(jiǎn)單了解Java類(lèi)成員初始化順序

    這篇文章主要介紹了簡(jiǎn)單了解Java類(lèi)成員初始化順序,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下
    2019-11-11
  • java并發(fā)編程專(zhuān)題(八)----(JUC)實(shí)例講解CountDownLatch

    java并發(fā)編程專(zhuān)題(八)----(JUC)實(shí)例講解CountDownLatch

    這篇文章主要介紹了java CountDownLatch的相關(guān)資料,文中示例代碼非常詳細(xì),幫助大家理解和學(xué)習(xí),感興趣的朋友可以了解下
    2020-07-07

最新評(píng)論

利津县| 桦甸市| 蕉岭县| 滦平县| 千阳县| 林西县| 平陆县| 克东县| 清新县| 威信县| 滦平县| 开远市| 随州市| 盐边县| 昭通市| 荔波县| 金阳县| 册亨县| 灌云县| 湖北省| 武功县| 定陶县| 新建县| 泰兴市| 滦平县| 信宜市| 石首市| 祁门县| 通化县| 宁明县| 平安县| 桐梓县| 赤城县| 陇西县| 屯昌县| 曲麻莱县| 碌曲县| 紫金县| 高雄市| 小金县| 四平市|