RabbitMQ streams with Spring Boot

Viewed 534

Im trying to utilise RabbitMQ streams (as of 3.9+) with Spring boot's org.springframework.boot:spring-boot-starter-amqp starter.

First I started the RabbitMQ docker compose by

  rabbitmq:
    image: rabbitmq:management
    ports:
      - "5672:5672"
      - "15672:15672"

and then run the bash script to enable stream plugin:

docker exec docker_rabbitmq_1 rabbitmq-plugins enable rabbitmq_stream

By declaring the queue

    @Bean
    Queue queue2() {
        Map<String, Object> args = new HashMap<>();
        args.put("x-queue-type", "stream");
        return new Queue("test3-queue",true, false, false, args);
    }

I configure the queue as stream and the broker says in logs

2021-09-22 21:06:46.400520+00:00 [warn] <0.1090.0> rabbit_stream_coordinator: started writer __test3-queue_1632344806396284388 on rabbit@1115e1f022e2 in 1

This should be enough config to start using Rabbit's new feature Streams with the given queue.

When I send message to the rabbitTemplate

template.convertAndSend("test3-stream", request.getMessage());

all my listeners

    @RabbitListener(id = "listener1", queues = "#{queue2}")
    public void listen1(String in) {
        log.info("AMQP listener 1: {}", in);
    }

    @RabbitListener(id = "listener2", queues = "#{queue2}")
    public void listen2(String in) {
        log.info("AMQP listener 2: {}", in);
    }

    @RabbitListener(id = "listener3", queues = "#{queue2}")
    public void listen3(String in) {
        log.info("AMQP listener 3: {}", in);
    }

receive the message and print its log. According the doc, referencing the queue by SpEL receives the queue configured with x-queue-type: stream.

As SB uses the amqp starter with the amqp client 0.9.1, it should work, but id does not:

But when I add new listener, old messages are not being processed with it.

Is it even possible to use the @RabbitListener with the new Streams feature or Im too early with using the append-only log with Rabbit broker? Should I use Kafka instead just because Spring has not yet implemented support for RabbitMQ streams?

Rabbit comes with its own Java library handling the stream events, which works well but it is missing the simplicity of the Spring underlying heavy lifting..

EDIT: I've got a bit further by configuring the RabbitListenerContainerFactory with customizer:

    @Bean
    RabbitListenerContainerFactory rabbitListenerContainerFactory(
            SimpleRabbitListenerContainerFactoryConfigurer configurer, ConnectionFactory connectionFactory) {
        SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory();
        configurer.configure(factory, connectionFactory);
        factory.setContainerCustomizer(c -> {
            Map<String, Object> args = new HashMap<>(c.getConsumerArguments());
            args.put("x-stream-offset", "first");
            c.setConsumerArguments(args);
        });
        return factory;
    }

Which causes the new listener to read the messages from the beginning (offset 0). The odd is that ALL listeners now read from offset 0. I guess its because the consumers (listeners) do not have the name set correctly?

0 Answers
Related