启用多Vhost动态支持后默认Vhost的RabbitMQ监听器无法消费消息
问题:Spring Boot RabbitMQ 单队列监听多Vhost后默认Vhost监听器无法消费消息
我的应用基于Spring Boot + RabbitMQ,需求是单队列监听多个Vhost。但启用该功能后,默认虚拟主机(/)的监听器无法消费消息。
现有配置代码
public class RabbitMqConfig implements RabbitListenerConfigurer { @Bean public RabbitMessagingTemplate rabbitMessagingTemplate( RabbitTemplate rabbitTemplate, ConnectionFactory connectionFactory) { RabbitMessagingTemplate rabbitMessagingTemplate = new RabbitMessagingTemplate(); SimpleRoutingConnectionFactory simpleRoutingConnectionFactory = new SimpleRoutingConnectionFactory(); simpleRoutingConnectionFactory.setDefaultTargetConnectionFactory(connectionFactory); rabbitTemplate.setConnectionFactory(simpleRoutingConnectionFactory); rabbitMessagingTemplate.setRabbitTemplate(rabbitTemplate); return rabbitMessagingTemplate; } @Override public void configureRabbitListeners(RabbitListenerEndpointRegistrar registrar) { clients .stream() .forEach(client -> { var endpoint = createEndpoint(client.getVhost()); var listenerContainerFactory = simpleRabbitListenerContainerFactory(client.getVhost()); listenerContainerFactory.setConnectionFactory( annotationConfigWebApplicationContext.getBean((client.getVhost()), ConnectionFactory.class)); registrar.registerEndpoint(endpoint, listenerContainerFactory); }); } public SimpleRabbitListenerContainerFactory simpleRabbitListenerContainerFactory(String vHost) { SimpleRabbitListenerContainerFactory simpleRabbitListenerContainerFactory = new SimpleRabbitListenerContainerFactory(); simpleRabbitListenerContainerFactory.setAcknowledgeMode(AcknowledgeMode.AUTO); CachingConnectionFactory cf = new CachingConnectionFactory(); cf.setUsername(rabbitProperties.getUsername()); cf.setPassword(rabbitProperties.getPassword()); cf.setVirtualHost(vHost); cf.setHost(rabbitProperties.getHost()); cf.setBeanName(vHost); simpleRabbitListenerContainerFactory.setConnectionFactory(cf); annotationConfigWebApplicationContext.getBeanFactory().registerSingleton(vHost, cf); simpleRabbitListenerContainerFactory.setContainerCustomizer(simpleMessageListenerContainer -> { simpleMessageListenerContainer.setQueueNames(graphqlMutationQueue); simpleMessageListenerContainer.setLookupKeyQualifier(graphqlMutationQueue); simpleMessageListenerContainer.setMessageListener(msg -> { log.info(msg + " from VH: " + msg.getMessageProperties().getHeader(REPLY_TO_VHOST)); updateGraphqlListener.onMessage(msg); }); simpleMessageListenerContainer.setAfterReceivePostProcessors(msg -> { msg.getMessageProperties().setHeader(REPLY_TO_VHOST, vHost); return msg; }); simpleMessageListenerContainer.start(); }); return simpleRabbitListenerContainerFactory; } private SimpleRabbitListenerEndpoint createEndpoint(String vHost) { var endpoint = new SimpleRabbitListenerEndpoint(); endpoint.setId(vHost); endpoint.setQueueNames(graphqlMutationQueue); endpoint.setMessageConverter(jsonMessageConverter()); endpoint.setMessageListener(new MessageListenerAdapter("OnMessage")); return endpoint; } }
消息发送代码
应用启动时,通过以下代码向交换机发送消息:
rabbitTemplate.convertAndSend(graphQLExchangeName, StringUtils.EMPTY, Collections.EMPTY_MAP);
执行该代码后,默认虚拟主机的监听器无法消费消息。
补充配置(编辑2)
@Bean RabbitTemplate rabbitTemplate(ConnectionFactory defaultConnectionFactory) { return new RabbitTemplate(defaultConnectionFactory); } @Bean RabbitAdmin rabbitAdminAdmin(ConnectionFactory defaultConnectionFactory) { return new RabbitAdmin(defaultConnectionFactory); } @Bean SimpleRabbitListenerContainerFactory defaultContFactory(ConnectionFactory defaultConnectionFactory) { SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory(); factory.setConnectionFactory(defaultConnectionFactory); return factory; } @RabbitListener(queues = "#{someQueue}", containerFactory = "defaultContFactory") public void onMessage(Map<?, ?> map) { } @SpringBootApplication(exclude = RabbitAutoConfiguration.class)
问题分析与解决方案
核心问题点
- 默认Vhost未被纳入监听注册逻辑:
configureRabbitListeners仅遍历clients列表创建监听器,默认Vhost(/)不在列表中,导致没有对应的监听端点。 - 连接工厂路由规则缺失:配置了
SimpleRoutingConnectionFactory但未设置路由映射,发送消息时无法匹配到默认Vhost的连接工厂。 - 容器生命周期管理错误:手动调用
simpleMessageListenerContainer.start(),跳过了Spring的生命周期管理,可能导致容器初始化异常。
修复步骤
1. 为默认Vhost注册监听端点
修改configureRabbitListeners,先处理默认Vhost的监听注册:
@Override public void configureRabbitListeners(RabbitListenerEndpointRegistrar registrar) { // 注册默认Vhost的监听器 String defaultVhost = "/"; var defaultEndpoint = createEndpoint(defaultVhost); var defaultListenerFactory = simpleRabbitListenerContainerFactory(defaultVhost); registrar.registerEndpoint(defaultEndpoint, defaultListenerFactory); // 处理其他客户Vhost的监听器 clients.stream().forEach(client -> { var endpoint = createEndpoint(client.getVhost()); var listenerContainerFactory = simpleRabbitListenerContainerFactory(client.getVhost()); registrar.registerEndpoint(endpoint, listenerContainerFactory); }); }
2. 完善SimpleRoutingConnectionFactory路由配置
添加Vhost与连接工厂的映射,并配置路由解析规则:
@Bean public RabbitMessagingTemplate rabbitMessagingTemplate( RabbitTemplate rabbitTemplate, ConnectionFactory defaultConnectionFactory) { RabbitMessagingTemplate rabbitMessagingTemplate = new RabbitMessagingTemplate(); SimpleRoutingConnectionFactory simpleRoutingConnectionFactory = new SimpleRoutingConnectionFactory(); simpleRoutingConnectionFactory.setDefaultTargetConnectionFactory(defaultConnectionFactory); // 配置所有Vhost的连接工厂映射 Map<Object, ConnectionFactory> targetFactories = new HashMap<>(); targetFactories.put("/", defaultConnectionFactory); clients.forEach(client -> { String vhost = client.getVhost(); targetFactories.put(vhost, annotationConfigWebApplicationContext.getBean(vhost, ConnectionFactory.class)); }); simpleRoutingConnectionFactory.setTargetConnectionFactories(targetFactories); // 设置路由解析器,根据消息头选择对应Vhost的连接工厂 rabbitTemplate.setRoutingKeyResolver(message -> { String vhost = message.getMessageProperties().getHeader(REPLY_TO_VHOST); return vhost != null ? vhost : "/"; }); rabbitTemplate.setConnectionFactory(simpleRoutingConnectionFactory); rabbitMessagingTemplate.setRabbitTemplate(rabbitTemplate); return rabbitMessagingTemplate; }
3. 移除手动启动容器的代码
删除containerCustomizer中的simpleMessageListenerContainer.start(),交给Spring自动管理容器生命周期:
simpleRabbitListenerContainerFactory.setContainerCustomizer(simpleMessageListenerContainer -> { simpleMessageListenerContainer.setQueueNames(graphqlMutationQueue); simpleMessageListenerContainer.setLookupKeyQualifier(graphqlMutationQueue); simpleMessageListenerContainer.setMessageListener(msg -> { log.info(msg + " from VH: " + msg.getMessageProperties().getHeader(REPLY_TO_VHOST)); updateGraphqlListener.onMessage(msg); }); simpleMessageListenerContainer.setAfterReceivePostProcessors(msg -> { msg.getMessageProperties().setHeader(REPLY_TO_VHOST, vHost); return msg; }); // 移除该行代码 // simpleMessageListenerContainer.start(); });
4. 验证默认Vhost的交换机与队列绑定
检查默认Vhost中graphQLExchangeName交换机和graphqlMutationQueue队列是否已正确绑定,确保消息能正常路由到队列。
内容的提问来源于stack exchange,提问作者Saurabh Singh
相关产品推荐
相关产品推荐

