GH-1966 Ensure dynamic destination property works for Messages wrapped as iterable/array types

Resolves #1966
This commit is contained in:
Oleg Zhurakousky
2020-06-17 15:08:13 +02:00
parent 9d37f48028
commit 9d5679ae1d

View File

@@ -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));