diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java index b46a43234..07de384c2 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java @@ -461,6 +461,22 @@ public class FunctionConfiguration { producerProperties, applicationContext)) { @Override protected void sendOutputs(Object result, Message requestMessage) { + if (result instanceof Iterable) { + for (Object resultElement : (Iterable) result) { + this.doSendMessage(resultElement, requestMessage); + } + } + else if (ObjectUtils.isArray(result)) { + for (int i = 0; i < ((Object[]) result).length; i++) { + this.doSendMessage(((Object[]) result)[i], requestMessage); + } + } + else { + this.doSendMessage(result, requestMessage); + } + } + + private void doSendMessage(Object result, Message requestMessage) { if (result instanceof Message && ((Message) result).getHeaders().get("spring.cloud.stream.sendto.destination") != null) { String destinationName = (String) ((Message) result).getHeaders().get("spring.cloud.stream.sendto.destination"); SubscribableChannel outputChannel = streamBridge.resolveDestination(destinationName, producerProperties); @@ -479,6 +495,8 @@ public class FunctionConfiguration { return handler; } + + private boolean isReactiveOrMultipleInputOutput(BindableProxyFactory bindableProxyFactory, Type functionType) { boolean reactiveInputsOutputs = FunctionTypeUtils.isReactive(FunctionTypeUtils.getInputType(functionType, 0)) || FunctionTypeUtils.isReactive(FunctionTypeUtils.getOutputType(functionType, 0));