GH-1949 Fix logic to deal with multiple binder multi-initializatioin

Fixed logic which was creating an instance of every available binder configuration. This is primarily an issue for
Kafka Streams binder which as a module brings 3 different binder configurations - KTable, KStream, GlobalKTable.

Adding more log statements.

Resolves #1949
This commit is contained in:
Oleg Zhurakousky
2020-05-05 16:42:36 +02:00
committed by Soby Chacko
parent e856ca694d
commit 286c057a8e
2 changed files with 29 additions and 8 deletions

View File

@@ -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 <T> Binder<T, ConsumerProperties, ProducerProperties> doGetBinder(String name,
Class<? extends T> bindingTargetType) {
if (!MessageChannel.class.isAssignableFrom(bindingTargetType)
&& !PollableMessageSource.class.isAssignableFrom(bindingTargetType)) {
String bindingTargetTypeName = StringUtils.hasText(name) ? name : bindingTargetType.getSimpleName().toLowerCase();
Binder<T, ConsumerProperties, ProducerProperties> 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<String> candidatesForBindableType = new ArrayList<>();
for (String defaultCandidateConfiguration : defaultCandidateConfigurations) {
Binder<Object, ?, ?> binderInstance = getBinderInstance(
defaultCandidateConfiguration);
Binder<Object, ?, ?> 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<T, ConsumerProperties, ProducerProperties> binderInstance = getBinderInstance(
configurationName);
Binder<T, ConsumerProperties, ProducerProperties> 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 <T> Binder<T, ConsumerProperties, ProducerProperties> getBinderInstance(
String configurationName) {
private <T> Binder<T, ConsumerProperties, ProducerProperties> 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<T, ConsumerProperties, ProducerProperties>) this.binderInstanceCache
.get(configurationName).getKey();
}

View File

@@ -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 <T> Collection<Binding<T>> bindConsumer(T input, String inputName) {
Collection<Binding<T>> 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<T, ConsumerProperties, ?> binder = (Binder<T, ConsumerProperties, ?>) getBinder(
inputName, input.getClass());
inputName, inputClass);
ConsumerProperties consumerProperties = this.bindingServiceProperties
.getConsumerProperties(inputName);
if (binder instanceof ExtendedPropertiesBinder) {
@@ -254,8 +261,13 @@ public class BindingService {
public <T> Binding<T> 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<T, ?, ProducerProperties> binder = (Binder<T, ?, ProducerProperties>) getBinder(
outputName, output.getClass());
outputName, outputClass);
ProducerProperties producerProperties = this.bindingServiceProperties
.getProducerProperties(outputName);
if (binder instanceof ExtendedPropertiesBinder) {