最新国产好看的视频,伊人天堂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ì)這方面感興趣的朋友可以參考下

一、模式說明

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

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

二、代碼

Work Queues 與入門程序的 簡(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ì)列名稱 
		 * 參數(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ī)名稱,如果沒有指定則使用默認(rèn)Default Exchage 
			 * 參數(shù)2:路由key,簡(jiǎn)單模式可以傳遞隊(duì)列名稱 
			 * 參數(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ì)列名稱 
		 * 參數(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)聽消息 
		/**
		 * 參數(shù)1:隊(duì)列名稱
		 * 參數(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ì)列名稱 
		 * 參數(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)聽消息 
		/**
		 * 參數(shù)1:隊(duì)列名稱 
		 * 參數(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í)問題

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

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

    hbase訪問方式之java api

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

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

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

    Spring 加載多份配置文件的問題及解決方案

    在Spring項(xiàng)目中,有時(shí)候需要加載多份配置文件以簡(jiǎn)化復(fù)雜的配置管理,解決這一問題的方法是使用spring.config.import屬性,通過這種方式,可以在主配置文件中指定額外的配置文件路徑,支持文件、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),文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(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類成員初始化順序

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

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

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

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

最新評(píng)論

含山县| 阿合奇县| 电白县| 镇江市| 九龙坡区| 景东| 阳信县| 金堂县| 绿春县| 镇巴县| 瑞昌市| 温宿县| 宁乡县| 和龙市| 枞阳县| 从化市| 赤峰市| 镇坪县| 芜湖市| 石楼县| 留坝县| 嘉鱼县| 唐海县| 始兴县| 黎平县| 罗甸县| 易门县| 台北市| 哈密市| 织金县| 新绛县| 哈巴河县| 张家港市| 汝城县| 昌平区| 呼图壁县| 浦县| 丽江市| 巴青县| 锡林郭勒盟| 成安县|