From 25a3591bfdd43e9df11878a01f9c68e6574d8561 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Tue, 21 Aug 2018 20:25:47 +0200 Subject: [PATCH] Added support for type-less Message to FunctionInvoker --- .../cloud/stream/function/FunctionInvoker.java | 4 +++- .../stream/function/FunctionInvokerTests.java | 17 +++++++++++++++++ 2 files changed, 20 insertions(+), 1 deletion(-) 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 426a13aa6..8a3376d83 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 @@ -87,6 +87,7 @@ class FunctionInvoker implements Function>, Flux onError(x, (Message) y)) .transform(this.userFunction::apply) // invoke user function + .onErrorContinue((x, y) -> onError(x, (Message) y)) .map(resultMessage -> toMessage(resultMessage, originalMessageRef.get())); // create output message } @@ -136,7 +137,8 @@ class FunctionInvoker implements Function>, Flux message) { - return !message.getPayload().getClass().isAssignableFrom(this.inputClass) && + return !this.inputClass.isAssignableFrom(Message.class) && + !message.getPayload().getClass().isAssignableFrom(this.inputClass) && !this.inputClass.isAssignableFrom(Object.class); } 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 5f397aee3..6ec512a3d 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 @@ -62,6 +62,11 @@ public class FunctionInvokerTests { new FunctionCatalogWrapper(context.getBean(FunctionCatalog.class)), context.getBean(FunctionInspector.class), context.getBean(CompositeMessageConverterFactory.class)); outputMessage = pojoToPojoSameType.apply(Flux.just(inputMessage)).blockFirst(); assertThat(inputMessage.getPayload()).isEqualTo(outputMessage.getPayload()); + + FunctionInvoker messageToMessageNoType = new FunctionInvoker<>("messageToMessageNoType", + new FunctionCatalogWrapper(context.getBean(FunctionCatalog.class)), context.getBean(FunctionInspector.class), context.getBean(CompositeMessageConverterFactory.class)); + outputMessage = messageToMessageNoType.apply(Flux.just(inputMessage)).blockFirst(); + assertThat(outputMessage).isInstanceOf(Message.class); } } @@ -72,6 +77,17 @@ public class FunctionInvokerTests { public Function, Message> messageToMessageDifferentType() { return x -> MessageBuilder.withPayload(new Bar()).copyHeaders(x.getHeaders()).build(); } + + @Bean + public Function, Message> messageToMessageAnyType() { + return x -> MessageBuilder.withPayload(new Bar()).copyHeaders(x.getHeaders()).build(); + } + + @Bean + public Function messageToMessageNoType() { + return x -> MessageBuilder.withPayload(new Bar()).copyHeaders(x.getHeaders()).build(); + } + @Bean public Function, Message> messageToMessageSameType() { return x -> x; @@ -81,6 +97,7 @@ public class FunctionInvokerTests { public Function pojoToPojoSameType() { return x -> x; } + } private static class Foo {