From 240ae8282e38ac35f9efa1794291fab0122e6042 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Wed, 27 May 2020 09:53:10 -0400 Subject: [PATCH] KafkaTopicProvisioner improvements Use a retry template for topic description method call through admin client when provisioning producer destinations. We are retrying because in the event this operation gets failed, it is retried with the default retry settings in the provisioner. Resolves https://github.com/spring-cloud/spring-cloud-stream-binder-kafka/issues/888 --- .../provisioning/KafkaTopicProvisioner.java | 35 +++++++++++-------- 1 file changed, 21 insertions(+), 14 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 b3fd86125..d7c5273c7 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 @@ -18,6 +18,7 @@ package org.springframework.cloud.stream.binder.kafka.provisioning; import java.util.Collection; import java.util.Collections; +import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.Set; @@ -147,21 +148,27 @@ public class KafkaTopicProvisioner implements createTopic(adminClient, name, properties.getPartitionCount(), false, properties.getExtension().getTopic()); int partitions = 0; + Map topicDescriptions = new HashMap<>(); if (this.configurationProperties.isAutoCreateTopics()) { - DescribeTopicsResult describeTopicsResult = adminClient - .describeTopics(Collections.singletonList(name)); - KafkaFuture> all = describeTopicsResult - .all(); - - Map topicDescriptions = null; - try { - topicDescriptions = all.get(this.operationTimeout, TimeUnit.SECONDS); - } - catch (Exception ex) { - throw new ProvisioningException( - "Problems encountered with partitions finding", ex); - } - TopicDescription topicDescription = topicDescriptions.get(name); + this.metadataRetryOperations.execute(context -> { + try { + if (logger.isDebugEnabled()) { + logger.debug("Attempting to retrieve the description for the topic: " + name); + } + DescribeTopicsResult describeTopicsResult = adminClient + .describeTopics(Collections.singletonList(name)); + KafkaFuture> all = describeTopicsResult + .all(); + topicDescriptions.putAll(all.get(this.operationTimeout, TimeUnit.SECONDS)); + } + catch (Exception ex) { + throw new ProvisioningException("Problems encountered with partitions finding", ex); + } + return null; + }); + } + TopicDescription topicDescription = topicDescriptions.get(name); + if (topicDescription != null) { partitions = topicDescription.partitions().size(); } return new KafkaProducerDestination(name, partitions);