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 <grussell@vmware.com> * Update core/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java Co-authored-by: Gary Russell <grussell@vmware.com> * Update core/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java Co-authored-by: Gary Russell <grussell@vmware.com> --------- Co-authored-by: Gary Russell <grussell@vmware.com>
This commit is contained in:
@@ -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");
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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 -> {
|
||||
|
||||
Reference in New Issue
Block a user