GH-1784 Fix support for Consumer-type handlers

Resolves #1784
This commit is contained in:
Oleg Zhurakousky
2019-08-19 19:01:38 +02:00
parent 99581a6158
commit 433bea8f26
4 changed files with 32 additions and 64 deletions

View File

@@ -111,7 +111,7 @@
-->
<module>spring-cloud-stream-schema</module>
<module>spring-cloud-stream-schema-server</module>
<module>docs</module>
</modules>
<build>

View File

@@ -218,15 +218,8 @@ public class FunctionConfiguration {
}
else {
FunctionInvocationWrapper function = functionCatalog.lookup(functionProperties.getDefinition(), "application/json");
if (!function.isSupplier()) {
if (function.isConsumer()) {
throw new UnsupportedOperationException("Consumers are not currently supported");
}
else if (function.isFunction()) {
if ("input".equals(channelName)) {
this.postProcessForStandAloneFunction(function, messageChannel);
}
}
if (!function.isSupplier() && "input".equals(channelName)) {
this.postProcessForStandAloneFunction(function, messageChannel);
}
}
}
@@ -243,7 +236,9 @@ public class FunctionConfiguration {
ServiceActivatingHandler handler = new ServiceActivatingHandler(new FunctionWrapper(function));
handler.setBeanFactory(context);
handler.afterPropertiesSet();
handler.setOutputChannelName("output");
if (!FunctionTypeUtils.isConsumer(functionType)) {
handler.setOutputChannelName("output");
}
SubscribableChannel subscribeChannel = (SubscribableChannel) inputChannel;
subscribeChannel.subscribe(handler);
}

View File

@@ -67,7 +67,6 @@ import static org.assertj.core.api.Assertions.assertThat;
public class GreenfieldFunctionEnableBindingTests {
@Test
@Ignore
public void testSourceFromSupplier() {
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
TestChannelBinderConfiguration
@@ -108,7 +107,6 @@ public class GreenfieldFunctionEnableBindingTests {
}
@Test
@Ignore
public void testSinkFromConsumer() {
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
TestChannelBinderConfiguration

View File

@@ -26,17 +26,14 @@ import reactor.core.publisher.Flux;
import org.springframework.boot.WebApplicationType;
import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
import org.springframework.boot.builder.SpringApplicationBuilder;
import org.springframework.cloud.stream.annotation.EnableBinding;
import org.springframework.cloud.stream.binder.test.InputDestination;
import org.springframework.cloud.stream.binder.test.OutputDestination;
import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration;
import org.springframework.cloud.stream.messaging.Processor;
import org.springframework.context.ConfigurableApplicationContext;
import org.springframework.context.annotation.Bean;
import org.springframework.integration.support.MessageBuilder;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.GenericMessage;
import static org.assertj.core.api.Assertions.assertThat;
@@ -54,7 +51,7 @@ public class ImplicitFunctionBindingTests {
}
@Test
public void testBindingWithNoEnableBindingConfiguration() {
public void testSimpleFunctionWithStreamProperty() {
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
TestChannelBinderConfiguration.getCompleteConfiguration(
@@ -78,7 +75,7 @@ public class ImplicitFunctionBindingTests {
}
@Test
public void testBindingWithNoEnableBindingConfigurationWithFunctionNativeDefinitionProperty() {
public void testSimpleFunctionWithNativeProperty() {
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
TestChannelBinderConfiguration.getCompleteConfiguration(
@@ -102,31 +99,7 @@ public class ImplicitFunctionBindingTests {
}
@Test
public void testBindingWithEnableBindingConfiguration() {
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
TestChannelBinderConfiguration.getCompleteConfiguration(
EnableBindingConfiguration.class))
.web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false",
"--spring.cloud.stream.function.definition=func")) {
InputDestination inputDestination = context.getBean(InputDestination.class);
OutputDestination outputDestination = context
.getBean(OutputDestination.class);
Message<byte[]> inputMessage = MessageBuilder
.withPayload("Hello".getBytes()).build();
inputDestination.send(inputMessage);
Message<byte[]> outputMessage = outputDestination.receive();
assertThat(outputMessage.getPayload()).isEqualTo("Hello".getBytes());
}
}
@Test
public void testBindingWithNoEnableBindingAndNoDefinitionPropertyConfiguration() {
public void testSimpleFunctionWithoutDefinitionProperty() {
System.clearProperty("spring.cloud.stream.function.definition");
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
TestChannelBinderConfiguration.getCompleteConfiguration(
@@ -148,6 +121,20 @@ public class ImplicitFunctionBindingTests {
}
}
@Test
public void testConsumer() {
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
TestChannelBinderConfiguration
.getCompleteConfiguration(SingleFunctionConfiguration.class))
.web(WebApplicationType.NONE)
.run("--spring.cloud.stream.function.definition=consumer",
"--spring.jmx.enabled=false")) {
InputDestination source = context.getBean(InputDestination.class);
source.send(new GenericMessage<byte[]>("John Doe".getBytes()));
}
}
@Test
public void testBindingWithReactiveFunction() {
System.clearProperty("spring.cloud.stream.function.definition");
@@ -195,26 +182,6 @@ public class ImplicitFunctionBindingTests {
}
}
@EnableAutoConfiguration
@EnableBinding(Processor.class)
public static class EnableBindingConfiguration {
@Bean
public Function<String, String> func() {
return x -> {
System.out.println("Function");
return x;
};
}
@Bean
public Consumer<String> cons() {
return x -> {
System.out.println("Consumer");
};
}
}
@EnableAutoConfiguration
public static class SingleFunctionConfiguration {
@@ -225,6 +192,14 @@ public class ImplicitFunctionBindingTests {
return x;
};
}
@Bean
public Consumer<String> consumer() {
return value -> {
System.out.println(value);
};
}
}
@EnableAutoConfiguration