diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultBinderFactory.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultBinderFactory.java index f63c4bf89..3eb5c6a51 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultBinderFactory.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/DefaultBinderFactory.java @@ -48,6 +48,7 @@ import org.springframework.core.convert.support.GenericConversionService; import org.springframework.core.env.ConfigurableEnvironment; import org.springframework.core.env.MapPropertySource; import org.springframework.core.env.StandardEnvironment; +import org.springframework.messaging.MessageChannel; import org.springframework.util.Assert; import org.springframework.util.CollectionUtils; import org.springframework.util.StringUtils; @@ -143,6 +144,13 @@ public class DefaultBinderFactory implements BinderFactory, DisposableBean, Appl private Binder doGetBinder(String name, Class bindingTargetType) { + + if (!MessageChannel.class.isAssignableFrom(bindingTargetType) + && !PollableMessageSource.class.isAssignableFrom(bindingTargetType)) { + String bindingTargetTypeName = StringUtils.hasText(name) ? name : bindingTargetType.getSimpleName().toLowerCase(); + Binder binderInstance = getBinderInstance(bindingTargetTypeName); + return binderInstance; + } String configurationName; // Fall back to a default if no argument is provided if (StringUtils.isEmpty(name)) { @@ -165,8 +173,8 @@ public class DefaultBinderFactory implements BinderFactory, DisposableBean, Appl else { List candidatesForBindableType = new ArrayList<>(); for (String defaultCandidateConfiguration : defaultCandidateConfigurations) { - Binder binderInstance = getBinderInstance( - defaultCandidateConfiguration); + Binder binderInstance = getBinderInstance(defaultCandidateConfiguration); + Class binderType = GenericsUtils.getParameterType( binderInstance.getClass(), Binder.class, 0); if (binderType.isAssignableFrom(bindingTargetType)) { @@ -198,8 +206,7 @@ public class DefaultBinderFactory implements BinderFactory, DisposableBean, Appl else { configurationName = name; } - Binder binderInstance = getBinderInstance( - configurationName); + Binder binderInstance = getBinderInstance(configurationName); Assert.state(verifyBinderTypeMatchesTarget(binderInstance, bindingTargetType), "The binder '" + configurationName + "' cannot bind a " + bindingTargetType.getName()); @@ -225,9 +232,9 @@ public class DefaultBinderFactory implements BinderFactory, DisposableBean, Appl } @SuppressWarnings("unchecked") - private Binder getBinderInstance( - String configurationName) { + private Binder getBinderInstance(String configurationName) { if (!this.binderInstanceCache.containsKey(configurationName)) { + logger.info("Creating binder: " + configurationName); BinderConfiguration binderConfiguration = this.binderConfigurations .get(configurationName); Assert.state(binderConfiguration != null, @@ -327,9 +334,11 @@ public class DefaultBinderFactory implements BinderFactory, DisposableBean, Appl binderProducingContext); } } + logger.info("Caching the binder: " + configurationName); this.binderInstanceCache.put(configurationName, new SimpleImmutableEntry<>(binder, binderProducingContext)); } + logger.info("Retrieving cached binder: " + configurationName); return (Binder) this.binderInstanceCache .get(configurationName).getKey(); } diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/BindingService.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/BindingService.java index 124662d0f..04193794a 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/BindingService.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/BindingService.java @@ -24,11 +24,13 @@ import java.util.HashMap; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; +import java.util.stream.Stream; import com.fasterxml.jackson.databind.ObjectMapper; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; +import org.springframework.aop.framework.Advised; import org.springframework.beans.BeanUtils; import org.springframework.cloud.stream.binder.Binder; import org.springframework.cloud.stream.binder.BinderFactory; @@ -92,8 +94,13 @@ public class BindingService { @SuppressWarnings({ "unchecked", "rawtypes" }) public Collection> bindConsumer(T input, String inputName) { Collection> bindings = new ArrayList<>(); + Class inputClass = input.getClass(); + if (input instanceof Advised) { + inputClass = Stream.of(((Advised) input).getProxiedInterfaces()).filter(c -> !c.getName().contains("org.springframework")).findFirst() + .orElse(inputClass); + } Binder binder = (Binder) getBinder( - inputName, input.getClass()); + inputName, inputClass); ConsumerProperties consumerProperties = this.bindingServiceProperties .getConsumerProperties(inputName); if (binder instanceof ExtendedPropertiesBinder) { @@ -254,8 +261,13 @@ public class BindingService { public Binding bindProducer(T output, String outputName) { String bindingTarget = this.bindingServiceProperties .getBindingDestination(outputName); + Class outputClass = output.getClass(); + if (output instanceof Advised) { + outputClass = Stream.of(((Advised) output).getProxiedInterfaces()).filter(c -> !c.getName().contains("org.springframework")).findFirst() + .orElse(outputClass); + } Binder binder = (Binder) getBinder( - outputName, output.getClass()); + outputName, outputClass); ProducerProperties producerProperties = this.bindingServiceProperties .getProducerProperties(outputName); if (binder instanceof ExtendedPropertiesBinder) {