From 5867636d17123876da3263691fb064af67f7694b Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Thu, 4 Jun 2020 11:28:27 +0200 Subject: [PATCH] GH-1977 Add failure when consumer concurrency is > 1 for reactive consumers Providing that project reactor maintains it's own concurrency mechanisms, using stream's concurrency setting can lead to unexpected results including loss of messages. This fixes it by disallowing concurrency setting when reactive consumers are used. Resolves #1977 --- .../stream/function/FunctionConfiguration.java | 7 +++++++ .../function/ImplicitFunctionBindingTests.java | 17 +++++++++++++++++ 2 files changed, 24 insertions(+) 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 cc78e3c7c..b46a43234 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 {