一、背景
最近项目中需要用到了RabbitMQ来监听消息队列,监听的消息队列的 虚拟主机(virtualHost)和队列名(queueName)是不一致的,但是接收到的消息格式相同的。而且可能还存在程序不停机的情况下,动态的增加新的队列(queue)的监听,因此就需要我们自己在程序中实现一种方法实现动态配置RabbitMQ。
二、需求
我们有2个RabbitMQ的配置,在程序启动的时候,动态的配置好这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}
五、实现效果
六、代码
https://gitee.com/huan1993/rabbitmq/tree/master/rabbitmq-springboot-multi