From b5a0013e1e83516bb50ec556aa2837229ec62719 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Thu, 8 Mar 2018 09:11:06 -0500 Subject: [PATCH] GH-330 Polishing Resolves #330 Resolves #334 --- .../provisioning/KafkaTopicProvisioner.java | 33 +++++++------------ .../stream/binder/kafka/AdminConfigTests.java | 20 ++--------- 2 files changed, 14 insertions(+), 39 deletions(-) diff --git a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java index a4bd6424f..f1aaaec03 100644 --- a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java +++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java @@ -69,6 +69,7 @@ import org.springframework.util.StringUtils; * @author Gary Russell * @author Ilayaperumal Gopinathan * @author Simon Flandergan + * @author Oleg Zhurakousky */ public class KafkaTopicProvisioner implements ProvisioningProvider, ExtendedProducerProperties>, InitializingBean { @@ -128,23 +129,22 @@ public class KafkaTopicProvisioner implements ProvisioningProvider> all = describeTopicsResult.all(); + Map topicDescriptions = null; try { - Map topicDescriptions = all.get(this.operationTimeout, TimeUnit.SECONDS); - TopicDescription topicDescription = topicDescriptions.get(name); - int partitions = topicDescription.partitions().size(); - return new KafkaProducerDestination(name, partitions); + topicDescriptions = all.get(this.operationTimeout, TimeUnit.SECONDS); } catch (Exception e) { throw new ProvisioningException("Problems encountered with partitions finding", e); } + TopicDescription topicDescription = topicDescriptions.get(name); + partitions = topicDescription.partitions().size(); } - else { - return new KafkaProducerDestination(name); - } + return new KafkaProducerDestination(name, partitions); } } @@ -160,6 +160,7 @@ public class KafkaTopicProvisioner implements ProvisioningProvider topicDescriptions = all.get(operationTimeout, TimeUnit.SECONDS); TopicDescription topicDescription = topicDescriptions.get(name); int partitions = topicDescription.partitions().size(); - ConsumerDestination dlqTopic = createDlqIfNeedBe(adminClient, name, group, properties, anonymous, - partitions); - if (dlqTopic != null) { - return dlqTopic; + consumerDestination = createDlqIfNeedBe(adminClient, name, group, properties, anonymous, partitions); + if (consumerDestination == null) { + consumerDestination = new KafkaConsumerDestination(name, partitions); } - return new KafkaConsumerDestination(name, partitions); } catch (Exception e) { throw new ProvisioningException("provisioning exception", e); } } } - return new KafkaConsumerDestination(name); + return consumerDestination; } AdminClient createAdminClient() { @@ -399,10 +398,6 @@ public class KafkaTopicProvisioner implements ProvisioningProvider