From 8bba5746903f092b82f5a96b39f901bb9520b683 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 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 {