From 4760a4a066b454bf5ea98eb5e7d40d9e693b6466 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Thu, 29 Apr 2021 20:11:26 +0200 Subject: [PATCH] interim --- .../stream/binder/DefaultBinderFactory.java | 24 ++++++++++++------- .../BinderFactoryAutoConfiguration.java | 10 ++++++++ 2 files changed, 25 insertions(+), 9 deletions(-) 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 965be3954..04b714ab6 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 @@ -140,7 +140,8 @@ public class DefaultBinderFactory implements BinderFactory, DisposableBean, Appl String[] bindingTargetTypeNameForLambda = new String[]{bindingTargetTypeName}; final Optional foundBinder = binders.keySet().stream().filter(b -> b.equals(bindingTargetTypeNameForLambda[0])).findFirst(); if (foundBinder.isPresent()) { // found a match - now do all the customizations. - binder = binders.get(foundBinder.get()); +// binder = binders.get(foundBinder.get()); + binder = this.doGetBinder(binderName, bindingTargetType); if (!CollectionUtils.isEmpty(this.listeners)) { for (Listener binderFactoryListener : this.listeners) { binderFactoryListener.afterBinderContextInitialized(bindingTargetTypeName, @@ -153,15 +154,18 @@ public class DefaultBinderFactory implements BinderFactory, DisposableBean, Appl binder = (Binder) this.context .getBean(binderName); } - else if (binders.size() == 1) { + else if (binders.size() == 1 && name == null) { binder = binders.values().iterator().next(); } - else if (binders.size() > 1) { - throw new IllegalStateException( - "Multiple binders are available, however neither default nor " - + "per-destination binder name is provided. Available binders are " - + binders.keySet()); + else if (binders.size() == 1 && binders.keySet().iterator().next().equals(name)) { + binder = binders.values().iterator().next(); } +// else if (binders.size() > 1) { +// throw new IllegalStateException( +// "Multiple binders are available, however neither default nor " +// + "per-destination binder name is provided. Available binders are " +// + binders.keySet()); +// } else { /* * This is the fall back to the old bootstrap that relies on spring.binders. @@ -311,7 +315,7 @@ public class DefaultBinderFactory implements BinderFactory, DisposableBean, Appl && this.context != null; if (useApplicationContextAsParent) { - springApplicationBuilder.parent(this.context); +// springApplicationBuilder.parent(this.context); } else { this.customizeParentChildContextRelationship(springApplicationBuilder, this.context); @@ -355,9 +359,12 @@ public class DefaultBinderFactory implements BinderFactory, DisposableBean, Appl springApplicationBuilder.environment(binderEnvironment); } + + ConfigurableApplicationContext binderProducingContext = springApplicationBuilder .run(args.toArray(new String[0])); + Binder binder = binderProducingContext.getBean(Binder.class); Map messageConverters = binderProducingContext.getBeansOfType(MessageConverter.class); if (!CollectionUtils.isEmpty(messageConverters) && !ObjectUtils.isEmpty(context.getBeansOfType(FunctionCatalog.class))) { @@ -373,7 +380,6 @@ public class DefaultBinderFactory implements BinderFactory, DisposableBean, Appl } } - Binder binder = binderProducingContext.getBean(Binder.class); /* * This will ensure that application defined errorChannel and other beans are * accessible within binder's context (see diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BinderFactoryAutoConfiguration.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BinderFactoryAutoConfiguration.java index 642333bdf..e002e0a65 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BinderFactoryAutoConfiguration.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BinderFactoryAutoConfiguration.java @@ -195,6 +195,16 @@ public class BinderFactoryAutoConfiguration { binderTypes.put(binderType.getDefaultName(), binderType); } } + + try { + BinderType kafkaType = new BinderType("kafka", new Class[] {Class.forName("org.springframework.cloud.stream.binder.kafka.config.KafkaBinderConfiguration")}); + binderTypes.put("kafka", kafkaType); + } + catch (Exception e) { + // TODO: handle exception + } +// kafka:\ +// } }