rabbitmq 学习三-spring-boot中配置交换器和队列

2024-03-14 21:48

本文主要是介绍rabbitmq 学习三-spring-boot中配置交换器和队列,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

rabbitmq主要的交换器类型有fanout、direct、topic、headers
下面分别介绍三种常用的交换器使用方法
1.fanout交换器
将所有发送到该交换器的消息会路由到所有与该交换器绑定的队列中
2.direct交换器
将消息路由到RoutingKey完全匹配的队列中
3.topic 交换器
将消息路由到RoutingKey匹配的队列中,匹配的规则支持特殊字符 ”*“和“#”
一、fanout交换器使用和配置
1.声明队列、交换器并将队列绑定到交换器

@Configuration
public class RabbitFanoutExchangeConfiguration {/*** 声明队列* @return*/@Bean(name = "q.test1")public Queue queue() {return new Queue("q.test1");}@Bean("q.test2")public Queue queue2() {return new Queue("q.test2");}/*** 声明交换器* @return*/@Beanpublic FanoutExchange fanoutExchange() {return new FanoutExchange("x.test1");}/*** 绑定队列到交换器* @return*/@Beanpublic Binding bindingQueue2Exchange() {return BindingBuilder.bind(queue()).to(fanoutExchange());}@Beanpublic Binding bindingQueue2FanoutExchange() {return BindingBuilder.bind(queue2()).to(fanoutExchange());}
}

2.创建消息生产者

@Component
public class RabbitProducer {@AutowiredAmqpTemplate amqpTemplate;public void sendMessage(Object object) {try {Message message= MessageBuilder.withBody(JSON.toJSONString(object).getBytes("UTF-8")).setContentType(MessageProperties.CONTENT_TYPE_JSON).setContentEncoding("UTF-8").setMessageId(UUID.randomUUID().toString()).build();// 指定exchange 为x.test1amqpTemplate.send("x.test1","",message);} catch (UnsupportedEncodingException e) {e.printStackTrace();}}
}

3.创建消息消费者

@Component
@RabbitListener(queues = "q.test1")
public class FanoutExchangeConsumer1 {@RabbitHandlerpublic void receiveMessage(Object object) {Message message=(Message)object;byte bytes[]=null;if (message != null) {bytes=message.getBody();}try {String msg=new String(bytes,"UTF-8");System.out.println("q.test1------------------------"+msg);} catch (UnsupportedEncodingException e) {e.printStackTrace();}}
}@Component
@RabbitListener(queues = "q.test2")
public class FanoutExchangeConsumer2 {@RabbitHandlerpublic void receiveMessage(Object object) {Message message=(Message)object;byte bytes[]=null;if (message != null) {bytes=message.getBody();}try {String msg=new String(bytes,"UTF-8");System.out.println("q.test2------------------------"+msg);} catch (UnsupportedEncodingException e) {e.printStackTrace();}}
}

二、direct交换器使用和配置
1.声明队列、交换器并将队列绑定到交换器

@Configuration
public class RabbitDirectExchangeConfiguration {/*** 动态声明队列* @return*/@Bean(name = "q.direct1")public Queue queue1() {return new Queue("q.direct1");}/*** 动态声明队列* @return*/@Bean(name = "q.direct2")public Queue queue2() {return new Queue("q.direct2");}/*** 动态声明交换器* @return*/@Beanpublic DirectExchange directExchange() {return new DirectExchange("x.direct");}/*** 使用路由键r.direct.routingKey1将交换器绑定到队列 q.direct1* @return*/@Beanpublic Binding bindingQueue1Exchange() {return BindingBuilder.bind(queue1()).to(directExchange()).with("r.direct.routingKey1");}/*** 使用路由键r.direct.routingKey2 将交换器绑定到队列 q.direct1* @return*/@Beanpublic Binding bindingExchange2Queue2() {return BindingBuilder.bind(queue1()).to(directExchange()).with("r.direct.routingKey2");}/*** 使用路由键 r.direct.routingKey1将交换器绑定到队列 q.direct2* @return*/@Beanpublic Binding bindingExchange2Queue() {return BindingBuilder.bind(queue2()).to(directExchange()).with("r.direct.routingKey1");}
}

2.创建消息生产者

public void sendMessage2DirectExchange(Object object) {try {Message message= MessageBuilder.withBody(JSON.toJSONString(object).getBytes("UTF-8")).setContentType(MessageProperties.CONTENT_TYPE_JSON).setContentEncoding("UTF-8").setMessageId(UUID.randomUUID().toString()).build();// 指定exchange 为x.direct, 路由键为 r.direct.routingKey1amqpTemplate.send("x.direct","r.direct.routingKey1",message);} catch (UnsupportedEncodingException e) {e.printStackTrace();}}public void sendMessage2DirectDirectExchange(Object object) {try {Message message= MessageBuilder.withBody(JSON.toJSONString(object).getBytes("UTF-8")).setContentType(MessageProperties.CONTENT_TYPE_JSON).setContentEncoding("UTF-8").setMessageId(UUID.randomUUID().toString()).build();// 指定exchange 为x.direct, 路由键为 r.direct.routingKey2amqpTemplate.send("x.direct","r.direct.routingKey2",message);} catch (UnsupportedEncodingException e) {e.printStackTrace();}}

3.创建消息消费者

@Component
@RabbitListener(queues = "q.direct1")
public class DirectExchangeConsumer1 {@RabbitHandlerpublic void receiveMessage(Object object) {Message message=(Message)object;byte bytes[]=null;if (message != null) {bytes=message.getBody();}try {String msg=new String(bytes,"UTF-8");System.out.println("q.direct1------------------------"+msg);} catch (UnsupportedEncodingException e) {e.printStackTrace();}}
}@Component
@RabbitListener(queues = "q.direct2")
public class DirectExchangeConsumer2 {@RabbitHandlerpublic void receiveMessage(Object object) {Message message=(Message)object;byte bytes[]=null;if (message != null) {bytes=message.getBody();}try {String msg=new String(bytes,"UTF-8");System.out.println("q.direct2------------------------"+msg);} catch (UnsupportedEncodingException e) {e.printStackTrace();}}
}

三、topic交换器配置和使用
1.声明队列、交换器并将队列绑定到交换器

@Configuration
public class TopicExchangeConfiguration {/*** 动态声明队列* @return*/@Beanpublic Queue topicQueue1() {return new Queue("q.topic1");}/*** 动态声明队列* @return*/@Beanpublic Queue topicQueue2() {return new Queue("q.topic2");}/*** 动态声明交换器* @return*/@Beanpublic TopicExchange topicExchange() {return new TopicExchange("x.topic");}@Beanpublic Binding bindingTopicExchange2Queue1() {return BindingBuilder.bind(topicQueue1()).to(topicExchange()).with("*.*.test");}@Beanpublic Binding bindingTopicExchange2Queue2() {return BindingBuilder.bind(topicQueue2()).to(topicExchange()).with("*.topic.*");}@Beanpublic Binding bindingTopicExchange2Queue3() {return BindingBuilder.bind(topicQueue1()).to(topicExchange()).with("com.#");}
}

2.创建消息生产者

public void sendMessage2TopicMessage(Object object) {try {Message message= MessageBuilder.withBody(JSON.toJSONString(object).getBytes("UTF-8")).setContentType(MessageProperties.CONTENT_TYPE_JSON).setContentEncoding("UTF-8").setMessageId(UUID.randomUUID().toString()).build();// 指定exchange 为x.direct, 路由键为 r.direct.routingKey2amqpTemplate.send("x.topic","com.test",message);} catch (UnsupportedEncodingException e) {e.printStackTrace();}}public void sendMessage2TopicMessage2(Object object) {try {Message message= MessageBuilder.withBody(JSON.toJSONString(object).getBytes("UTF-8")).setContentType(MessageProperties.CONTENT_TYPE_JSON).setContentEncoding("UTF-8").setMessageId(UUID.randomUUID().toString()).build();// 指定exchange 为x.direct, 路由键为 r.direct.routingKey2amqpTemplate.send("x.topic","client.topic.test",message);} catch (UnsupportedEncodingException e) {e.printStackTrace();}}

3.创建消息消费者

@Component
@RabbitListener(queues = "q.topic1")
public class TopicExchangeConsumer1 {@RabbitHandlerpublic void receiveMessage(Object object) {Message message=(Message)object;byte bytes[]=null;if (message != null) {bytes=message.getBody();}try {String msg=new String(bytes,"UTF-8");System.out.println("q.topic1------------------------"+msg);} catch (UnsupportedEncodingException e) {e.printStackTrace();}}
}@Component
@RabbitListener(queues = "q.topic2")
public class TopicExchangeConsumer2 {@RabbitHandlerpublic void receiveMessage(Object object) {Message message=(Message)object;byte bytes[]=null;if (message != null) {bytes=message.getBody();}try {String msg=new String(bytes,"UTF-8");System.out.println("q.topic2------------------------"+msg);} catch (UnsupportedEncodingException e) {e.printStackTrace();}}
}

四、创建消息发送测试controller

@RestController
@RequestMapping(value = "/messages")
public class MessageController {@AutowiredRabbitProducer rabbitProducer;@PostMapping(value = "/fanout")public Map<String,Object> sendMessage() {Map<String,Object> map=new HashMap<>();map.put("username","test1");map.put("password","123456");map.put("name","我是fanout 交换器测试人员");rabbitProducer.sendMessage(map);Map<String,Object> resultMap=new HashMap<>();resultMap.put("code","200");resultMap.put("message","success");return resultMap;}@PostMapping(value = "/direct")public Map<String,Object> sendMessage2DirectExchange() {Map<String,Object> map=new HashMap<>();map.put("username","test2");map.put("password","123456");map.put("name","我是direct 交换器测试人员");rabbitProducer.sendMessage2DirectExchange(map);Map<String,Object> resultMap=new HashMap<>();resultMap.put("code","200");resultMap.put("message","success");return resultMap;}@PostMapping(value = "/topic")public Map<String,Object> sendMessage2TopicExchange() {Map<String,Object> map=new HashMap<>();map.put("username","test3");map.put("password","123456");map.put("name","我是topic 交换器测试人员");rabbitProducer.sendMessage2TopicMessage(map);Map<String,Object> resultMap=new HashMap<>();resultMap.put("code","200");resultMap.put("message","success");return resultMap;}@PostMapping(value = "/topic1")public Map<String,Object> sendMessage2TopicExchange1() {Map<String,Object> map=new HashMap<>();map.put("username","test3");map.put("password","123456");map.put("name","我是topic 交换器测试人员");rabbitProducer.sendMessage2TopicMessage2(map);Map<String,Object> resultMap=new HashMap<>();resultMap.put("code","200");resultMap.put("message","success");return resultMap;}
}

源码下载地址
https://github.com/tangyajun/spring-boot-rabbit-consumer
https://github.com/tangyajun/rabbitmq-spring-demo

这篇关于rabbitmq 学习三-spring-boot中配置交换器和队列的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/809812

相关文章

IDEA运行spring项目时,控制台未出现的解决方案

《IDEA运行spring项目时,控制台未出现的解决方案》文章总结了在使用IDEA运行代码时,控制台未出现的问题和解决方案,问题可能是由于点击图标或重启IDEA后控制台仍未显示,解决方案提供了解决方法... 目录问题分析解决方案总结问题js使用IDEA,点击运行按钮,运行结束,但控制台未出现http://

解决Spring运行时报错:Consider defining a bean of type ‘xxx.xxx.xxx.Xxx‘ in your configuration

《解决Spring运行时报错:Considerdefiningabeanoftype‘xxx.xxx.xxx.Xxx‘inyourconfiguration》该文章主要讲述了在使用S... 目录问题分析解决方案总结问题Description:Parameter 0 of constructor in x

解决IDEA使用springBoot创建项目,lombok标注实体类后编译无报错,但是运行时报错问题

《解决IDEA使用springBoot创建项目,lombok标注实体类后编译无报错,但是运行时报错问题》文章详细描述了在使用lombok的@Data注解标注实体类时遇到编译无误但运行时报错的问题,分析... 目录问题分析问题解决方案步骤一步骤二步骤三总结问题使用lombok注解@Data标注实体类,编译时

JSON字符串转成java的Map对象详细步骤

《JSON字符串转成java的Map对象详细步骤》:本文主要介绍如何将JSON字符串转换为Java对象的步骤,包括定义Element类、使用Jackson库解析JSON和添加依赖,文中通过代码介绍... 目录步骤 1: 定义 Element 类步骤 2: 使用 Jackson 库解析 jsON步骤 3: 添

VScode连接远程Linux服务器环境配置图文教程

《VScode连接远程Linux服务器环境配置图文教程》:本文主要介绍如何安装和配置VSCode,包括安装步骤、环境配置(如汉化包、远程SSH连接)、语言包安装(如C/C++插件)等,文中给出了详... 目录一、安装vscode二、环境配置1.中文汉化包2.安装remote-ssh,用于远程连接2.1安装2

Java中注解与元数据示例详解

《Java中注解与元数据示例详解》Java注解和元数据是编程中重要的概念,用于描述程序元素的属性和用途,:本文主要介绍Java中注解与元数据的相关资料,文中通过代码介绍的非常详细,需要的朋友可以参... 目录一、引言二、元数据的概念2.1 定义2.2 作用三、Java 注解的基础3.1 注解的定义3.2 内

Java中使用Java Mail实现邮件服务功能示例

《Java中使用JavaMail实现邮件服务功能示例》:本文主要介绍Java中使用JavaMail实现邮件服务功能的相关资料,文章还提供了一个发送邮件的示例代码,包括创建参数类、邮件类和执行结... 目录前言一、历史背景二编程、pom依赖三、API说明(一)Session (会话)(二)Message编程客

Java中List转Map的几种具体实现方式和特点

《Java中List转Map的几种具体实现方式和特点》:本文主要介绍几种常用的List转Map的方式,包括使用for循环遍历、Java8StreamAPI、ApacheCommonsCollect... 目录前言1、使用for循环遍历:2、Java8 Stream API:3、Apache Commons

JavaScript中的isTrusted属性及其应用场景详解

《JavaScript中的isTrusted属性及其应用场景详解》在现代Web开发中,JavaScript是构建交互式应用的核心语言,随着前端技术的不断发展,开发者需要处理越来越多的复杂场景,例如事件... 目录引言一、问题背景二、isTrusted 属性的来源与作用1. isTrusted 的定义2. 为

Java循环创建对象内存溢出的解决方法

《Java循环创建对象内存溢出的解决方法》在Java中,如果在循环中不当地创建大量对象而不及时释放内存,很容易导致内存溢出(OutOfMemoryError),所以本文给大家介绍了Java循环创建对象... 目录问题1. 解决方案2. 示例代码2.1 原始版本(可能导致内存溢出)2.2 修改后的版本问题在