You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

启用多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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.26 00:57:13