SpringBoot整合多个RabbitMQ

一、背景

​ 最近项目中需要用到了RabbitMQ来监听消息队列,监听的消息队列的 虚拟主机(virtualHost)和队列名(queueName)是不一致的,但是接收到的消息格式相同的。而且可能还存在程序不停机的情况下,动态的增加新的队列(queue)的监听,因此就需要我们自己在程序中实现一种方法实现动态配置RabbitMQ

二、需求

我们有2RabbitMQ的配置,在程序启动的时候,动态的配置好这2个RabbitMQ,实现消息的监听。

RabbitMQ的配置信息

host

port

username

password

virtualHost

queueName

47.101.130.164

5672

rabbit-multi-01

rabbit-multi-01

/rabbit-multi-01

queue-rabbit-multi-01

47.101.130.164

5672

rabbit-multi-02

rabbit-multi-02

/rabbit-multi-02

queue-rabbit-multi-02

三、实现思路

1、动态配置RabbitMQ

包括 ConnectionFactory,RabbitAdmin,RabbitTemplate,SimpleMessageListenerContainer

2、将上方配置好的Bean注入到Spring容器中,之后可能需要用到

Spring容器中注入Bean的方法

1DefaultListableBeanFactory#registerSingleton 2 3 4DefaultListableBeanFactory#registerBeanDefinition

四、实现步骤

1、引入maven依赖

1<dependencies> 2 <dependency> 3 <groupId>org.springframework.boot</groupId> 4 <artifactId>spring-boot-starter-amqp</artifactId> 5 </dependency> 6 <dependency> 7 <groupId>org.springframework.boot</groupId> 8 <artifactId>spring-boot-starter-web</artifactId> 9 </dependency> 10 <dependency> 11 <groupId>org.projectlombok</groupId> 12 <artifactId>lombok</artifactId> 13 <optional>true</optional> 14 </dependency> 15 <dependency> 16 <groupId>org.springframework.boot</groupId> 17 <artifactId>spring-boot-starter-test</artifactId> 18 <scope>test</scope> 19 </dependency> 20</dependencies>

2、创建RabbitProperties用来表示RabbitMQ的配置信息

1@Data 2@NoArgsConstructor 3@AllArgsConstructor 4@Builder 5public class RabbitProperties { 6 private String host; 7 private Integer port; 8 private String username; 9 private String password; 10 private String virtualHost; 11 private String queueName; 12}

3、配置RabbitMQ

配置 ConnectionFactory,RabbitAdmin,RabbitTemplate,SimpleMessageListenerContainer等,并动态注入到Spring容器中

1@Configuration 2@RequiredArgsConstructor 3@Slf4j 4public class MultiRabbitMqConfig { 5 6 private final DefaultListableBeanFactory defaultListableBeanFactory; 7 8 private static Map<String, RabbitProperties> multiMqPropertiesMap = new HashMap<String, RabbitProperties>() { 9 { 10 put("first", RabbitProperties.builder() 11 .host("47.101.130.164") 12 .port(5672) 13 .username("rabbit-multi-01") 14 .password("rabbit-multi-01") 15 .virtualHost("/rabbit-multi-01") 16 .queueName("queue-rabbit-multi-01").build()); 17 put("second", RabbitProperties.builder() 18 .host("47.101.130.164") 19 .port(5672) 20 .username("rabbit-multi-02") 21 .password("rabbit-multi-02") 22 .virtualHost("/rabbit-multi-02") 23 .queueName("queue-rabbit-multi-02").build()); 24 } 25 }; 26 27 @PostConstruct 28 public void initRabbitmq() { 29 multiMqPropertiesMap.forEach((key, rabbitProperties) -> { 30 31 AbstractBeanDefinition beanDefinition = BeanDefinitionBuilder.genericBeanDefinition(CachingConnectionFactory.class) 32 .addPropertyValue("cacheMode", CachingConnectionFactory.CacheMode.CHANNEL) 33 .addPropertyValue("host", rabbitProperties.getHost()) 34 .addPropertyValue("port", rabbitProperties.getPort()) 35 .addPropertyValue("username", rabbitProperties.getUsername()) 36 .addPropertyValue("password", rabbitProperties.getPassword()) 37 .addPropertyValue("virtualHost", rabbitProperties.getVirtualHost()) 38 .getBeanDefinition(); 39 String connectionFactoryName = String.format("%s%s", key, "ConnectionFactory"); 40 defaultListableBeanFactory.registerBeanDefinition(connectionFactoryName, beanDefinition); 41 CachingConnectionFactory connectionFactory = defaultListableBeanFactory.getBean(connectionFactoryName, CachingConnectionFactory.class); 42 43 String rabbitAdminName = String.format("%s%s", key, "RabbitAdmin"); 44 AbstractBeanDefinition rabbitAdminBeanDefinition = BeanDefinitionBuilder.genericBeanDefinition(RabbitAdmin.class) 45 .addConstructorArgValue(connectionFactory) 46 .addPropertyValue("autoStartup", true) 47 .getBeanDefinition(); 48 defaultListableBeanFactory.registerBeanDefinition(rabbitAdminName, rabbitAdminBeanDefinition); 49 RabbitAdmin rabbitAdmin = defaultListableBeanFactory.getBean(rabbitAdminName, RabbitAdmin.class); 50 log.info("rabbitAdmin:[{}]", rabbitAdmin); 51 52 RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory); 53 defaultListableBeanFactory.registerSingleton(String.format("%s%s", key, "RabbitTemplate"), rabbitTemplate); 54 55 SimpleMessageListenerContainer simpleMessageListenerContainer = new SimpleMessageListenerContainer(connectionFactory); 56 // 设置监听的队列 57 simpleMessageListenerContainer.setQueueNames(rabbitProperties.getQueueName()); 58 // 指定要创建的并发使用者的数量,默认值是1,当并发高时可以增加这个的数值,同时下方max的数值也要增加 59 simpleMessageListenerContainer.setConcurrentConsumers(3); 60 // 最大的并发消费者 61 simpleMessageListenerContainer.setMaxConcurrentConsumers(10); 62 // 设置是否重回队列 63 simpleMessageListenerContainer.setDefaultRequeueRejected(false); 64 // 设置签收模式 65 simpleMessageListenerContainer.setAcknowledgeMode(AcknowledgeMode.MANUAL); 66 // 设置非独占模式 67 simpleMessageListenerContainer.setExclusive(false); 68 // 设置consumer未被 ack 的消息个数 69 simpleMessageListenerContainer.setPrefetchCount(1); 70 // 设置消息监听器 71 simpleMessageListenerContainer.setMessageListener((ChannelAwareMessageListener) (message, channel) -> { 72 try { 73 log.info("============> Thread:[{}] 接收到消息:[{}] ", Thread.currentThread().getName(), new String(message.getBody())); 74 log.info("====>connection:[{}]", channel.getConnection()); 75 channel.basicAck(message.getMessageProperties().getDeliveryTag(), false); 76 } catch (Exception e) { 77 log.error(e.getMessage(), e); 78 // 发生异常此处需要捕获到 79 channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, true); 80 } 81 }); 82 defaultListableBeanFactory.registerSingleton(String.format("%s%s", key, "SimpleMessageListenerContainer"), simpleMessageListenerContainer); 83 }); 84 new Thread(() -> { 85 try { 86 TimeUnit.SECONDS.sleep(3); 87 } catch (InterruptedException e) { 88 e.printStackTrace(); 89 } 90 RabbitTemplate firstRabbitTemplate = (RabbitTemplate) defaultListableBeanFactory.getBean("firstRabbitTemplate"); 91 firstRabbitTemplate.convertAndSend("exchange-rabbit-multi-01", "", "first queue message"); 92 log.info("over..."); 93 }).start(); 94 } 95}

五、实现效果

Multi-RabbitMQ

六、代码

https://gitee.com/huan1993/rabbitmq/tree/master/rabbitmq-springboot-multi

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

2020年前端实用代码段,为你的工作保驾护航

有空的时候,自己总结了几个代码段,在开发中也经常使用,谢谢。1、使用解构获取json数据let jsonData  id: 1,status: "OK",data: 'a', 'b';let  id, status, data: number   jsonData;console.log(id, status, number )