diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java index 1e84264e7..e33403a04 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java @@ -363,6 +363,13 @@ public class FunctionConfiguration { if (isReactiveOrMultipleInputOutput(bindableProxyFactory, functionType)) { Publisher[] inputPublishers = inputBindingNames.stream().map(inputBindingName -> { + BindingProperties bindingProperties = this.serviceProperties.getBindings().get(inputBindingName); + ConsumerProperties consumerProperties = bindingProperties == null ? null : bindingProperties.getConsumer(); + if (consumerProperties != null) { + Assert.isTrue(consumerProperties.getConcurrency() <= 1, "Concurrency > 1 is not supported by reactive " + + "consumer, given that project reactor maintains its own concurrency mechanism. Was '..." + + inputBindingName + ".consumer.concurrency=" + consumerProperties.getConcurrency() + "'"); + } SubscribableChannel inputChannel = this.applicationContext.getBean(inputBindingName, SubscribableChannel.class); return IntegrationReactiveUtils.messageChannelToFlux(inputChannel); }).toArray(Publisher[]::new); diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java index 0481926a6..786a14518 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java @@ -32,6 +32,7 @@ import org.junit.Test; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; +import org.springframework.beans.factory.BeanCreationException; import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.boot.builder.SpringApplicationBuilder; @@ -758,6 +759,22 @@ public class ImplicitFunctionBindingTests { } } + @Test(expected = BeanCreationException.class) + public void testReactiveConsumerWithConcurrencyFailureConfiguration() { + System.clearProperty("spring.cloud.function.definition"); + new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration(ReactiveConsumerWithConcurrencyFailureConfiguration.class)) + .web(WebApplicationType.NONE).run("--spring.jmx.enabled=false", + "--spring.cloud.stream.bindings.input-in-0.consumer.concurrency=2"); + } + + @EnableAutoConfiguration + public static class ReactiveConsumerWithConcurrencyFailureConfiguration { + @Bean + public Consumer>> input() { + return flux -> flux.subscribe(System.out::println); + } + } @EnableAutoConfiguration public static class NoEnableBindingConfiguration {