diff --git a/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java b/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java index b25894c2e0452..c30ca43eb99bb 100644 --- a/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java +++ b/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java @@ -980,6 +980,7 @@ public Map getMainConsumerConfigs(final String groupId, // add admin retries configs for creating topics final AdminClientConfig adminClientDefaultConfig = new AdminClientConfig(getClientPropsWithPrefix(ADMIN_CLIENT_PREFIX, AdminClientConfig.configNames())); consumerProps.put(adminClientPrefix(AdminClientConfig.RETRIES_CONFIG), adminClientDefaultConfig.getInt(AdminClientConfig.RETRIES_CONFIG)); + consumerProps.put(adminClientPrefix(AdminClientConfig.RETRY_BACKOFF_MS_CONFIG), adminClientDefaultConfig.getLong(AdminClientConfig.RETRY_BACKOFF_MS_CONFIG)); // verify that producer batch config is no larger than segment size, then add topic configs required for creating topics final Map topicProps = originalsWithPrefix(TOPIC_PREFIX, false); diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/InternalTopicManager.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/InternalTopicManager.java index 6159ee25d6feb..7e35126d26310 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/InternalTopicManager.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/InternalTopicManager.java @@ -57,6 +57,7 @@ private InternalAdminClientConfig(final Map props) { private final AdminClient adminClient; private final int retries; + private final long retryBackOffMs; public InternalTopicManager(final AdminClient adminClient, final StreamsConfig streamsConfig) { @@ -67,7 +68,9 @@ public InternalTopicManager(final AdminClient adminClient, replicationFactor = streamsConfig.getInt(StreamsConfig.REPLICATION_FACTOR_CONFIG).shortValue(); windowChangeLogAdditionalRetention = streamsConfig.getLong(StreamsConfig.WINDOW_STORE_CHANGE_LOG_ADDITIONAL_RETENTION_MS_CONFIG); - retries = new InternalAdminClientConfig(streamsConfig.getAdminConfigs("dummy")).getInt(AdminClientConfig.RETRIES_CONFIG); + final InternalAdminClientConfig dummyAdmin = new InternalAdminClientConfig(streamsConfig.getAdminConfigs("dummy")); + retries = dummyAdmin.getInt(AdminClientConfig.RETRIES_CONFIG); + retryBackOffMs = dummyAdmin.getLong(AdminClientConfig.RETRY_BACKOFF_MS_CONFIG); log.debug("Configs:" + Utils.NL, "\t{} = {}" + Utils.NL, @@ -115,17 +118,22 @@ public void makeReady(final Map topics) { // TODO: KAFKA-6928. should not need retries in the outer caller as it will be retried internally in admin client int remainingRetries = retries; + boolean retryBackOff = false; boolean retry; do { retry = false; final CreateTopicsResult createTopicsResult = adminClient.createTopics(newTopics); - final Set createTopicNames = new HashSet<>(); + final Set createdTopicNames = new HashSet<>(); for (final Map.Entry> createTopicResult : createTopicsResult.values().entrySet()) { try { + if (retryBackOff) { + retryBackOff = false; + Thread.sleep(retryBackOffMs); + } createTopicResult.getValue().get(); - createTopicNames.add(createTopicResult.getKey()); + createdTopicNames.add(createTopicResult.getKey()); } catch (final ExecutionException couldNotCreateTopic) { final Throwable cause = couldNotCreateTopic.getCause(); final String topicName = createTopicResult.getKey(); @@ -135,10 +143,23 @@ public void makeReady(final Map topics) { log.debug("Could not get number of partitions for topic {} due to timeout. " + "Will try again (remaining retries {}).", topicName, remainingRetries - 1); } else if (cause instanceof TopicExistsException) { - createTopicNames.add(createTopicResult.getKey()); - log.info("Topic {} exist already: {}", - topicName, - couldNotCreateTopic.toString()); + // This topic didn't exist earlier, it might be marked for deletion or it might differ + // from the desired setup. It needs re-validation. + final Map existingTopicPartition = getNumPartitions(Collections.singleton(topicName)); + + if (existingTopicPartition.containsKey(topicName) + && validateTopicPartitions(Collections.singleton(topics.get(topicName)), existingTopicPartition).isEmpty()) { + createdTopicNames.add(createTopicResult.getKey()); + log.info("Topic {} exists already and has the right number of partitions: {}", + topicName, + couldNotCreateTopic.toString()); + } else { + retry = true; + retryBackOff = true; + log.info("Could not create topic {}. Topic is probably marked for deletion (number of partitions is unknown).\n" + + "Will retry to create this topic in {} ms (to let broker finish async delete operation first).\n" + + "Error message was: {}", topicName, retryBackOffMs, couldNotCreateTopic.toString()); + } } else { throw new StreamsException(String.format("Could not create topic %s.", topicName), couldNotCreateTopic); @@ -151,7 +172,7 @@ public void makeReady(final Map topics) { } if (retry) { - newTopics.removeIf(newTopic -> createTopicNames.contains(newTopic.name())); + newTopics.removeIf(newTopic -> createdTopicNames.contains(newTopic.name())); continue; } diff --git a/streams/src/test/java/org/apache/kafka/streams/StreamsConfigTest.java b/streams/src/test/java/org/apache/kafka/streams/StreamsConfigTest.java index a86c38946df77..c510f7d83ad13 100644 --- a/streams/src/test/java/org/apache/kafka/streams/StreamsConfigTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/StreamsConfigTest.java @@ -129,9 +129,11 @@ public void consumerConfigMustContainStreamPartitionAssignorConfig() { } @Test - public void consumerConfigMustUseAdminClientConfigForRetries() { + public void consumerConfigShouldContainAdminClientConfigsForRetriesAndRetryBackOffMsWithAdminPrefix() { props.put(StreamsConfig.adminClientPrefix(StreamsConfig.RETRIES_CONFIG), 20); + props.put(StreamsConfig.adminClientPrefix(StreamsConfig.RETRY_BACKOFF_MS_CONFIG), 200L); props.put(StreamsConfig.RETRIES_CONFIG, 10); + props.put(StreamsConfig.RETRY_BACKOFF_MS_CONFIG, 100L); final StreamsConfig streamsConfig = new StreamsConfig(props); final String groupId = "example-application"; @@ -139,6 +141,7 @@ public void consumerConfigMustUseAdminClientConfigForRetries() { final Map returnedProps = streamsConfig.getMainConsumerConfigs(groupId, clientId); assertEquals(20, returnedProps.get(StreamsConfig.adminClientPrefix(StreamsConfig.RETRIES_CONFIG))); + assertEquals(200L, returnedProps.get(StreamsConfig.adminClientPrefix(StreamsConfig.RETRY_BACKOFF_MS_CONFIG))); } @Test