GH-1617 added support for Flux<Message> functions to FunctionInvoker

Resolves #1617
This commit is contained in:
Oleg Zhurakousky
2019-02-19 14:49:10 +01:00
parent 7e0f42ed34
commit 8edbe6ef0c
2 changed files with 120 additions and 6 deletions

View File

@@ -169,7 +169,8 @@ class FunctionInvoker<I, O> implements Function<Flux<Message<I>>, Flux<Message<O
}
else {
returnMessage = (Message<O>) (value instanceof Message ? value
: this.messageConverter.toMessage(value, originalMessage.getHeaders()));
: this.messageConverter.toMessage(value,
originalMessage.getHeaders()));
if (returnMessage == null
&& value.getClass().isAssignableFrom(this.outputClass)) {
returnMessage = wrapOutputToMessage(value, originalMessage);
@@ -212,7 +213,11 @@ class FunctionInvoker<I, O> implements Function<Flux<Message<I>>, Flux<Message<O
? this.messageConverter.fromMessage(message, this.inputClass) : message);
Assert.notNull(argument, "Failed to resolve argument type '" + this.inputClass
+ "' from message: " + message);
if (!this.isInputArgumentMessage && argument instanceof Message) {
if (this.isInputArgumentMessage && !(argument instanceof Message)) {
argument = (T) MessageBuilder.withPayload(argument)
.copyHeaders(message.getHeaders()).build();
}
else if (!this.isInputArgumentMessage && argument instanceof Message) {
argument = ((Message<T>) argument).getPayload();
}
return argument;

View File

@@ -49,7 +49,6 @@ import org.springframework.util.ReflectionUtils;
import static org.assertj.core.api.Assertions.assertThat;
/**
* @author Oleg Zhurakousky
* @author Tolga Kavukcu
@@ -63,7 +62,30 @@ public class FunctionInvokerTests {
public void testSimpleEchoConfiguration() {
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
TestChannelBinderConfiguration.getCompleteConfiguration(
SimpleEchoConfiguration.class))
SimpleEchoConfiguration.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("{\"name\":\"bob\"}".getBytes()).build();
inputDestination.send(inputMessage);
Message<byte[]> outputMessage = outputDestination.receive();
assertThat(outputMessage.getPayload())
.isEqualTo("{\"name\":\"bob\"}".getBytes());
}
}
@Test
public void testFluxPojoFunction() {
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
TestChannelBinderConfiguration
.getCompleteConfiguration(SimpleFluxFunctionConfiguration.class))
.web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false",
"--spring.cloud.stream.function.definition=func")) {
@@ -77,12 +99,33 @@ public class FunctionInvokerTests {
inputDestination.send(inputMessage);
Message<byte[]> outputMessage = outputDestination.receive();
System.out.println("Received: " + new String(outputMessage.getPayload()));
assertThat(outputMessage.getPayload()).isEqualTo("{\"name\":\"bob\"}".getBytes());
assertThat(outputMessage.getPayload()).isEqualTo("Person: bob".getBytes());
}
}
@Test
public void testFluxMessagePojoFunction() {
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
TestChannelBinderConfiguration.getCompleteConfiguration(
SimpleFluxMessageFunctionConfiguration.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("{\"name\":\"bob\"}".getBytes()).build();
inputDestination.send(inputMessage);
Message<byte[]> outputMessage = outputDestination.receive();
assertThat(outputMessage.getPayload()).isEqualTo("Person: bob".getBytes());
}
}
@Test
public void testFunctionHonorsOutboundBindingContentType() {
try (ConfigurableApplicationContext context = new SpringApplicationBuilder(
@@ -339,6 +382,7 @@ public class FunctionInvokerTests {
}
public static class Person {
private String name;
public String getName() {
@@ -348,7 +392,72 @@ public class FunctionInvokerTests {
public void setName(String name) {
this.name = name;
}
}
}
@EnableAutoConfiguration
@EnableBinding(Processor.class)
public static class SimpleFluxFunctionConfiguration {
@Bean
public Function<Flux<Person>, Flux<String>> func() {
return x -> x.map(person -> person.toString());
}
public static class Person {
private String name;
public String getName() {
return name;
}
public void setName(String name) {
this.name = name;
}
public String toString() {
return "Person: " + name;
}
}
}
@EnableAutoConfiguration
@EnableBinding(Processor.class)
public static class SimpleFluxMessageFunctionConfiguration {
@Bean
public Function<Flux<Message<Person>>, Flux<Message<String>>> func() {
return x -> x.map(personMessage -> {
Person person = personMessage.getPayload();
Message<String> message = MessageBuilder.withPayload(person.toString())
.copyHeaders(personMessage.getHeaders()).build();
return message;
});
}
public static class Person {
private String name;
public String getName() {
return name;
}
public void setName(String name) {
this.name = name;
}
public String toString() {
return "Person: " + name;
}
}
}
@EnableAutoConfiguration