diff --git a/pom.xml b/pom.xml
index 89efbd080..25ad5c116 100644
--- a/pom.xml
+++ b/pom.xml
@@ -25,7 +25,7 @@
1.8
2020.0.0-RC2
2.1
- 3.1.0-M5
+ 3.1.0-SNAPSHOT
true
true
true
diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java
index de2e624e1..e1b3ccedb 100644
--- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java
+++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java
@@ -1361,7 +1361,10 @@ public class ImplicitFunctionBindingTests {
@Bean
public Function>>, Message>>> funcB() {
- return v -> v;
+ return v -> {
+ System.out.println(v);
+ return v;
+ };
}
}
diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/MultipleInputOutputFunctionTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/MultipleInputOutputFunctionTests.java
index 4281cdf95..80bff1bfb 100644
--- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/MultipleInputOutputFunctionTests.java
+++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/MultipleInputOutputFunctionTests.java
@@ -151,6 +151,31 @@ public class MultipleInputOutputFunctionTests {
}
}
+ @Test
+ public void testMultiInputMessageSingleOutput() {
+ try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
+ TestChannelBinderConfiguration.getCompleteConfiguration(
+ ReactiveFunctionConfiguration.class))
+ .web(WebApplicationType.NONE)
+ .run("--spring.jmx.enabled=false",
+ "--spring.cloud.function.definition=multiInputSingleOutputMessage")) {
+ context.getBean(InputDestination.class);
+
+ InputDestination inputDestination = context.getBean(InputDestination.class);
+ OutputDestination outputDestination = context.getBean(OutputDestination.class);
+
+ Message stringInputMessage = MessageBuilder.withPayload("one".getBytes()).build();
+ Message integerInputMessage = MessageBuilder.withPayload("1".getBytes()).build();
+ inputDestination.send(stringInputMessage, "multiInputSingleOutputMessage-in-0");
+ inputDestination.send(integerInputMessage, "multiInputSingleOutputMessage-in-1");
+
+ Message outputMessage = outputDestination.receive();
+ assertThat(outputMessage.getPayload()).isEqualTo("one".getBytes());
+ outputMessage = outputDestination.receive(0, "multiInputSingleOutputMessage-out-0");
+ assertThat(outputMessage.getPayload()).isEqualTo("1".getBytes());
+ }
+ }
+
@Test
public void testSingleInputMultiOutput() {
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
@@ -336,6 +361,18 @@ public class MultipleInputOutputFunctionTests {
};
}
+ @Bean
+ public Function>, Flux>>, Flux> multiInputSingleOutputMessage() {
+ return tuple -> {
+ Flux stringStream = tuple.getT1().map(m -> m.getPayload());
+ Flux intStream = tuple.getT2().map(i -> {
+ int v = i.getPayload();
+ return String.valueOf(v);
+ });
+ return Flux.merge(stringStream, intStream);
+ };
+ }
+
@Bean
@SuppressWarnings({ "unchecked", "rawtypes" })
public static Function, Tuple2, Flux>> singleInputMultipleOutputs() {