Listen to RabbitMQ Queues at same time in different Virtual Host using SpringBoot

Viewed 72

I have multiple Clients each client has virtual host defined.

I need to provide ability to listen queues which are located in different virtual host.

SpringBoot , RabbitMQ I am using in the application.

1 Answers

Spring Boot can only auto configure one connection factory.

You will need to configure two or more sets of infrastructure beans (connection factory, container factory, template, etc).

You have to configure both because when Boot detects one, it disables its auto configuration.

EDIT

@SpringBootApplication(exclude = RabbitAutoConfiguration.class)
@EnableRabbit
public class So72953705Application {

    public static void main(String[] args) {
        SpringApplication.run(So72953705Application.class, args);
    }

    @Bean
    ConnectionFactory fooConn() {
        CachingConnectionFactory ccf = new CachingConnectionFactory("localhost", 5672);
        ccf.setVirtualHost("foo");
        return ccf;
    }

    @Bean
    RabbitTemplate fooTemplate(ConnectionFactory fooConn) {
        return new RabbitTemplate(fooConn);
    }

    @Bean
    RabbitAdmin fooAdmin(ConnectionFactory fooConn) {
        return new RabbitAdmin(fooConn);
    }

    @Bean
    SimpleRabbitListenerContainerFactory fooContFactory(ConnectionFactory fooConn) {
        SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory();
        factory.setConnectionFactory(fooConn);
        return factory;
    }

    @Bean
    ConnectionFactory barConn() {
        CachingConnectionFactory ccf = new CachingConnectionFactory("localhost", 5672);
        ccf.setVirtualHost("bar");
        return ccf;
    }

    @Bean
    RabbitTemplate barTemplate(ConnectionFactory barConn) {
        return new RabbitTemplate(barConn);
    }

    @Bean
    RabbitAdmin barAdmin(ConnectionFactory barConn) {
        return new RabbitAdmin(barConn);
    }

    @Bean
    SimpleRabbitListenerContainerFactory barContFactory(ConnectionFactory barConn) {
        SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory();
        factory.setConnectionFactory(barConn);
        return factory;
    }

}

@Component
class FooListener {

    @RabbitListener(queues = "fooQ", containerFactory = "fooContFactory")
    public void listen(String in) {
        System.out.println(in);
    }

}


@Component
class BarListener {

    @RabbitListener(queues = "barQ", containerFactory = "barContFactory")
    public void listen(String in) {
        System.out.println(in);
    }

}

EDIT2

Dynamic registration of infrastructure beans.

@SpringBootApplication(exclude = RabbitAutoConfiguration.class)
@EnableRabbit
public class So72953705Application {

    public static void main(String[] args) {
        SpringApplication.run(So72953705Application.class, args);
    }

    @Bean
    ApplicationRunner runner(DynamicRegistrar registrar) {
        return args -> {
            registrar.registerInfrastructure("foo");
            registrar.registerInfrastructure("bar");
        };
    }

}

@Component
class DynamicRegistrar {

    private final ConfigurableListableBeanFactory beanFactory;

    DynamicRegistrar(ConfigurableListableBeanFactory beanFactory) {
        this.beanFactory = beanFactory;
    }

    void registerInfrastructure(String virtualHost) {
        CachingConnectionFactory ccf = new CachingConnectionFactory("localhost", 5672);
        ccf.setVirtualHost(virtualHost);
        this.beanFactory.registerSingleton(virtualHost + ".cf", ccf);
        this.beanFactory.initializeBean(ccf, virtualHost + ".cf");

        RabbitTemplate template = new RabbitTemplate(ccf);
        this.beanFactory.registerSingleton(virtualHost + ".template", template);
        this.beanFactory.initializeBean(template, virtualHost + ".template");

        RabbitAdmin admin = new RabbitAdmin(ccf);
        this.beanFactory.registerSingleton(virtualHost + ".admin", admin);
        this.beanFactory.initializeBean(admin, virtualHost + ".admin");

        SimpleMessageListenerContainer container = new SimpleMessageListenerContainer(ccf);
        container.setQueueNames("same.queue.name");
        container.setMessageListener(msg -> {
            System.out.println(msg + " from VH: " + msg.getMessageProperties().getHeader("fromVirtualHost"));
        });
        container.setAfterReceivePostProcessors(msg -> {
            msg.getMessageProperties().setHeader("fromVirtualHost", virtualHost);
            return msg;
        });
        this.beanFactory.registerSingleton(virtualHost + ".container", container);
        this.beanFactory.initializeBean(container, virtualHost + ".container");
        container.start();
    }

}
Related