在RabbitMQ中实现Work queues工作队列模式

 更新时间:2021年04月16日 15:00:37   作者:Java_Caiyo  
这篇文章主要介绍了如何在RabbitMQ中实现Work queues模式,代码详细,解释清晰,可以帮助大家更好理解java,对这方面感兴趣的朋友可以参考下

一、模式说明

Work Queues 与入门程序的简单模式相比,多了一个或一些消费端,多个消费端共同消费同一个队列中的消息。

应用场景 :对于任务过重或任务较多情况使用工作队列可以提高任务处理的速度。

二、代码

Work Queues 与入门程序的 简单模式 的代码是几乎一样的:可以完全复制,并复制多一个消费者进行多个消费者同时消费消息的测试。

①生产者

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 { 
		//创建连接 
		Connection connection = ConnectionUtil.getConnection(); 
		// 创建频道 
		Channel channel = connection.createChannel(); 
		// 声明(创建)队列 
		/**
		 * 参数1:队列名称 
		 * 参数2:是否定义持久化队列 
		 * 参数3:是否独占本次连接 
		 * 参数4:是否在不使用的时候自动删除队列 
		 * 参数5:队列其它参数 
		*/ 
		channel.queueDeclare(QUEUE_NAME, true, false, false, null); 
		for (int i = 1; i <= 30; i++) { 
			// 发送信息 
			String message = "你好;小兔子!work模式--" + i; 
			/**
			 * 参数1:交换机名称,如果没有指定则使用默认Default Exchage 
			 * 参数2:路由key,简单模式可以传递队列名称 
			 * 参数3:消息其它属性 
			 * 参数4:消息内容 
			*/ 
			channel.basicPublish("", QUEUE_NAME, null, message.getBytes()); 
			System.out.println("已发送消息:" + message); 
		}
		// 关闭资源 
		channel.close(); connection.close(); 
	} 
}

②消费者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(); 
		// 创建频道 
		Channel channel = connection.createChannel(); 
		// 声明(创建)队列 
		/**
		 * 参数1:队列名称 
		 * 参数2:是否定义持久化队列 
		 * 参数3:是否独占本次连接 
		 * 参数4:是否在不使用的时候自动删除队列 
		 * 参数5:队列其它参数 
		*/ 
		channel.queueDeclare(Producer.QUEUE_NAME, true, false, false, null); 
		//一次只能接收并处理一个消息 
		channel.basicQos(1); 
		//创建消费者;并设置消息处理 
		DefaultConsumer consumer = new DefaultConsumer(channel){ 
			@Override 
			/**
			 * consumerTag 消息者标签,在channel.basicConsume时候可以指定 
			 * envelope 消息包的内容,可从中获取消息id,消息routingkey,交换机,消息和重传标志(收到消息失败后是否需要重新发送) 
			 * properties 属性信息 
			 * body 消息 
			*/ 
			public void handleDelivery(String consumerTag, Envelope envelope, 
					AMQP.BasicProperties properties, byte[] body) throws IOException { 
				try {
					//路由key 
					System.out.println("路由key为:" + envelope.getRoutingKey()); 
					//交换机 
					System.out.println("交换机为:" + envelope.getExchange()); 
					//消息id 
					System.out.println("消息id为:" + envelope.getDeliveryTag()); 
					//收到的消息 
					System.out.println("消费者1-接收到的消息为:" + new String(body, "utf-8")); 
					Thread.sleep(1000); 
					//确认消息 
					channel.basicAck(envelope.getDeliveryTag(), false); 
				} 
				catch (InterruptedException e) { 
					e.printStackTrace(); 
				} 
			} 
		};
		//监听消息 
		/**
		 * 参数1:队列名称
		 * 参数2:是否自动确认,设置为true为表示消息接收到自动向mq回复接收到了,mq接收到回复会删除消息,设置为false则需要手动确认 
		 * 参数3:消息接收到后回调 
		*/ 
		channel.basicConsume(Producer.QUEUE_NAME, false, consumer); 
	} 
}

③消费者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(); 
		// 创建频道 
		Channel channel = connection.createChannel(); 
		// 声明(创建)队列 
		/**
		 * 参数1:队列名称 
		 * 参数2:是否定义持久化队列 
		 * 参数3:是否独占本次连接 
		 * 参数4:是否在不使用的时候自动删除队列 
		 * 参数5:队列其它参数 
		*/ 
		channel.queueDeclare(Producer.QUEUE_NAME, true, false, false, null); 
		//一次只能接收并处理一个消息 
		channel.basicQos(1); 
		//创建消费者;并设置消息处理 
		DefaultConsumer consumer = new DefaultConsumer(channel){ 
			@Override 
			/**
			 * consumerTag 消息者标签,在channel.basicConsume时候可以指定 
			 * envelope 消息包的内容,可从中获取消息id,消息routingkey,交换机,消息和重传标志(收到消息失败后是否需要重新发送) 
			 * properties 属性信息 
			 * body 消息 
			*/ 
			public void handleDelivery(String consumerTag, Envelope envelope, 
					AMQP.BasicProperties properties, byte[] body) throws IOException { 
				try {
					//路由key 
					System.out.println("路由key为:" + envelope.getRoutingKey()); 
					//交换机 
					System.out.println("交换机为:" + envelope.getExchange()); 
					//消息id 
					System.out.println("消息id为:" + envelope.getDeliveryTag());
					//收到的消息 
					System.out.println("消费者2-接收到的消息为:" + new String(body, "utf-8")); 
					Thread.sleep(1000); 
					//确认消息 
					channel.basicAck(envelope.getDeliveryTag(), false); 
				} catch (InterruptedException e) { 
					e.printStackTrace(); 
				} 
			} 
		};
		//监听消息 
		/**
		 * 参数1:队列名称 
		 * 参数2:是否自动确认,设置为true为表示消息接收到自动向mq回复接收到了,mq接收到回复会删除消息,设置为false则需要手动确认 
		 * 参数3:消息接收到后回调 
		*/ 
		channel.basicConsume(Producer.QUEUE_NAME, false, consumer); 
	} 
}

三、测试

启动两个消费者,然后再启动生产者发送消息;到IDEA的两个消费者对应的控制台查看是否竞争性的接收到消息。

总结

在一个队列中如果有多个消费者,那么消费者之间对于同一个消息的关系是竞争的关系。

到此这篇关于如何在RabbitMQ中实现Work queues模式的文章就介绍到这了,希望对你有所帮助,更多相关RabbitMQ内容请搜索脚本之家以前的文章或继续浏览下面的相关文章,希望大家以后多多支持脚本之家!

相关文章

  • 基于swing实现窗体拖拽和拉伸

    基于swing实现窗体拖拽和拉伸

    这篇文章主要为大家详细介绍了基于swing实现窗体拖拽和拉伸,具有一定的参考价值,感兴趣的小伙伴们可以参考一下
    2017-12-12
  • Springmvc Controller接口代码示例

    Springmvc Controller接口代码示例

    这篇文章主要介绍了Springmvc Controller接口代码示例,具有一定参考价值,需要的朋友可以了解下。
    2017-11-11
  • 使用ObjectMapper解析json不用一直new了

    使用ObjectMapper解析json不用一直new了

    这篇文章主要为大家介绍了使用ObjectMapper解析json不用一直new了的方法详解,有需要的朋友可以借鉴参考下,希望能够有所帮助,祝大家多多进步,早日升职加薪
    2022-06-06
  • 轻松掌握Java观察者模式

    轻松掌握Java观察者模式

    这篇文章主要帮助大家轻松掌握Java观察者模式,具有一定的参考价值,感兴趣的小伙伴们可以参考一下
    2016-09-09
  • Spring Security使用Lambda DSL配置流程详解

    Spring Security使用Lambda DSL配置流程详解

    Spring Security 5.2 对 Lambda DSL 语法的增强,允许使用lambda配置HttpSecurity、ServerHttpSecurity,重要提醒,之前的配置方法仍然有效。lambda的添加旨在提供更大的灵活性,但是用法是可选的。让我们看一下HttpSecurity的lambda配置与以前的配置样式相比
    2023-02-02
  • Elasticsearch8.1中的Script使用实例深入解读

    Elasticsearch8.1中的Script使用实例深入解读

    这篇文章主要为大家介绍了Elasticsearch8.1中的Script使用实例深入解读,有需要的朋友可以借鉴参考下,希望能够有所帮助,祝大家多多进步,早日升职加薪
    2023-10-10
  • 基于java构造方法Vector修改元素源码分析

    基于java构造方法Vector修改元素源码分析

    本篇文章是关于ava构造方法Vector源码分析系列文章,本文主要介绍了Vector修改元素的源码分析,有需要的朋友可以借鉴参考下,希望可以有所帮助
    2021-09-09
  • Spring Boot统一异常处理详解

    Spring Boot统一异常处理详解

    我们在做Web应用的时候,请求处理过程中发生错误是非常常见的情况。这个时候就需要统一异常处理了,这篇文章主要给大家介绍了Spring Boot如何进行统一异常处理的相关资料,需要的朋友可以参考借鉴,下面来一起看看吧。
    2017-02-02
  • Java并发编程之ConcurrentLinkedQueue解读

    Java并发编程之ConcurrentLinkedQueue解读

    这篇文章主要介绍了Java并发编程之ConcurrentLinkedQueue解读,非阻塞的实现方式则可以使用循环CAS的方式来实现,而ConcurrentLinkedQueue就是juc包中自带的经典非堵塞方式实现的工具类,需要的朋友可以参考下
    2023-12-12
  • SpringBoot项目中访问HTML页面的三种方法

    SpringBoot项目中访问HTML页面的三种方法

    这篇文章主要介绍了SpringBoot项目中访问HTML页面的三种方法,文中通过代码示例和图文结合的方式讲解的非常详细,对大家的学习或工作有一定的帮助,需要的朋友可以参考下
    2024-07-07

最新评论