From 219d790910af8203efeec6c0c3ef1290eae9ff63 Mon Sep 17 00:00:00 2001 From: Ilayaperumal Gopinathan Date: Thu, 23 Jul 2015 10:50:50 -0700 Subject: [PATCH] Set channel resolver for binding adapter - Set `ChannelBindingAdapter`'s channel resolver to use `BinderAwareChannelResolver`. - Previously the runtime `Binder` singleton object was lazily instantiated. But the instantiation actually happened at the constructor of BinderAwareRouterBeanPostProcessor itself. Hence, removed the lazy instantiation of the Binder object. Also, I don't see any impact of creating the `Binder` object at the time of BeanPostProcessor creation. --- .../stream/adapter/ChannelBindingAdapter.java | 6 +- .../BinderAwareRouterBeanPostProcessor.java | 6 +- .../ChannelBindingAdapterConfiguration.java | 82 ++++--------------- 3 files changed, 21 insertions(+), 73 deletions(-) diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/adapter/ChannelBindingAdapter.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/adapter/ChannelBindingAdapter.java index 72dc8abd5..9bf731fb5 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/adapter/ChannelBindingAdapter.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/adapter/ChannelBindingAdapter.java @@ -28,7 +28,9 @@ import java.util.concurrent.atomic.AtomicBoolean; import org.slf4j.Logger; import org.slf4j.LoggerFactory; + import org.springframework.beans.BeansException; +import org.springframework.cloud.stream.binder.Binder; import org.springframework.cloud.stream.binder.BinderHeaders; import org.springframework.cloud.stream.config.ChannelBindingProperties; import org.springframework.context.ApplicationContext; @@ -40,7 +42,6 @@ import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.channel.interceptor.WireTap; import org.springframework.integration.support.DefaultMessageBuilderFactory; import org.springframework.integration.support.MessageBuilderFactory; -import org.springframework.integration.support.channel.BeanFactoryChannelResolver; import org.springframework.jmx.export.annotation.ManagedAttribute; import org.springframework.jmx.export.annotation.ManagedOperation; import org.springframework.jmx.export.annotation.ManagedResource; @@ -50,7 +51,6 @@ import org.springframework.messaging.core.DestinationResolver; import org.springframework.messaging.support.ChannelInterceptorAdapter; import org.springframework.util.Assert; import org.springframework.util.StringUtils; -import org.springframework.cloud.stream.binder.Binder; /** * Binds input/output channels. @@ -91,7 +91,6 @@ public class ChannelBindingAdapter implements Lifecycle, ApplicationContextAware public ChannelBindingAdapter(ChannelBindingProperties module, Binder binder) { this.module = module; this.binder = binder; - this.channelLocator = new DefaultChannelLocator(module); } public void setChannelLocator(ChannelLocator channelLocator) { @@ -101,7 +100,6 @@ public class ChannelBindingAdapter implements Lifecycle, ApplicationContextAware @Override public void setApplicationContext(ApplicationContext applicationContext) throws BeansException { this.applicationContext = (ConfigurableApplicationContext) applicationContext; - this.channelResolver = new BeanFactoryChannelResolver(applicationContext); } public void setChannelResolver(DestinationResolver channelResolver) { diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/BinderAwareRouterBeanPostProcessor.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/BinderAwareRouterBeanPostProcessor.java index 0d0dffd1e..996ee7dc8 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/BinderAwareRouterBeanPostProcessor.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/BinderAwareRouterBeanPostProcessor.java @@ -16,8 +16,6 @@ package org.springframework.cloud.stream.binder; -import java.util.Properties; - import org.springframework.beans.BeansException; import org.springframework.beans.factory.BeanFactory; import org.springframework.beans.factory.BeanFactoryAware; @@ -35,8 +33,8 @@ public class BinderAwareRouterBeanPostProcessor implements BeanPostProcessor, Be private final BinderAwareChannelResolver channelResolver; - public BinderAwareRouterBeanPostProcessor(Binder binder, Properties producerProperties) { - this.channelResolver = new BinderAwareChannelResolver(binder, producerProperties); + public BinderAwareRouterBeanPostProcessor(BinderAwareChannelResolver channelResolver) { + this.channelResolver = channelResolver; } @Override diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/ChannelBindingAdapterConfiguration.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/ChannelBindingAdapterConfiguration.java index 7c7aebde5..3d8183f55 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/ChannelBindingAdapterConfiguration.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/ChannelBindingAdapterConfiguration.java @@ -15,40 +15,37 @@ */ package org.springframework.cloud.stream.config; -import java.util.Arrays; import java.util.Collection; import java.util.LinkedHashSet; import java.util.Properties; import java.util.Set; -import org.aopalliance.intercept.MethodInterceptor; -import org.aopalliance.intercept.MethodInvocation; -import org.springframework.aop.framework.ProxyFactory; -import org.springframework.aop.target.LazyInitTargetSource; import org.springframework.beans.factory.BeanFactoryUtils; -import org.springframework.beans.factory.ListableBeanFactory; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.config.BeanDefinition; import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; import org.springframework.beans.factory.support.AbstractBeanDefinition; import org.springframework.cloud.stream.adapter.ChannelBindingAdapter; -import org.springframework.cloud.stream.adapter.ChannelLocator; +import org.springframework.cloud.stream.adapter.DefaultChannelLocator; import org.springframework.cloud.stream.adapter.InputChannelBinding; import org.springframework.cloud.stream.adapter.OutputChannelBinding; import org.springframework.cloud.stream.annotation.Input; import org.springframework.cloud.stream.annotation.Output; +import org.springframework.cloud.stream.binder.Binder; +import org.springframework.cloud.stream.binder.BinderAwareChannelResolver; import org.springframework.cloud.stream.binder.BinderAwareRouterBeanPostProcessor; import org.springframework.cloud.stream.endpoint.ChannelsEndpoint; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.messaging.MessageChannel; -import org.springframework.util.Assert; -import org.springframework.cloud.stream.binder.Binder; /** + * Configuration class that provides necessary beans for {@link MessageChannel} binding. + * * @author Dave Syer * @author David Turanski * @author Marius Bogoevici + * @author Ilayaperumal Gopinathan */ @Configuration public class ChannelBindingAdapterConfiguration { @@ -59,19 +56,16 @@ public class ChannelBindingAdapterConfiguration { @Autowired private ConfigurableListableBeanFactory beanFactory; - private ChannelLocator channelLocator; - @Autowired private Binder binder; @Bean public ChannelBindingAdapter bindingAdapter() { ChannelBindingAdapter adapter = new ChannelBindingAdapter(this.module, this.binder); + adapter.setChannelLocator(new DefaultChannelLocator(module)); adapter.setOutputChannels(getOutputChannels()); adapter.setInputChannels(getInputChannels()); - if (this.channelLocator != null) { - adapter.setChannelLocator(this.channelLocator); - } + adapter.setChannelResolver(binderAwareChannelResolver()); return adapter; } @@ -80,13 +74,7 @@ public class ChannelBindingAdapterConfiguration { return new ChannelsEndpoint(adapter); } - public void refresh() { - ChannelBindingAdapter adapter = bindingAdapter(); - adapter.setOutputChannels(getOutputChannels()); - adapter.setInputChannels(getInputChannels()); - } - - protected Collection getOutputChannels() { + Collection getOutputChannels() { Set channels = new LinkedHashSet<>(); String[] names = this.beanFactory.getBeanNamesForType(MessageChannel.class); for (String name : names) { @@ -100,7 +88,7 @@ public class ChannelBindingAdapterConfiguration { return channels; } - protected Collection getInputChannels() { + Collection getInputChannels() { Set channels = new LinkedHashSet<>(); String[] names = this.beanFactory.getBeanNamesForType(MessageChannel.class); for (String name : names) { @@ -114,49 +102,13 @@ public class ChannelBindingAdapterConfiguration { return channels; } - // Nested class to avoid instantiating all of the above early - @Configuration - protected static class BinderAwareRouterConfiguration { + @Bean + public BinderAwareChannelResolver binderAwareChannelResolver() { + return new BinderAwareChannelResolver(BeanFactoryUtils.beanOfType(beanFactory, Binder.class), new Properties()); + } - @Autowired - private ListableBeanFactory beanFactory; - - @Bean - public BinderAwareRouterBeanPostProcessor binderAwareRouterBeanPostProcessor() { - - return new BinderAwareRouterBeanPostProcessor(createLazyProxy( - this.beanFactory, Binder.class), new Properties()); - } - - private T createLazyProxy(ListableBeanFactory beanFactory, Class type) { - ProxyFactory factory = new ProxyFactory(); - LazyInitTargetSource source = new LazyInitTargetSource(); - source.setTargetClass(type); - source.setTargetBeanName(getBeanNameFor(beanFactory, Binder.class)); - source.setBeanFactory(beanFactory); - factory.setTargetSource(source); - factory.addAdvice(new PassthruAdvice()); - factory.setInterfaces(new Class[] {type}); - @SuppressWarnings("unchecked") - T proxy = (T) factory.getProxy(); - return proxy; - } - - private String getBeanNameFor(ListableBeanFactory beanFactory, Class type) { - String[] names = BeanFactoryUtils.beanNamesForTypeIncludingAncestors( - beanFactory, type, false, false); - Assert.state(names.length == 1, "No unique Binder (found " + names.length - + ": " + Arrays.asList(names) + ")"); - return names[0]; - } - - private class PassthruAdvice implements MethodInterceptor { - - @Override - public Object invoke(MethodInvocation invocation) throws Throwable { - return invocation.proceed(); - } - - } + @Bean + public BinderAwareRouterBeanPostProcessor binderAwareRouterBeanPostProcessor(BinderAwareChannelResolver resolver) { + return new BinderAwareRouterBeanPostProcessor(resolver); } }