From e4a82107d86ae065a86015df598be02e28a9d4e0 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Sun, 9 Sep 2018 14:32:00 +0200 Subject: [PATCH] GH-1465 Polishing function support for channel binder Removed BinderFunctionSupport in favor of private method (at least for now) Removed 'context.getBean()` from IntegrationFlowFunctionSupport in favor of propper DI Resolves #1465 --- .../binder/AbstractMessageChannelBinder.java | 34 +++++++-------- .../function/BinderFunctionSupport.java | 42 ------------------- .../IntegrationFlowFunctionSupport.java | 20 ++------- 3 files changed, 20 insertions(+), 76 deletions(-) delete mode 100644 spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/BinderFunctionSupport.java diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java index 431886634..9508c9375 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java @@ -18,17 +18,17 @@ package org.springframework.cloud.stream.binder; import java.util.LinkedHashMap; import java.util.Map; +import java.util.function.Function; import com.fasterxml.jackson.databind.ObjectMapper; import org.apache.commons.logging.Log; import org.springframework.beans.factory.DisposableBean; import org.springframework.beans.factory.InitializingBean; -import org.springframework.beans.factory.NoSuchBeanDefinitionException; +import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; import org.springframework.beans.factory.support.DefaultSingletonBeanRegistry; import org.springframework.cloud.stream.config.ListenerContainerCustomizer; -import org.springframework.cloud.stream.function.BinderFunctionSupport; import org.springframework.cloud.stream.function.IntegrationFlowFunctionSupport; import org.springframework.cloud.stream.provisioning.ConsumerDestination; import org.springframework.cloud.stream.provisioning.ProducerDestination; @@ -39,6 +39,8 @@ import org.springframework.context.ApplicationEventPublisher; import org.springframework.context.ApplicationEventPublisherAware; import org.springframework.context.Lifecycle; import org.springframework.integration.channel.AbstractMessageChannel; +import org.springframework.integration.channel.DirectChannel; +import org.springframework.integration.channel.MessageChannelReactiveUtils; import org.springframework.integration.channel.PublishSubscribeChannel; import org.springframework.integration.context.IntegrationContextUtils; import org.springframework.integration.core.MessageProducer; @@ -103,6 +105,7 @@ public abstract class AbstractMessageChannelBinder doGetExtendedInfo(Object destination, Object properties) { Map extendedInfo = new LinkedHashMap<>(); extendedInfo.put("bindingDestination", destination.toString()); @@ -777,6 +767,16 @@ public abstract class AbstractMessageChannelBinder inputPublisher = Flux.from(publisher); subscribeToOutput(outputProcessor, functionInvoker.apply((Flux>) inputPublisher)).subscribe(); } - - @Override - public void setApplicationContext(ApplicationContext applicationContext) throws BeansException { - this.applicationContext = applicationContext; - } - - @Override - public void afterPropertiesSet() throws Exception { - this.errorChannel = (MessageChannel) applicationContext.getBean(IntegrationContextUtils.ERROR_CHANNEL_BEAN_NAME); - } }