From e2a41e78718a605b3d770e6278dcd1bd39b0a041 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Thu, 31 Oct 2019 20:05:29 +0100 Subject: [PATCH] GH-1837 Fix dynamic destination resolution for Supplier Resolves #1837 --- .../cloud/stream/function/FunctionConfiguration.java | 11 +++++++++-- 1 file changed, 9 insertions(+), 2 deletions(-) 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 a4a544e60..1ba3a8985 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 @@ -141,7 +141,7 @@ public class FunctionConfiguration { @Bean InitializingBean supplierInitializer(FunctionCatalog functionCatalog, StreamFunctionProperties functionProperties, GenericApplicationContext context, BindingServiceProperties serviceProperties, - @Nullable BindableFunctionProxyFactory[] proxyFactories) { + @Nullable BindableFunctionProxyFactory[] proxyFactories, BinderAwareChannelResolver dynamicDestinationResolver) { if (!ObjectUtils.isEmpty(context.getBeanNamesForAnnotation(EnableBinding.class)) || proxyFactories == null) { return null; @@ -173,8 +173,15 @@ public class FunctionConfiguration { if (!functionProperties.isComposeFrom() && !functionProperties.isComposeTo()) { String integrationFlowName = proxyFactory.getFunctionDefinition() + "_integrationflow"; PollableBean pollable = extractPollableAnnotation(functionProperties, context, proxyFactory); + IntegrationFlow integrationFlow = integrationFlowFromProvidedSupplier(functionWrapper, beginPublishingTrigger, pollable) - .channel(outputName).get(); + .route(Message.class, message -> { + if (message.getHeaders().get("spring.cloud.stream.sendto.destination") != null) { + String destinationName = (String) message.getHeaders().get("spring.cloud.stream.sendto.destination"); + return dynamicDestinationResolver.resolveDestination(destinationName); + } + return outputName; + }).get(); IntegrationFlow postProcessedFlow = (IntegrationFlow) context.getAutowireCapableBeanFactory() .applyBeanPostProcessorsBeforeInitialization(integrationFlow, integrationFlowName); context.registerBean(integrationFlowName, IntegrationFlow.class, () -> postProcessedFlow);