From 1b3fc7074bd8b1de8ca2568e89b9476109a96c05 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Mon, 15 May 2023 16:09:45 -0400 Subject: [PATCH] Reactive Kafka Binder errors when concurrency > 1 (#2734) * Reactive Kafka Binder errors when concurrency > 1 When using Reactive Kafka binder, it is allowed to have concurrency > 1. There is a check in FunctionConfiguration that throws an error if concurrency is > 1, when using reactive types. Since it is allowed to do so with Reative Kafka binder, switch this conversion into a warning log message. Resolves https://github.com/spring-cloud/spring-cloud-stream/issues/2726 * Update core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java Co-authored-by: Gary Russell * Update core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java Co-authored-by: Gary Russell * Update core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java Co-authored-by: Gary Russell --------- Co-authored-by: Gary Russell --- .../function/ImplicitFunctionBindingTests.java | 17 ++++++++--------- .../stream/function/FunctionConfiguration.java | 8 +++++--- 2 files changed, 13 insertions(+), 12 deletions(-) diff --git a/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java b/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java index 409b403de..ddae8e574 100644 --- a/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java +++ b/core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java @@ -33,6 +33,7 @@ import java.util.function.Supplier; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.Disabled; import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.ValueSource; import reactor.core.publisher.Flux; @@ -40,10 +41,11 @@ import reactor.core.publisher.Mono; import reactor.core.publisher.Sinks; import reactor.core.publisher.Sinks.Many; -import org.springframework.beans.factory.BeanCreationException; import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.boot.builder.SpringApplicationBuilder; +import org.springframework.boot.test.system.CapturedOutput; +import org.springframework.boot.test.system.OutputCaptureExtension; import org.springframework.cloud.function.context.FunctionCatalog; import org.springframework.cloud.function.context.FunctionRegistration; import org.springframework.cloud.function.context.catalog.FunctionAroundWrapper; @@ -961,17 +963,14 @@ public class ImplicitFunctionBindingTests { } @Test - void testReactiveConsumerWithConcurrencyFailureConfiguration() { + @ExtendWith(OutputCaptureExtension.class) + void testReactiveConsumerWithConcurrencyGreaterThanOneLogsWarning(CapturedOutput output) { System.clearProperty("spring.cloud.function.definition"); - try { - new SpringApplicationBuilder( + try (ConfigurableApplicationContext ignored = new SpringApplicationBuilder( TestChannelBinderConfiguration.getCompleteConfiguration(ReactiveConsumerWithConcurrencyFailureConfiguration.class)) .web(WebApplicationType.NONE).run("--spring.jmx.enabled=false", - "--spring.cloud.stream.bindings.input-in-0.consumer.concurrency=2"); - fail(); - } - catch (BeanCreationException e) { - // good + "--spring.cloud.stream.bindings.input-in-0.consumer.concurrency=2")) { + assertThat(output).contains("When using concurrency > 1 in reactive contexts, please make sure that you are using a reactive binder"); } } diff --git a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java index cf3c83430..ff67b23c9 100644 --- a/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java +++ b/core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java @@ -517,9 +517,11 @@ public class FunctionConfiguration { ConsumerProperties consumerProperties = bindingProperties == null ? null : bindingProperties.getConsumer(); if (consumerProperties != null) { function.setSkipInputConversion(consumerProperties.isUseNativeDecoding()); - 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() + "'"); + if (consumerProperties.getConcurrency() > 1) { + this.logger.warn("When using concurrency > 1 in reactive contexts, please make sure that you are using a " + + "reactive binder that supports concurrency settings. Otherwise, concurrency settings > 1 will be ignored when " + + "using reactive types."); + } } MessageChannel inputChannel = this.applicationContext.getBean(inputBindingName, MessageChannel.class); return IntegrationReactiveUtils.messageChannelToFlux(inputChannel).map(m -> {