From c7becca82f772b3c36dbdc87867479ec19d8f223 Mon Sep 17 00:00:00 2001 From: Greg Harris Date: Tue, 18 Apr 2023 14:48:11 -0700 Subject: [PATCH 1/2] KAFKA-14905: Reduce flakiness in MM2 ForwardingAdmin test due to admin timeouts Signed-off-by: Greg Harris --- .../FakeForwardingAdminWithLocalMetadata.java | 27 +++++++++---------- 1 file changed, 12 insertions(+), 15 deletions(-) diff --git a/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/clients/admin/FakeForwardingAdminWithLocalMetadata.java b/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/clients/admin/FakeForwardingAdminWithLocalMetadata.java index 535dcaca9ee82..2f650796c9450 100644 --- a/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/clients/admin/FakeForwardingAdminWithLocalMetadata.java +++ b/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/clients/admin/FakeForwardingAdminWithLocalMetadata.java @@ -38,8 +38,6 @@ import java.util.Collection; import java.util.Map; import java.util.concurrent.ExecutionException; -import java.util.concurrent.TimeUnit; -import java.util.concurrent.TimeoutException; /** Customised ForwardingAdmin for testing only. * The class create/alter topics, partitions and ACLs in Kafka then store metadata in {@link FakeLocalMetadataStore}. @@ -47,7 +45,6 @@ public class FakeForwardingAdminWithLocalMetadata extends ForwardingAdmin { private static final Logger log = LoggerFactory.getLogger(FakeForwardingAdminWithLocalMetadata.class); - private final long timeout = 1000L; public FakeForwardingAdminWithLocalMetadata(Map configs) { super(configs); @@ -60,14 +57,14 @@ public CreateTopicsResult createTopics(Collection newTopics, CreateTop try { log.info("Add topic '{}' to cluster and metadata store", newTopic); // Wait for topic to be created before edit the fake local store - createTopicsResult.values().get(newTopic.name()).get(timeout, TimeUnit.MILLISECONDS); + createTopicsResult.values().get(newTopic.name()).get(); FakeLocalMetadataStore.addTopicToLocalMetadataStore(newTopic); - } catch (InterruptedException | ExecutionException | TimeoutException e) { + } catch (InterruptedException | ExecutionException e) { if (e.getCause() instanceof TopicExistsException) { log.warn("Topic '{}' already exists. Update the local metadata store if absent", newTopic.name()); FakeLocalMetadataStore.addTopicToLocalMetadataStore(newTopic); } else - log.error(e.getMessage()); + log.error("Unable to intercept admin client operation", e); } }); return createTopicsResult; @@ -79,10 +76,10 @@ public CreatePartitionsResult createPartitions(Map newPar newPartitions.forEach((topic, newPartition) -> { try { // Wait for topic partition to be created before edit the fake local store - createPartitionsResult.values().get(topic).get(timeout, TimeUnit.MILLISECONDS); + createPartitionsResult.values().get(topic).get(); FakeLocalMetadataStore.updatePartitionCount(topic, newPartition.totalCount()); - } catch (InterruptedException | ExecutionException | TimeoutException e) { - log.error(e.getMessage()); + } catch (InterruptedException | ExecutionException e) { + log.error("Unable to intercept admin client operation", e); } }); return createPartitionsResult; @@ -96,11 +93,11 @@ public AlterConfigsResult alterConfigs(Map configs, Alte try { if (configResource.type() == ConfigResource.Type.TOPIC) { // Wait for config to be altered before edit the fake local store - alterConfigsResult.values().get(configResource).get(timeout, TimeUnit.MILLISECONDS); + alterConfigsResult.values().get(configResource).get(); FakeLocalMetadataStore.updateTopicConfig(configResource.name(), newConfigs); } - } catch (InterruptedException | ExecutionException | TimeoutException e) { - log.error(e.getMessage()); + } catch (InterruptedException | ExecutionException e) { + log.error("Unable to intercept admin client operation", e); } }); return alterConfigsResult; @@ -112,12 +109,12 @@ public CreateAclsResult createAcls(Collection acls, CreateAclsOption CreateAclsResult aclsResult = super.createAcls(acls, options); try { // Wait for acls to be created before edit the fake local store - aclsResult.all().get(timeout, TimeUnit.MILLISECONDS); + aclsResult.all().get(); acls.forEach(aclBinding -> { FakeLocalMetadataStore.addACLs(aclBinding.entry().principal(), aclBinding); }); - } catch (InterruptedException | ExecutionException | TimeoutException e) { - log.error(e.getMessage()); + } catch (InterruptedException | ExecutionException e) { + log.error("Unable to intercept admin client operation", e); } return aclsResult; } From 38329c3674f043c9cd9ea4c70d66419552b46c9c Mon Sep 17 00:00:00 2001 From: Greg Harris Date: Fri, 21 Apr 2023 09:45:29 -0700 Subject: [PATCH 2/2] fixup: use whenComplete to avoid infinite blocking Signed-off-by: Greg Harris --- .../FakeForwardingAdminWithLocalMetadata.java | 59 ++++++++----------- 1 file changed, 24 insertions(+), 35 deletions(-) diff --git a/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/clients/admin/FakeForwardingAdminWithLocalMetadata.java b/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/clients/admin/FakeForwardingAdminWithLocalMetadata.java index 2f650796c9450..3ac8a8b17f00d 100644 --- a/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/clients/admin/FakeForwardingAdminWithLocalMetadata.java +++ b/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/clients/admin/FakeForwardingAdminWithLocalMetadata.java @@ -37,7 +37,6 @@ import java.util.Collection; import java.util.Map; -import java.util.concurrent.ExecutionException; /** Customised ForwardingAdmin for testing only. * The class create/alter topics, partitions and ACLs in Kafka then store metadata in {@link FakeLocalMetadataStore}. @@ -53,35 +52,29 @@ public FakeForwardingAdminWithLocalMetadata(Map configs) { @Override public CreateTopicsResult createTopics(Collection newTopics, CreateTopicsOptions options) { CreateTopicsResult createTopicsResult = super.createTopics(newTopics, options); - newTopics.forEach(newTopic -> { - try { - log.info("Add topic '{}' to cluster and metadata store", newTopic); - // Wait for topic to be created before edit the fake local store - createTopicsResult.values().get(newTopic.name()).get(); + newTopics.forEach(newTopic -> createTopicsResult.values().get(newTopic.name()).whenComplete((ignored, error) -> { + if (error == null) { FakeLocalMetadataStore.addTopicToLocalMetadataStore(newTopic); - } catch (InterruptedException | ExecutionException e) { - if (e.getCause() instanceof TopicExistsException) { - log.warn("Topic '{}' already exists. Update the local metadata store if absent", newTopic.name()); - FakeLocalMetadataStore.addTopicToLocalMetadataStore(newTopic); - } else - log.error("Unable to intercept admin client operation", e); + } else if (error.getCause() instanceof TopicExistsException) { + log.warn("Topic '{}' already exists. Update the local metadata store if absent", newTopic.name()); + FakeLocalMetadataStore.addTopicToLocalMetadataStore(newTopic); + } else { + log.error("Unable to intercept admin client operation", error); } - }); + })); return createTopicsResult; } @Override public CreatePartitionsResult createPartitions(Map newPartitions, CreatePartitionsOptions options) { CreatePartitionsResult createPartitionsResult = super.createPartitions(newPartitions, options); - newPartitions.forEach((topic, newPartition) -> { - try { - // Wait for topic partition to be created before edit the fake local store - createPartitionsResult.values().get(topic).get(); + newPartitions.forEach((topic, newPartition) -> createPartitionsResult.values().get(topic).whenComplete((ignored, error) -> { + if (error == null) { FakeLocalMetadataStore.updatePartitionCount(topic, newPartition.totalCount()); - } catch (InterruptedException | ExecutionException e) { - log.error("Unable to intercept admin client operation", e); + } else { + log.error("Unable to intercept admin client operation", error); } - }); + })); return createPartitionsResult; } @@ -89,17 +82,15 @@ public CreatePartitionsResult createPartitions(Map newPar @Override public AlterConfigsResult alterConfigs(Map configs, AlterConfigsOptions options) { AlterConfigsResult alterConfigsResult = super.alterConfigs(configs, options); - configs.forEach((configResource, newConfigs) -> { - try { + configs.forEach((configResource, newConfigs) -> alterConfigsResult.values().get(configResource).whenComplete((ignored, error) -> { + if (error == null) { if (configResource.type() == ConfigResource.Type.TOPIC) { - // Wait for config to be altered before edit the fake local store - alterConfigsResult.values().get(configResource).get(); FakeLocalMetadataStore.updateTopicConfig(configResource.name(), newConfigs); } - } catch (InterruptedException | ExecutionException e) { - log.error("Unable to intercept admin client operation", e); + } else { + log.error("Unable to intercept admin client operation", error); } - }); + })); return alterConfigsResult; } @@ -107,15 +98,13 @@ public AlterConfigsResult alterConfigs(Map configs, Alte @Override public CreateAclsResult createAcls(Collection acls, CreateAclsOptions options) { CreateAclsResult aclsResult = super.createAcls(acls, options); - try { - // Wait for acls to be created before edit the fake local store - aclsResult.all().get(); - acls.forEach(aclBinding -> { + aclsResult.values().forEach((aclBinding, future) -> future.whenComplete((ignored, error) -> { + if (error == null) { FakeLocalMetadataStore.addACLs(aclBinding.entry().principal(), aclBinding); - }); - } catch (InterruptedException | ExecutionException e) { - log.error("Unable to intercept admin client operation", e); - } + } else { + log.error("Unable to intercept admin client operation", error); + } + })); return aclsResult; } }