From 8edbe6ef0cee1948ed30570c01ad404290941fbf Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Tue, 19 Feb 2019 14:49:10 +0100 Subject: [PATCH] GH-1617 added support for Flux functions to FunctionInvoker Resolves #1617 --- .../stream/function/FunctionInvoker.java | 9 +- .../stream/function/FunctionInvokerTests.java | 117 +++++++++++++++++- 2 files changed, 120 insertions(+), 6 deletions(-) diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionInvoker.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionInvoker.java index f18f28fe8..5319a8db9 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionInvoker.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionInvoker.java @@ -169,7 +169,8 @@ class FunctionInvoker implements Function>, Flux) (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 implements Function>, Flux) argument).getPayload(); } return argument; diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/FunctionInvokerTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/FunctionInvokerTests.java index 99be5c010..f8a17e38a 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/FunctionInvokerTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/FunctionInvokerTests.java @@ -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 inputMessage = MessageBuilder + .withPayload("{\"name\":\"bob\"}".getBytes()).build(); + inputDestination.send(inputMessage); + + Message 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 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 inputMessage = MessageBuilder + .withPayload("{\"name\":\"bob\"}".getBytes()).build(); + inputDestination.send(inputMessage); + + Message 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> 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>> func() { + return x -> x.map(personMessage -> { + Person person = personMessage.getPayload(); + Message 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