From f3c82b888b89ee55014e8618125635338e9ade3b Mon Sep 17 00:00:00 2001 From: omercelikceng Date: Sun, 16 Oct 2022 14:58:06 +0300 Subject: [PATCH] KafkaTopidProvisioner retry refactoring Use a retry template for topic description method call through admin client when provisioning consumer 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/issues/2520 --- .../provisioning/KafkaTopicProvisioner.java | 129 +++++++++--------- 1 file changed, 62 insertions(+), 67 deletions(-) diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java index 36081b5ac..48b9e339f 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java @@ -85,6 +85,7 @@ import org.springframework.util.StringUtils; * @author Oleg Zhurakousky * @author Aldo Sinanaj * @author Yi Liu + * @author Omer Celik */ public class KafkaTopicProvisioner implements // @checkstyle:off @@ -169,47 +170,27 @@ public class KafkaTopicProvisioner implements @Override public ProducerDestination provisionProducerDestination(final String name, - ExtendedProducerProperties properties) { + ExtendedProducerProperties properties) { if (logger.isInfoEnabled()) { logger.info("Using kafka topic for outbound: " + name); } - KafkaTopicUtils.validateTopicName(name); - try (AdminClient adminClient = createAdminClient()) { - createTopic(adminClient, name, properties.getPartitionCount(), false, + if (this.configurationProperties.isAutoCreateTopics()) { + KafkaTopicUtils.validateTopicName(name); + try (AdminClient adminClient = createAdminClient()) { + createTopic(adminClient, name, properties.getPartitionCount(), false, properties.getExtension().getTopic()); - int partitions = 0; - Map topicDescriptions = new HashMap<>(); - if (this.configurationProperties.isAutoCreateTopics()) { - 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 for: " + name, ex); - } - return null; - }); + int partitions = getPartitionsForTopic(name, adminClient); + return new KafkaProducerDestination(name, partitions); } - TopicDescription topicDescription = topicDescriptions.get(name); - if (topicDescription != null) { - partitions = topicDescription.partitions().size(); - } - return new KafkaProducerDestination(name, partitions); } + return new KafkaProducerDestination(name, 0); } @Override public ConsumerDestination provisionConsumerDestination(final String name, - final String group, - ExtendedConsumerProperties properties) { + final String group, + ExtendedConsumerProperties properties) { if (!properties.isMultiplex()) { return doProvisionConsumerDestination(name, group, properties); } @@ -221,56 +202,70 @@ public class KafkaTopicProvisioner implements return new KafkaConsumerDestination(name); } } - private ConsumerDestination doProvisionConsumerDestination(final String name, - final String group, - ExtendedConsumerProperties properties) { - + final String group, + ExtendedConsumerProperties properties) { + final KafkaConsumerDestination kafkaConsumerDestination = new KafkaConsumerDestination(name); if (properties.getExtension().isDestinationIsPattern()) { Assert.isTrue(!properties.getExtension().isEnableDlq(), - "enableDLQ is not allowed when listening to topic patterns"); + "enableDLQ is not allowed when listening to topic patterns"); if (logger.isDebugEnabled()) { logger.debug("Listening to a topic pattern - " + name - + " - no provisioning performed"); + + " - no provisioning performed"); } - return new KafkaConsumerDestination(name); + return kafkaConsumerDestination; } - KafkaTopicUtils.validateTopicName(name); - boolean anonymous = !StringUtils.hasText(group); - Assert.isTrue(!anonymous || !properties.getExtension().isEnableDlq(), + if (this.configurationProperties.isAutoCreateTopics()) { + KafkaTopicUtils.validateTopicName(name); + boolean anonymous = !StringUtils.hasText(group); + Assert.isTrue(!anonymous || !properties.getExtension().isEnableDlq(), "DLQ support is not available for anonymous subscriptions"); - if (properties.getInstanceCount() == 0) { - throw new IllegalArgumentException("Instance count cannot be zero"); - } - int partitionCount = properties.getInstanceCount() * properties.getConcurrency(); - ConsumerDestination consumerDestination = new KafkaConsumerDestination(name); - try (AdminClient adminClient = createAdminClient()) { - createTopic(adminClient, name, partitionCount, + if (properties.getInstanceCount() == 0) { + throw new IllegalArgumentException("Instance count cannot be zero"); + } + int partitionCount = properties.getInstanceCount() * properties.getConcurrency(); + ConsumerDestination consumerDestination; + try (AdminClient adminClient = createAdminClient()) { + createTopic(adminClient, name, partitionCount, properties.getExtension().isAutoRebalanceEnabled(), properties.getExtension().getTopic()); - if (this.configurationProperties.isAutoCreateTopics()) { - DescribeTopicsResult describeTopicsResult = adminClient - .describeTopics(Collections.singletonList(name)); - KafkaFuture> all = describeTopicsResult - .all(); - try { - Map topicDescriptions = all - .get(this.operationTimeout, TimeUnit.SECONDS); - TopicDescription topicDescription = topicDescriptions.get(name); - int partitions = topicDescription.partitions().size(); - consumerDestination = createDlqIfNeedBe(adminClient, name, group, - properties, anonymous, partitions); - if (consumerDestination == null) { - consumerDestination = new KafkaConsumerDestination(name, - partitions); - } - } - catch (Exception ex) { - throw new ProvisioningException("Provisioning exception encountered for " + name, ex); + int partitions = getPartitionsForTopic(name, adminClient); + consumerDestination = createDlqIfNeedBe(adminClient, name, group, + properties, anonymous, partitions); + if (consumerDestination == null) { + consumerDestination = new KafkaConsumerDestination(name, + partitions); } + return consumerDestination; } } - return consumerDestination; + return kafkaConsumerDestination; + } + + private int getPartitionsForTopic(String topicName, AdminClient adminClient) { + int partitions = 0; + Map topicDescriptions = retrieveTopicDescriptions(topicName, adminClient); + TopicDescription topicDescription = topicDescriptions.get(topicName); + if (topicDescription != null) { + partitions = topicDescription.partitions().size(); + } + return partitions; + } + + private Map retrieveTopicDescriptions(String topicName, AdminClient adminClient) { + return this.metadataRetryOperations.execute(context -> { + try { + if (logger.isDebugEnabled()) { + logger.debug("Attempting to retrieve the description for the topic: " + topicName); + } + DescribeTopicsResult describeTopicsResult = adminClient + .describeTopics(Collections.singletonList(topicName)); + KafkaFuture> all = describeTopicsResult.all(); + return all.get(this.operationTimeout, TimeUnit.SECONDS); + } catch (Exception ex) { + throw new ProvisioningException("Problems encountered with partitions finding for: " + topicName, ex); + } + }); } AdminClient createAdminClient() {