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:
@@ -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);
|
||||
|
||||
@@ -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 {
|
||||
|
||||
Reference in New Issue
Block a user