Concurrency > 1 is not supported by reactive consumer, given that project reactor maintains its own concurrency mechanism

Viewed 272

I'm migrating to the new spring cloud stream version. Spring Cloud Stream 3.1.0

And I have the following consumer configuration:

someChannel-in-0:
  destination: Response1
  group: response-channel
  consumer:
    concurrency: 4

 someChannel-out-0:
   destination: response2

I have connected this channel to the new function binders

// just sample dummy code
@Bean
    public Function<Flux<String>, Flux<String>> aggregate() {
        return inbound -> inbound.
                .map(w -> w.toLowerCase());
    }

And when I'm starting the app I'm getting the following error:

Concurrency > 1 is not supported by reactive consumer, given that project reactor maintains its own concurrency mechanism.

My question is what is the equivalent of concurrency: 4 in project reactor and how do I implement this ?

1 Answers

Basically, Spring cloud stream manage consumers through MessageListenerContainer and it provides a hook which allow users create a bean and inject some advanced configurations. And so here comes to the solution if you are using RabbitMQ as messaging middleware.

@Bean
public ListenerContainerCustomizer<AbstractMessageListenerContainer> listenerContainerCustomizer() {
    return (container, destinationName, group) -> {
        if (container instanceof SimpleMessageListenerContainer) {
            ((SimpleMessageListenerContainer) container).setConcurrency("3");
        }
    };
}
Related