From ebcdae1bcee6fc149bbe3c4a9220e2e959e59fbb Mon Sep 17 00:00:00 2001 From: Deqi Hu Date: Sat, 1 Jul 2023 17:00:10 +0800 Subject: [PATCH] KAFKA-15139:Optimize the performance of `Set.removeAll(List)` in `MirrorCheckpointConnector` --- .../mirror/MirrorCheckpointConnector.java | 22 +++++++------- .../mirror/MirrorCheckpointConnectorTest.java | 30 ++++++++++--------- 2 files changed, 26 insertions(+), 26 deletions(-) diff --git a/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorCheckpointConnector.java b/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorCheckpointConnector.java index 80154e8211f79..07e7b49a44ea4 100644 --- a/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorCheckpointConnector.java +++ b/connect/mirror/src/main/java/org/apache/kafka/connect/mirror/MirrorCheckpointConnector.java @@ -29,10 +29,10 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.util.ArrayList; import java.util.Collection; import java.util.Collections; import java.util.HashSet; -import java.util.LinkedList; import java.util.List; import java.util.Map; import java.util.Set; @@ -55,14 +55,14 @@ public class MirrorCheckpointConnector extends SourceConnector { private Admin sourceAdminClient; private Admin targetAdminClient; private SourceAndTarget sourceAndTarget; - private List knownConsumerGroups = Collections.emptyList(); + private Set knownConsumerGroups = Collections.emptySet(); public MirrorCheckpointConnector() { // nop } // visible for testing - MirrorCheckpointConnector(List knownConsumerGroups, MirrorCheckpointConfig config) { + MirrorCheckpointConnector(Set knownConsumerGroups, MirrorCheckpointConfig config) { this.knownConsumerGroups = knownConsumerGroups; this.config = config; } @@ -116,7 +116,7 @@ public List> taskConfigs(int maxTasks) { return Collections.emptyList(); } int numTasks = Math.min(maxTasks, knownConsumerGroups.size()); - List> groupsPartitioned = ConnectorUtils.groupPartitions(knownConsumerGroups, numTasks); + List> groupsPartitioned = ConnectorUtils.groupPartitions(new ArrayList<>(knownConsumerGroups), numTasks); return IntStream.range(0, numTasks) .mapToObj(i -> config.taskConfigForConsumerGroups(groupsPartitioned.get(i), i)) .collect(Collectors.toList()); @@ -134,12 +134,10 @@ public String version() { private void refreshConsumerGroups() throws InterruptedException, ExecutionException { - List consumerGroups = findConsumerGroups(); - Set newConsumerGroups = new HashSet<>(); - newConsumerGroups.addAll(consumerGroups); + Set consumerGroups = findConsumerGroups(); + Set newConsumerGroups = new HashSet<>(consumerGroups); newConsumerGroups.removeAll(knownConsumerGroups); - Set deadConsumerGroups = new HashSet<>(); - deadConsumerGroups.addAll(knownConsumerGroups); + Set deadConsumerGroups = new HashSet<>(knownConsumerGroups); deadConsumerGroups.removeAll(consumerGroups); if (!newConsumerGroups.isEmpty() || !deadConsumerGroups.isEmpty()) { log.info("Found {} consumer groups for {}. {} are new. {} were removed. Previously had {}.", @@ -156,15 +154,15 @@ private void loadInitialConsumerGroups() knownConsumerGroups = findConsumerGroups(); } - List findConsumerGroups() + Set findConsumerGroups() throws InterruptedException, ExecutionException { List filteredGroups = listConsumerGroups().stream() .map(ConsumerGroupListing::groupId) .filter(this::shouldReplicateByGroupFilter) .collect(Collectors.toList()); - List checkpointGroups = new LinkedList<>(); - List irrelevantGroups = new LinkedList<>(); + Set checkpointGroups = new HashSet<>(); + Set irrelevantGroups = new HashSet<>(); for (String group : filteredGroups) { Set consumedTopics = listConsumerGroupOffsets(group).keySet().stream() diff --git a/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/MirrorCheckpointConnectorTest.java b/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/MirrorCheckpointConnectorTest.java index be64567d78c05..5481802915090 100644 --- a/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/MirrorCheckpointConnectorTest.java +++ b/connect/mirror/src/test/java/org/apache/kafka/connect/mirror/MirrorCheckpointConnectorTest.java @@ -21,7 +21,6 @@ import org.apache.kafka.common.TopicPartition; import org.junit.jupiter.api.Test; -import java.util.ArrayList; import java.util.Arrays; import java.util.Collection; import java.util.Collections; @@ -49,7 +48,7 @@ public void testMirrorCheckpointConnectorDisabled() { MirrorCheckpointConfig config = new MirrorCheckpointConfig( makeProps("emit.checkpoints.enabled", "false")); - List knownConsumerGroups = new ArrayList<>(); + Set knownConsumerGroups = new HashSet<>(); knownConsumerGroups.add(CONSUMER_GROUP); // MirrorCheckpointConnector as minimum to run taskConfig() MirrorCheckpointConnector connector = new MirrorCheckpointConnector(knownConsumerGroups, @@ -65,7 +64,7 @@ public void testMirrorCheckpointConnectorEnabled() { MirrorCheckpointConfig config = new MirrorCheckpointConfig( makeProps("emit.checkpoints.enabled", "true")); - List knownConsumerGroups = new ArrayList<>(); + Set knownConsumerGroups = new HashSet<>(); knownConsumerGroups.add(CONSUMER_GROUP); // MirrorCheckpointConnector as minimum to run taskConfig() MirrorCheckpointConnector connector = new MirrorCheckpointConnector(knownConsumerGroups, @@ -81,7 +80,7 @@ public void testMirrorCheckpointConnectorEnabled() { @Test public void testNoConsumerGroup() { MirrorCheckpointConfig config = new MirrorCheckpointConfig(makeProps()); - MirrorCheckpointConnector connector = new MirrorCheckpointConnector(new ArrayList<>(), config); + MirrorCheckpointConnector connector = new MirrorCheckpointConnector(new HashSet<>(), config); List> output = connector.taskConfigs(1); // expect no task will be created assertEquals(0, output.size(), "ConsumerGroup shouldn't exist"); @@ -92,7 +91,7 @@ public void testReplicationDisabled() { // disable the replication MirrorCheckpointConfig config = new MirrorCheckpointConfig(makeProps("enabled", "false")); - List knownConsumerGroups = new ArrayList<>(); + Set knownConsumerGroups = new HashSet<>(); knownConsumerGroups.add(CONSUMER_GROUP); // MirrorCheckpointConnector as minimum to run taskConfig() MirrorCheckpointConnector connector = new MirrorCheckpointConnector(knownConsumerGroups, config); @@ -106,7 +105,7 @@ public void testReplicationEnabled() { // enable the replication MirrorCheckpointConfig config = new MirrorCheckpointConfig(makeProps("enabled", "true")); - List knownConsumerGroups = new ArrayList<>(); + Set knownConsumerGroups = new HashSet<>(); knownConsumerGroups.add(CONSUMER_GROUP); // MirrorCheckpointConnector as minimum to run taskConfig() MirrorCheckpointConnector connector = new MirrorCheckpointConnector(knownConsumerGroups, config); @@ -120,7 +119,7 @@ public void testReplicationEnabled() { @Test public void testFindConsumerGroups() throws Exception { MirrorCheckpointConfig config = new MirrorCheckpointConfig(makeProps()); - MirrorCheckpointConnector connector = new MirrorCheckpointConnector(Collections.emptyList(), config); + MirrorCheckpointConnector connector = new MirrorCheckpointConnector(Collections.emptySet(), config); connector = spy(connector); Collection groups = Arrays.asList( @@ -132,21 +131,21 @@ public void testFindConsumerGroups() throws Exception { doReturn(true).when(connector).shouldReplicateByTopicFilter(anyString()); doReturn(true).when(connector).shouldReplicateByGroupFilter(anyString()); doReturn(offsets).when(connector).listConsumerGroupOffsets(anyString()); - List groupFound = connector.findConsumerGroups(); + Set groupFound = connector.findConsumerGroups(); Set expectedGroups = groups.stream().map(ConsumerGroupListing::groupId).collect(Collectors.toSet()); - assertEquals(expectedGroups, new HashSet<>(groupFound), + assertEquals(expectedGroups, groupFound, "Expected groups are not the same as findConsumerGroups"); doReturn(false).when(connector).shouldReplicateByTopicFilter(anyString()); - List topicFilterGroupFound = connector.findConsumerGroups(); - assertEquals(Collections.emptyList(), topicFilterGroupFound); + Set topicFilterGroupFound = connector.findConsumerGroups(); + assertEquals(Collections.emptySet(), topicFilterGroupFound); } @Test public void testFindConsumerGroupsInCommonScenarios() throws Exception { MirrorCheckpointConfig config = new MirrorCheckpointConfig(makeProps()); - MirrorCheckpointConnector connector = new MirrorCheckpointConnector(Collections.emptyList(), config); + MirrorCheckpointConnector connector = new MirrorCheckpointConnector(Collections.emptySet(), config); connector = spy(connector); Collection groups = Arrays.asList( @@ -177,8 +176,11 @@ public void testFindConsumerGroupsInCommonScenarios() throws Exception { doReturn(offsetsForGroup3).when(connector).listConsumerGroupOffsets("g3"); doReturn(offsetsForGroup4).when(connector).listConsumerGroupOffsets("g4"); - List groupFound = connector.findConsumerGroups(); - assertEquals(groupFound, Arrays.asList("g1", "g2")); + Set groupFound = connector.findConsumerGroups(); + Set verifiedSet = new HashSet<>(); + verifiedSet.add("g1"); + verifiedSet.add("g2"); + assertEquals(groupFound, verifiedSet); } }