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
This commit is contained in:
Oleg Zhurakousky
2020-06-04 11:28:27 +02:00
parent a8f8a1659e
commit 5867636d17
2 changed files with 24 additions and 0 deletions

View File

@@ -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);

View File

@@ -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<Flux<Message<String>>> input() {
return flux -> flux.subscribe(System.out::println);
}
}
@EnableAutoConfiguration
public static class NoEnableBindingConfiguration {