From 7cf648930fbe49f1423549ef9c7d91a427258720 Mon Sep 17 00:00:00 2001 From: Deqi Hu Date: Sun, 25 Jun 2023 20:26:36 +0800 Subject: [PATCH 1/5] KAFKA-15119:Support incremental syncTopicAcls in MirrorSourceConnector --- .../connect/mirror/MirrorSourceConnector.java | 21 +++++++++++++------ 1 file changed, 15 insertions(+), 6 deletions(-) diff --git a/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorSourceConnector.java b/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorSourceConnector.java index b97fc7bc577ee..459ddc4a08cfb 100644 --- a/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorSourceConnector.java +++ b/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorSourceConnector.java @@ -104,6 +104,7 @@ public class MirrorSourceConnector extends SourceConnector { private Admin offsetSyncsAdminClient; private volatile boolean useIncrementalAlterConfigs; private AtomicBoolean noAclAuthorizer = new AtomicBoolean(false); + private List knownTopicAclBindings = Collections.emptyList(); public MirrorSourceConnector() { // nop @@ -582,12 +583,20 @@ void incrementalAlterConfigs(Map topicConfigs) { } private void updateTopicAcls(List bindings) { - log.trace("Syncing {} topic ACL bindings.", bindings.size()); - targetAdminClient.createAcls(bindings).values().forEach((k, v) -> v.whenComplete((x, e) -> { - if (e != null) { - log.warn("Could not sync ACL of topic {}.", k.pattern().name(), e); - } - })); + Set addBindings = new HashSet<>(); + addBindings.addAll(bindings); + addBindings.removeAll(knownTopicAclBindings); + if (!addBindings.isEmpty()) { + log.info("Syncing new found {} topic ACL bindings.", addBindings.size()); + targetAdminClient.createAcls(addBindings).values().forEach((k, v) -> v.whenComplete((x, e) -> { + if (e != null) { + log.warn("Could not sync ACL of topic {}.", k.pattern().name(), e); + } + })); + knownTopicAclBindings = bindings; + } else { + log.debug("Not found new topic Acl info, skip sync!"); + } } private static Stream expandTopicDescription(TopicDescription description) { From ee8904644868a58725c78751d5b44906985afd39 Mon Sep 17 00:00:00 2001 From: Deqi Hu Date: Sat, 1 Jul 2023 15:13:14 +0800 Subject: [PATCH 2/5] add unit test --- .../connect/mirror/MirrorSourceConnector.java | 19 +++++++---- .../mirror/MirrorSourceConnectorTest.java | 34 +++++++++++++++++++ 2 files changed, 47 insertions(+), 6 deletions(-) diff --git a/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorSourceConnector.java b/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorSourceConnector.java index 459ddc4a08cfb..16564f33a44f7 100644 --- a/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorSourceConnector.java +++ b/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorSourceConnector.java @@ -104,7 +104,7 @@ public class MirrorSourceConnector extends SourceConnector { private Admin offsetSyncsAdminClient; private volatile boolean useIncrementalAlterConfigs; private AtomicBoolean noAclAuthorizer = new AtomicBoolean(false); - private List knownTopicAclBindings = Collections.emptyList(); + private Set knownTopicAclBindings = Collections.emptySet(); public MirrorSourceConnector() { // nop @@ -582,21 +582,23 @@ void incrementalAlterConfigs(Map topicConfigs) { })); } - private void updateTopicAcls(List bindings) { - Set addBindings = new HashSet<>(); - addBindings.addAll(bindings); + // Visible for testing + int updateTopicAcls(List bindings) { + Set addBindings = new HashSet<>(bindings); addBindings.removeAll(knownTopicAclBindings); + int newBindCount = addBindings.size(); if (!addBindings.isEmpty()) { - log.info("Syncing new found {} topic ACL bindings.", addBindings.size()); + log.info("Syncing new found {} topic ACL bindings.", newBindCount); targetAdminClient.createAcls(addBindings).values().forEach((k, v) -> v.whenComplete((x, e) -> { if (e != null) { log.warn("Could not sync ACL of topic {}.", k.pattern().name(), e); } })); - knownTopicAclBindings = bindings; + knownTopicAclBindings = new HashSet<>(bindings); } else { log.debug("Not found new topic Acl info, skip sync!"); } + return newBindCount; } private static Stream expandTopicDescription(TopicDescription description) { @@ -681,4 +683,9 @@ boolean isCycle(String topic) { String formatRemoteTopic(String topic) { return replicationPolicy.formatRemoteTopic(sourceAndTarget.source(), topic); } + + // Visible for testing + Set knownTopicAclBindings() { + return knownTopicAclBindings; + } } diff --git a/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/MirrorSourceConnectorTest.java b/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/MirrorSourceConnectorTest.java index 5e626679e9df7..1a15a1dae9225 100644 --- a/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/MirrorSourceConnectorTest.java +++ b/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/MirrorSourceConnectorTest.java @@ -18,6 +18,7 @@ import org.apache.kafka.clients.admin.AlterConfigOp; import org.apache.kafka.clients.admin.Admin; +import org.apache.kafka.clients.admin.CreateAclsResult; import org.apache.kafka.clients.admin.DescribeAclsResult; import org.apache.kafka.common.KafkaFuture; import org.apache.kafka.common.TopicPartition; @@ -55,6 +56,7 @@ import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; +import static org.mockito.ArgumentMatchers.anySet; import static org.mockito.ArgumentMatchers.isA; import static org.mockito.Mockito.any; import static org.mockito.Mockito.doAnswer; @@ -74,6 +76,7 @@ import java.util.Collection; import java.util.Collections; import java.util.HashMap; +import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Optional; @@ -683,4 +686,35 @@ private Optional validateProperty(String name, Map assertNotNull(result, "Connector should not have record null config value for '" + name + "' property"); return Optional.of(result); } + + @Test + public void testUpdateIncrementTopicAcls() { + Admin sourceAdmin = mock(Admin.class); + Admin targetAdmin = mock(Admin.class); + MirrorSourceConnector connector = new MirrorSourceConnector(sourceAdmin, targetAdmin); + + List filteredBindings = new ArrayList<>(); + AclBinding binding1 = mock(AclBinding.class); + AclBinding binding2 = mock(AclBinding.class); + filteredBindings.add(binding1); + filteredBindings.add(binding2); + doReturn(mock(CreateAclsResult.class)).when(targetAdmin).createAcls(anySet()); + + // First topic acl info update when starting `syncTopicAcls` thread + int newAddCount = connector.updateTopicAcls(filteredBindings); + assertEquals(connector.knownTopicAclBindings(), new HashSet<>(filteredBindings)); + assertTrue(newAddCount == filteredBindings.size()); + + List newAddBindings = new ArrayList<>(); + AclBinding binding3 = mock(AclBinding.class); + AclBinding binding4 = mock(AclBinding.class); + newAddBindings.add(binding3); + newAddBindings.add(binding4); + filteredBindings.addAll(newAddBindings); + + // The next increment topic acl info update + newAddCount = connector.updateTopicAcls(filteredBindings); + assertEquals(connector.knownTopicAclBindings(), new HashSet<>(filteredBindings)); + assertTrue(newAddCount == newAddBindings.size()); + } } From 373579ea0cb98016fa808047cc702e86bdca01cf Mon Sep 17 00:00:00 2001 From: Deqi Hu Date: Tue, 11 Jul 2023 12:12:21 +0800 Subject: [PATCH 3/5] fix --- .../kafka/clients/admin/CreateAclsResult.java | 3 +- .../connect/mirror/MirrorSourceConnector.java | 5 ++- .../mirror/MirrorSourceConnectorTest.java | 33 +++++++++++++++++-- 3 files changed, 37 insertions(+), 4 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/CreateAclsResult.java b/clients/src/main/java/org/apache/kafka/clients/admin/CreateAclsResult.java index 6e69554635efc..efe1e08a9e407 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/CreateAclsResult.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/CreateAclsResult.java @@ -33,7 +33,8 @@ public class CreateAclsResult { private final Map> futures; - CreateAclsResult(Map> futures) { + // Visible for testing + public CreateAclsResult(Map> futures) { this.futures = futures; } diff --git a/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorSourceConnector.java b/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorSourceConnector.java index 16564f33a44f7..f8ae0546701af 100644 --- a/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorSourceConnector.java +++ b/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorSourceConnector.java @@ -585,6 +585,7 @@ void incrementalAlterConfigs(Map topicConfigs) { // Visible for testing int updateTopicAcls(List bindings) { Set addBindings = new HashSet<>(bindings); + Set failedBindings = new HashSet<>(); addBindings.removeAll(knownTopicAclBindings); int newBindCount = addBindings.size(); if (!addBindings.isEmpty()) { @@ -592,13 +593,15 @@ int updateTopicAcls(List bindings) { targetAdminClient.createAcls(addBindings).values().forEach((k, v) -> v.whenComplete((x, e) -> { if (e != null) { log.warn("Could not sync ACL of topic {}.", k.pattern().name(), e); + failedBindings.add(k); } })); + bindings.removeAll(failedBindings); knownTopicAclBindings = new HashSet<>(bindings); } else { log.debug("Not found new topic Acl info, skip sync!"); } - return newBindCount; + return newBindCount - failedBindings.size(); } private static Stream expandTopicDescription(TopicDescription description) { diff --git a/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/MirrorSourceConnectorTest.java b/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/MirrorSourceConnectorTest.java index 1a15a1dae9225..f7bf63334656d 100644 --- a/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/MirrorSourceConnectorTest.java +++ b/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/MirrorSourceConnectorTest.java @@ -28,8 +28,10 @@ import org.apache.kafka.common.acl.AclPermissionType; import org.apache.kafka.common.config.ConfigResource; import org.apache.kafka.common.config.ConfigValue; +import org.apache.kafka.common.errors.ApiException; import org.apache.kafka.common.errors.UnsupportedVersionException; import org.apache.kafka.common.errors.SecurityDisabledException; +import org.apache.kafka.common.internals.KafkaFutureImpl; import org.apache.kafka.common.resource.PatternType; import org.apache.kafka.common.resource.ResourcePattern; import org.apache.kafka.common.resource.ResourceType; @@ -42,6 +44,7 @@ import org.apache.kafka.connect.errors.ConnectException; import org.apache.kafka.connect.source.ExactlyOnceSupport; import org.junit.jupiter.api.Test; +import org.mockito.ArgumentMatchers; import static org.apache.kafka.clients.admin.AdminClientTestUtils.alterConfigsResult; import static org.apache.kafka.clients.consumer.ConsumerConfig.ISOLATION_LEVEL_CONFIG; @@ -80,6 +83,7 @@ import java.util.List; import java.util.Map; import java.util.Optional; +import java.util.Set; import java.util.concurrent.ExecutionException; import java.util.stream.Collectors; @@ -703,7 +707,7 @@ public void testUpdateIncrementTopicAcls() { // First topic acl info update when starting `syncTopicAcls` thread int newAddCount = connector.updateTopicAcls(filteredBindings); assertEquals(connector.knownTopicAclBindings(), new HashSet<>(filteredBindings)); - assertTrue(newAddCount == filteredBindings.size()); + assertEquals(filteredBindings.size(), newAddCount); List newAddBindings = new ArrayList<>(); AclBinding binding3 = mock(AclBinding.class); @@ -715,6 +719,31 @@ public void testUpdateIncrementTopicAcls() { // The next increment topic acl info update newAddCount = connector.updateTopicAcls(filteredBindings); assertEquals(connector.knownTopicAclBindings(), new HashSet<>(filteredBindings)); - assertTrue(newAddCount == newAddBindings.size()); + assertEquals(newAddBindings.size(), newAddCount); + + // The next increment topic acl info update, contains failed create + List newAddFailedBindings = new ArrayList<>(); + AclBinding binding5 = mock(AclBinding.class); + AclBinding binding6 = mock(AclBinding.class); + newAddFailedBindings.add(binding5); + newAddFailedBindings.add(binding6); + filteredBindings.addAll(newAddFailedBindings); + + Map> futures = new HashMap<>(); + KafkaFutureImpl futureForBinding5 = new KafkaFutureImpl<>(); + KafkaFutureImpl futureForBinding6 = new KafkaFutureImpl<>(); + futureForBinding5.complete(null); + futureForBinding6.completeExceptionally(new ApiException("mock create acl failure.")); + futures.put(binding5, futureForBinding5); + futures.put(binding6, futureForBinding6); + CreateAclsResult mockCreateAclsResult = new CreateAclsResult(new HashMap<>(futures)); + doReturn(new ResourcePattern(ResourceType.TOPIC, "topic6", PatternType.LITERAL)).when(binding6).pattern(); + doReturn(mockCreateAclsResult).when(targetAdmin).createAcls(ArgumentMatchers.eq(new HashSet<>(newAddFailedBindings))); + + int newAddSuccessCount = connector.updateTopicAcls(filteredBindings); + Set bindingToVerify = new HashSet<>(filteredBindings); + bindingToVerify.remove(binding6); + assertEquals(connector.knownTopicAclBindings(), bindingToVerify); + assertEquals(1, newAddSuccessCount); } } From dd25b78327b34fa119763558fe4c7ed93fc0b8c8 Mon Sep 17 00:00:00 2001 From: Deqi Hu Date: Tue, 11 Jul 2023 19:03:40 +0800 Subject: [PATCH 4/5] minor --- .../org/apache/kafka/connect/mirror/MirrorSourceConnector.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorSourceConnector.java b/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorSourceConnector.java index f8ae0546701af..5f1a077eb4dcb 100644 --- a/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorSourceConnector.java +++ b/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorSourceConnector.java @@ -596,8 +596,8 @@ int updateTopicAcls(List bindings) { failedBindings.add(k); } })); - bindings.removeAll(failedBindings); knownTopicAclBindings = new HashSet<>(bindings); + knownTopicAclBindings.removeAll(failedBindings); } else { log.debug("Not found new topic Acl info, skip sync!"); } From 585cbcdddd9d987f06809ab71fc99c0922b5d48b Mon Sep 17 00:00:00 2001 From: Deqi Hu Date: Wed, 12 Jul 2023 16:23:01 +0800 Subject: [PATCH 5/5] fix bug --- .../connect/mirror/MirrorSourceConnector.java | 9 ++++++++- .../mirror/MirrorSourceConnectorTest.java | 18 +++++++++++++----- 2 files changed, 21 insertions(+), 6 deletions(-) diff --git a/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorSourceConnector.java b/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorSourceConnector.java index 5f1a077eb4dcb..315ee6e414f89 100644 --- a/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorSourceConnector.java +++ b/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorSourceConnector.java @@ -20,10 +20,12 @@ import java.util.Locale; import java.util.Map.Entry; +import java.util.concurrent.Future; import java.util.concurrent.atomic.AtomicReference; import org.apache.kafka.clients.admin.Admin; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.common.IsolationLevel; +import org.apache.kafka.common.KafkaFuture; import org.apache.kafka.common.config.ConfigValue; import org.apache.kafka.common.errors.SecurityDisabledException; import org.apache.kafka.connect.connector.Task; @@ -74,6 +76,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import static org.apache.kafka.common.utils.Utils.sleep; import static org.apache.kafka.connect.mirror.MirrorSourceConfig.SYNC_TOPIC_ACLS_ENABLED; /** Replicate data, configuration, and ACLs between clusters. @@ -590,12 +593,16 @@ int updateTopicAcls(List bindings) { int newBindCount = addBindings.size(); if (!addBindings.isEmpty()) { log.info("Syncing new found {} topic ACL bindings.", newBindCount); - targetAdminClient.createAcls(addBindings).values().forEach((k, v) -> v.whenComplete((x, e) -> { + Map> futureMap = targetAdminClient.createAcls(addBindings).values(); + futureMap.forEach((k, v) -> v.whenComplete((x, e) -> { if (e != null) { log.warn("Could not sync ACL of topic {}.", k.pattern().name(), e); failedBindings.add(k); } })); + while (!futureMap.values().stream().allMatch(Future::isDone)) { + sleep(1000); + } knownTopicAclBindings = new HashSet<>(bindings); knownTopicAclBindings.removeAll(failedBindings); } else { diff --git a/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/MirrorSourceConnectorTest.java b/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/MirrorSourceConnectorTest.java index f7bf63334656d..793a32dc88caa 100644 --- a/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/MirrorSourceConnectorTest.java +++ b/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/MirrorSourceConnectorTest.java @@ -48,6 +48,7 @@ import static org.apache.kafka.clients.admin.AdminClientTestUtils.alterConfigsResult; import static org.apache.kafka.clients.consumer.ConsumerConfig.ISOLATION_LEVEL_CONFIG; +import static org.apache.kafka.common.utils.Utils.sleep; import static org.apache.kafka.connect.mirror.MirrorConnectorConfig.CONSUMER_CLIENT_PREFIX; import static org.apache.kafka.connect.mirror.MirrorConnectorConfig.SOURCE_PREFIX; import static org.apache.kafka.connect.mirror.MirrorSourceConfig.OFFSET_LAG_MAX; @@ -85,6 +86,8 @@ import java.util.Optional; import java.util.Set; import java.util.concurrent.ExecutionException; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; import java.util.stream.Collectors; public class MirrorSourceConnectorTest { @@ -706,7 +709,7 @@ public void testUpdateIncrementTopicAcls() { // First topic acl info update when starting `syncTopicAcls` thread int newAddCount = connector.updateTopicAcls(filteredBindings); - assertEquals(connector.knownTopicAclBindings(), new HashSet<>(filteredBindings)); + assertEquals(new HashSet<>(filteredBindings), connector.knownTopicAclBindings()); assertEquals(filteredBindings.size(), newAddCount); List newAddBindings = new ArrayList<>(); @@ -718,7 +721,7 @@ public void testUpdateIncrementTopicAcls() { // The next increment topic acl info update newAddCount = connector.updateTopicAcls(filteredBindings); - assertEquals(connector.knownTopicAclBindings(), new HashSet<>(filteredBindings)); + assertEquals(new HashSet<>(filteredBindings), connector.knownTopicAclBindings()); assertEquals(newAddBindings.size(), newAddCount); // The next increment topic acl info update, contains failed create @@ -732,8 +735,13 @@ public void testUpdateIncrementTopicAcls() { Map> futures = new HashMap<>(); KafkaFutureImpl futureForBinding5 = new KafkaFutureImpl<>(); KafkaFutureImpl futureForBinding6 = new KafkaFutureImpl<>(); - futureForBinding5.complete(null); - futureForBinding6.completeExceptionally(new ApiException("mock create acl failure.")); + ExecutorService singleThread = Executors.newSingleThreadExecutor(); + // Delayed completion of `createAclRequest` for simulating actual scenarios + singleThread.submit(() -> { + sleep(5000); + futureForBinding5.complete(null); + futureForBinding6.completeExceptionally(new ApiException("mock create acl failure.")); + }); futures.put(binding5, futureForBinding5); futures.put(binding6, futureForBinding6); CreateAclsResult mockCreateAclsResult = new CreateAclsResult(new HashMap<>(futures)); @@ -743,7 +751,7 @@ public void testUpdateIncrementTopicAcls() { int newAddSuccessCount = connector.updateTopicAcls(filteredBindings); Set bindingToVerify = new HashSet<>(filteredBindings); bindingToVerify.remove(binding6); - assertEquals(connector.knownTopicAclBindings(), bindingToVerify); + assertEquals(bindingToVerify, connector.knownTopicAclBindings()); assertEquals(1, newAddSuccessCount); } }