From 32cbadcac733148fcbfdc3e7cea0270a75b49686 Mon Sep 17 00:00:00 2001 From: Dhruvil Shah Date: Thu, 25 Feb 2021 21:14:42 -0800 Subject: [PATCH 1/7] KAFKA-12254: Ensure MM2 creates topics with source topic configs --- .../connect/mirror/MirrorSourceConnector.java | 30 ++++++++++++++----- .../mirror/MirrorSourceConnectorTest.java | 22 ++++++++++++-- 2 files changed, 43 insertions(+), 9 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 3e7f0c7bea77a..be37343b209de 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 @@ -309,18 +309,29 @@ private void createOffsetSyncsTopic() { // visible for testing void computeAndCreateTopicPartitions() throws InterruptedException, ExecutionException { - Map partitionCounts = knownSourceTopicPartitions.stream() + Map topicToParitionCount = knownSourceTopicPartitions.stream() .collect(Collectors.groupingBy(TopicPartition::topic, Collectors.counting())).entrySet().stream() .collect(Collectors.toMap(x -> formatRemoteTopic(x.getKey()), Entry::getValue)); Set knownTargetTopics = toTopics(knownTargetTopicPartitions); - List newTopics = partitionCounts.entrySet().stream() - .filter(x -> !knownTargetTopics.contains(x.getKey())) - .map(x -> new NewTopic(x.getKey(), x.getValue().intValue(), (short) replicationFactor)) + Set topicsToCreate = topicToParitionCount.keySet().stream() + .filter(topic -> !knownTargetTopics.contains(topic)) + .collect(Collectors.toSet()); + + Map configs = describeTopicConfigs(topicsToCreate); + List newTopics = topicToParitionCount.entrySet().stream() + .filter(topicAndPartitionCount -> !knownTargetTopics.contains(topicAndPartitionCount.getKey())) + .map(topicAndPartitionCount -> { + String topicName = topicAndPartitionCount.getKey(); + Map topicConfigs = configToMap(configs.get(topicName)); + return new NewTopic(topicName, topicAndPartitionCount.getValue().intValue(), (short) replicationFactor) + .configs(topicConfigs); + }) .collect(Collectors.toList()); - Map newPartitions = partitionCounts.entrySet().stream() + + Map newPartitions = topicToParitionCount.entrySet().stream() .filter(x -> knownTargetTopics.contains(x.getKey())) .collect(Collectors.toMap(Entry::getKey, x -> NewPartitions.increaseTo(x.getValue().intValue()))); - createTopicPartitions(partitionCounts, newTopics, newPartitions); + createTopicPartitions(topicToParitionCount, newTopics, newPartitions); } // visible for testing @@ -359,6 +370,11 @@ private static Collection describeTopics(AdminClient adminClie return adminClient.describeTopics(topics).all().get().values(); } + static Map configToMap(Config config) { + return config.entries().stream() + .collect(Collectors.toMap(configEntry -> configEntry.name(), configEntry -> configEntry.value())); + } + @SuppressWarnings("deprecation") // use deprecated alterConfigs API for broker compatibility back to 0.11.0 private void updateTopicConfigs(Map topicConfigs) @@ -390,7 +406,7 @@ private static Stream expandTopicDescription(TopicDescription de .map(x -> new TopicPartition(topic, x.partition())); } - private Map describeTopicConfigs(Set topics) + Map describeTopicConfigs(Set topics) throws InterruptedException, ExecutionException { Set resources = topics.stream() .map(x -> new ConfigResource(ConfigResource.Type.TOPIC, x)) 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 a9633915fb979..a14f87cdf92b0 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 @@ -16,6 +16,7 @@ */ package org.apache.kafka.connect.mirror; +import kafka.log.LogConfig; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.acl.AccessControlEntry; import org.apache.kafka.common.acl.AclBinding; @@ -183,9 +184,15 @@ public void testRefreshTopicPartitions() throws Exception { connector.initialize(mock(ConnectorContext.class)); connector = spy(connector); + Config topicConfig = new Config(Arrays.asList( + new ConfigEntry(LogConfig.CleanupPolicyProp(), "compact"), + new ConfigEntry(LogConfig.SegmentBytesProp(), "100"))); + Map configs = Collections.singletonMap("source.topic", topicConfig); + List sourceTopicPartitions = Collections.singletonList(new TopicPartition("topic", 0)); doReturn(sourceTopicPartitions).when(connector).findSourceTopicPartitions(); doReturn(Collections.emptyList()).when(connector).findTargetTopicPartitions(); + doReturn(configs).when(connector).describeTopicConfigs(Collections.singleton("source.topic")); doNothing().when(connector).createTopicPartitions(any(), any(), any()); connector.refreshTopicPartitions(); @@ -194,7 +201,12 @@ public void testRefreshTopicPartitions() throws Exception { Map expectedPartitionCounts = new HashMap<>(); expectedPartitionCounts.put("source.topic", 1L); - List expectedNewTopics = Arrays.asList(new NewTopic("source.topic", 1, (short) 0)); + Map configMap = MirrorSourceConnector.configToMap(topicConfig); + assertEquals(2, configMap.size()); + + List expectedNewTopics = Arrays.asList( + new NewTopic("source.topic", 1, (short) 0) + .configs(configMap)); verify(connector, times(2)).computeAndCreateTopicPartitions(); verify(connector, times(2)).createTopicPartitions( @@ -217,10 +229,17 @@ public void testRefreshTopicPartitionsTopicOnTargetFirst() throws Exception { connector.initialize(mock(ConnectorContext.class)); connector = spy(connector); + Config topicConfig = new Config(Arrays.asList( + new ConfigEntry(LogConfig.CleanupPolicyProp(), "compact"), + new ConfigEntry(LogConfig.SegmentBytesProp(), "100"))); + Map configs = Collections.singletonMap("source.topic", topicConfig); + List sourceTopicPartitions = Collections.emptyList(); List targetTopicPartitions = Collections.singletonList(new TopicPartition("source.topic", 0)); doReturn(sourceTopicPartitions).when(connector).findSourceTopicPartitions(); doReturn(targetTopicPartitions).when(connector).findTargetTopicPartitions(); + doReturn(configs).when(connector).describeTopicConfigs(Collections.singleton("source.topic")); + doReturn(Collections.emptyMap()).when(connector).describeTopicConfigs(Collections.emptySet()); doNothing().when(connector).createTopicPartitions(any(), any(), any()); // partitions appearing on the target cluster should not cause reconfiguration @@ -234,6 +253,5 @@ public void testRefreshTopicPartitionsTopicOnTargetFirst() throws Exception { // when partitions are added to the source cluster, reconfiguration is triggered connector.refreshTopicPartitions(); verify(connector, times(1)).computeAndCreateTopicPartitions(); - } } From caf6c4a0a3611c44bf979e3d803bca361bfd4792 Mon Sep 17 00:00:00 2001 From: Dhruvil Shah Date: Thu, 25 Feb 2021 21:25:08 -0800 Subject: [PATCH 2/7] minor style change --- .../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 be37343b209de..13f59d1718f2b 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 @@ -372,7 +372,7 @@ private static Collection describeTopics(AdminClient adminClie static Map configToMap(Config config) { return config.entries().stream() - .collect(Collectors.toMap(configEntry -> configEntry.name(), configEntry -> configEntry.value())); + .collect(Collectors.toMap(ConfigEntry::name, ConfigEntry::value)); } @SuppressWarnings("deprecation") From 70e7e396c6ed9552a2ffba6fb116dfdd85e7b82b Mon Sep 17 00:00:00 2001 From: Dhruvil Shah Date: Thu, 25 Feb 2021 21:37:35 -0800 Subject: [PATCH 3/7] add comments --- .../kafka/connect/mirror/MirrorSourceConnector.java | 10 +++++++++- 1 file changed, 9 insertions(+), 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 13f59d1718f2b..1e1f806c297de 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 @@ -313,11 +313,18 @@ void computeAndCreateTopicPartitions() .collect(Collectors.groupingBy(TopicPartition::topic, Collectors.counting())).entrySet().stream() .collect(Collectors.toMap(x -> formatRemoteTopic(x.getKey()), Entry::getValue)); Set knownTargetTopics = toTopics(knownTargetTopicPartitions); + + // get the set of topics that are not present on the target Set topicsToCreate = topicToParitionCount.keySet().stream() .filter(topic -> !knownTargetTopics.contains(topic)) .collect(Collectors.toSet()); + if (topicsToCreate.isEmpty()) + return; + // get corresponding topic configurations from the source Map configs = describeTopicConfigs(topicsToCreate); + + // construct collections for new topics and partitions List newTopics = topicToParitionCount.entrySet().stream() .filter(topicAndPartitionCount -> !knownTargetTopics.contains(topicAndPartitionCount.getKey())) .map(topicAndPartitionCount -> { @@ -327,10 +334,11 @@ void computeAndCreateTopicPartitions() .configs(topicConfigs); }) .collect(Collectors.toList()); - Map newPartitions = topicToParitionCount.entrySet().stream() .filter(x -> knownTargetTopics.contains(x.getKey())) .collect(Collectors.toMap(Entry::getKey, x -> NewPartitions.increaseTo(x.getValue().intValue()))); + + // create topic partitions on target createTopicPartitions(topicToParitionCount, newTopics, newPartitions); } From db09806e741f5108dbf5a5f9438ed81aeb81cdbd Mon Sep 17 00:00:00 2001 From: Dhruvil Shah Date: Thu, 25 Feb 2021 23:58:27 -0800 Subject: [PATCH 4/7] describe configs with source topic names --- .../connect/mirror/MirrorSourceConnector.java | 98 ++++++++++++------- .../mirror/MirrorSourceConnectorTest.java | 24 +++-- 2 files changed, 72 insertions(+), 50 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 1e1f806c297de..fbb06a3401a62 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 @@ -49,6 +49,7 @@ import java.util.HashSet; import java.util.Collection; import java.util.Collections; +import java.util.function.Function; import java.util.stream.Collectors; import java.util.stream.Stream; import java.util.concurrent.ExecutionException; @@ -306,59 +307,82 @@ private void createOffsetSyncsTopic() { MirrorUtils.createSinglePartitionCompactedTopic(config.offsetSyncsTopic(), config.offsetSyncsTopicReplicationFactor(), config.sourceAdminConfig()); } - // visible for testing - void computeAndCreateTopicPartitions() - throws InterruptedException, ExecutionException { - Map topicToParitionCount = knownSourceTopicPartitions.stream() - .collect(Collectors.groupingBy(TopicPartition::topic, Collectors.counting())).entrySet().stream() - .collect(Collectors.toMap(x -> formatRemoteTopic(x.getKey()), Entry::getValue)); - Set knownTargetTopics = toTopics(knownTargetTopicPartitions); - - // get the set of topics that are not present on the target - Set topicsToCreate = topicToParitionCount.keySet().stream() - .filter(topic -> !knownTargetTopics.contains(topic)) - .collect(Collectors.toSet()); - if (topicsToCreate.isEmpty()) - return; - - // get corresponding topic configurations from the source - Map configs = describeTopicConfigs(topicsToCreate); - - // construct collections for new topics and partitions - List newTopics = topicToParitionCount.entrySet().stream() - .filter(topicAndPartitionCount -> !knownTargetTopics.contains(topicAndPartitionCount.getKey())) - .map(topicAndPartitionCount -> { - String topicName = topicAndPartitionCount.getKey(); - Map topicConfigs = configToMap(configs.get(topicName)); - return new NewTopic(topicName, topicAndPartitionCount.getValue().intValue(), (short) replicationFactor) - .configs(topicConfigs); - }) - .collect(Collectors.toList()); - Map newPartitions = topicToParitionCount.entrySet().stream() - .filter(x -> knownTargetTopics.contains(x.getKey())) - .collect(Collectors.toMap(Entry::getKey, x -> NewPartitions.increaseTo(x.getValue().intValue()))); + void computeAndCreateTopicPartitions() throws ExecutionException, InterruptedException { + // get source and target topics with respective partition counts + Map sourceTopicToPartitionCounts = knownSourceTopicPartitions.stream() + .collect(Collectors.groupingBy(TopicPartition::topic, Collectors.counting())).entrySet().stream() + .collect(Collectors.toMap(Entry::getKey, Entry::getValue)); + Map targetTopicToPartitionCounts = knownTargetTopicPartitions.stream() + .collect(Collectors.groupingBy(TopicPartition::topic, Collectors.counting())).entrySet().stream() + .collect(Collectors.toMap(Entry::getKey, Entry::getValue)); + + Set knownSourceTopics = sourceTopicToPartitionCounts.keySet(); + Set knownTargetTopics = targetTopicToPartitionCounts.keySet(); + + // compute existing and new source topics + Map> partitionedSourceTopics = knownSourceTopics.stream() + .collect(Collectors.partitioningBy(sourceTopic -> knownTargetTopics.contains(formatRemoteTopic(sourceTopic)), + Collectors.toSet())); + Set existingSourceTopics = partitionedSourceTopics.get(true); + Set newSourceTopics = partitionedSourceTopics.get(false); + + // create new topics + if (!newSourceTopics.isEmpty()) + createNewTopics(newSourceTopics, sourceTopicToPartitionCounts); + + // compute topics with new partitions + Map sourceTopicsWithNewPartitions = existingSourceTopics.stream() + .filter(sourceTopic -> { + String targetTopic = formatRemoteTopic(sourceTopic); + return sourceTopicToPartitionCounts.get(sourceTopic) > targetTopicToPartitionCounts.get(targetTopic); + }) + .collect(Collectors.toMap(Function.identity(), sourceTopicToPartitionCounts::get)); + + // create new partitions + if (!sourceTopicsWithNewPartitions.isEmpty()) { + Map newTargetPartitions = sourceTopicsWithNewPartitions.entrySet().stream() + .collect(Collectors.toMap(sourceTopicAndPartitionCount -> formatRemoteTopic(sourceTopicAndPartitionCount.getKey()), + sourceTopicAndPartitionCount -> NewPartitions.increaseTo(sourceTopicAndPartitionCount.getValue().intValue()))); + createNewPartitions(newTargetPartitions); + } + } - // create topic partitions on target - createTopicPartitions(topicToParitionCount, newTopics, newPartitions); + private void createNewTopics(Set newSourceTopics, Map sourceTopicToPartitionCounts) + throws ExecutionException, InterruptedException { + Map sourceTopicToConfig = describeTopicConfigs(newSourceTopics); + List newTopics = newSourceTopics.stream() + .map(sourceTopic -> { + String remoteTopic = formatRemoteTopic(sourceTopic); + int partitionCount = sourceTopicToPartitionCounts.get(sourceTopic).intValue(); + Map configs = configToMap(sourceTopicToConfig.get(sourceTopic)); + return new NewTopic(remoteTopic, partitionCount, (short) replicationFactor) + .configs(configs); + }) + .collect(Collectors.toList()); + createNewTopics(newTopics); } // visible for testing - void createTopicPartitions(Map partitionCounts, List newTopics, - Map newPartitions) { + void createNewTopics(List newTopics) { + Map newTopicMap = newTopics.stream() + .collect(Collectors.toMap(NewTopic::name, Function.identity())); targetAdminClient.createTopics(newTopics, new CreateTopicsOptions()).values().forEach((k, v) -> v.whenComplete((x, e) -> { if (e != null) { log.warn("Could not create topic {}.", k, e); } else { - log.info("Created remote topic {} with {} partitions.", k, partitionCounts.get(k)); + log.info("Created remote topic {} with {} partitions.", k, newTopicMap.get(k).numPartitions()); } })); + } + + void createNewPartitions(Map newPartitions) { targetAdminClient.createPartitions(newPartitions).values().forEach((k, v) -> v.whenComplete((x, e) -> { if (e instanceof InvalidPartitionsException) { // swallow, this is normal } else if (e != null) { log.warn("Could not create topic-partitions for {}.", k, e); } else { - log.info("Increased size of {} to {} partitions.", k, partitionCounts.get(k)); + log.info("Increased size of {} to {} partitions.", k, newPartitions.get(k).totalCount()); } })); } 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 a14f87cdf92b0..66a7429e6d482 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 @@ -16,7 +16,6 @@ */ package org.apache.kafka.connect.mirror; -import kafka.log.LogConfig; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.acl.AccessControlEntry; import org.apache.kafka.common.acl.AclBinding; @@ -185,15 +184,15 @@ public void testRefreshTopicPartitions() throws Exception { connector = spy(connector); Config topicConfig = new Config(Arrays.asList( - new ConfigEntry(LogConfig.CleanupPolicyProp(), "compact"), - new ConfigEntry(LogConfig.SegmentBytesProp(), "100"))); - Map configs = Collections.singletonMap("source.topic", topicConfig); + new ConfigEntry("cleanup.policy", "compact"), + new ConfigEntry("segment.bytes", "100"))); + Map configs = Collections.singletonMap("topic", topicConfig); List sourceTopicPartitions = Collections.singletonList(new TopicPartition("topic", 0)); doReturn(sourceTopicPartitions).when(connector).findSourceTopicPartitions(); doReturn(Collections.emptyList()).when(connector).findTargetTopicPartitions(); - doReturn(configs).when(connector).describeTopicConfigs(Collections.singleton("source.topic")); - doNothing().when(connector).createTopicPartitions(any(), any(), any()); + doReturn(configs).when(connector).describeTopicConfigs(Collections.singleton("topic")); + doNothing().when(connector).createNewTopics(any()); connector.refreshTopicPartitions(); // if target topic is not created, refreshTopicPartitions() will call createTopicPartitions() again @@ -209,10 +208,8 @@ public void testRefreshTopicPartitions() throws Exception { .configs(configMap)); verify(connector, times(2)).computeAndCreateTopicPartitions(); - verify(connector, times(2)).createTopicPartitions( - eq(expectedPartitionCounts), - eq(expectedNewTopics), - eq(Collections.emptyMap())); + verify(connector, times(2)).createNewTopics(eq(expectedNewTopics)); + verify(connector, times(0)).createNewPartitions(any()); List targetTopicPartitions = Collections.singletonList(new TopicPartition("source.topic", 0)); doReturn(targetTopicPartitions).when(connector).findTargetTopicPartitions(); @@ -230,8 +227,8 @@ public void testRefreshTopicPartitionsTopicOnTargetFirst() throws Exception { connector = spy(connector); Config topicConfig = new Config(Arrays.asList( - new ConfigEntry(LogConfig.CleanupPolicyProp(), "compact"), - new ConfigEntry(LogConfig.SegmentBytesProp(), "100"))); + new ConfigEntry("cleanup.policy", "compact"), + new ConfigEntry("segment.bytes", "100"))); Map configs = Collections.singletonMap("source.topic", topicConfig); List sourceTopicPartitions = Collections.emptyList(); @@ -240,7 +237,8 @@ public void testRefreshTopicPartitionsTopicOnTargetFirst() throws Exception { doReturn(targetTopicPartitions).when(connector).findTargetTopicPartitions(); doReturn(configs).when(connector).describeTopicConfigs(Collections.singleton("source.topic")); doReturn(Collections.emptyMap()).when(connector).describeTopicConfigs(Collections.emptySet()); - doNothing().when(connector).createTopicPartitions(any(), any(), any()); + doNothing().when(connector).createNewTopics(any()); + doNothing().when(connector).createNewPartitions(any()); // partitions appearing on the target cluster should not cause reconfiguration connector.refreshTopicPartitions(); From 6bce2666e57af8717bc4ea0385668ce18d7d8790 Mon Sep 17 00:00:00 2001 From: Dhruvil Shah Date: Fri, 26 Feb 2021 00:23:46 -0800 Subject: [PATCH 5/7] cache source to remote topic mapping --- .../kafka/connect/mirror/MirrorSourceConnector.java | 8 +++++--- 1 file changed, 5 insertions(+), 3 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 fbb06a3401a62..015d535fff679 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 @@ -318,10 +318,12 @@ void computeAndCreateTopicPartitions() throws ExecutionException, InterruptedExc Set knownSourceTopics = sourceTopicToPartitionCounts.keySet(); Set knownTargetTopics = targetTopicToPartitionCounts.keySet(); + Map sourceToRemoteTopic = knownSourceTopics.stream() + .collect(Collectors.toMap(Function.identity(), sourceTopic -> formatRemoteTopic(sourceTopic))); // compute existing and new source topics Map> partitionedSourceTopics = knownSourceTopics.stream() - .collect(Collectors.partitioningBy(sourceTopic -> knownTargetTopics.contains(formatRemoteTopic(sourceTopic)), + .collect(Collectors.partitioningBy(sourceTopic -> knownTargetTopics.contains(sourceToRemoteTopic.get(sourceTopic)), Collectors.toSet())); Set existingSourceTopics = partitionedSourceTopics.get(true); Set newSourceTopics = partitionedSourceTopics.get(false); @@ -333,7 +335,7 @@ void computeAndCreateTopicPartitions() throws ExecutionException, InterruptedExc // compute topics with new partitions Map sourceTopicsWithNewPartitions = existingSourceTopics.stream() .filter(sourceTopic -> { - String targetTopic = formatRemoteTopic(sourceTopic); + String targetTopic = sourceToRemoteTopic.get(sourceTopic); return sourceTopicToPartitionCounts.get(sourceTopic) > targetTopicToPartitionCounts.get(targetTopic); }) .collect(Collectors.toMap(Function.identity(), sourceTopicToPartitionCounts::get)); @@ -341,7 +343,7 @@ void computeAndCreateTopicPartitions() throws ExecutionException, InterruptedExc // create new partitions if (!sourceTopicsWithNewPartitions.isEmpty()) { Map newTargetPartitions = sourceTopicsWithNewPartitions.entrySet().stream() - .collect(Collectors.toMap(sourceTopicAndPartitionCount -> formatRemoteTopic(sourceTopicAndPartitionCount.getKey()), + .collect(Collectors.toMap(sourceTopicAndPartitionCount -> sourceToRemoteTopic.get(sourceTopicAndPartitionCount.getKey()), sourceTopicAndPartitionCount -> NewPartitions.increaseTo(sourceTopicAndPartitionCount.getValue().intValue()))); createNewPartitions(newTargetPartitions); } From 5b4925d78831546dce505aac282bae4e44f68384 Mon Sep 17 00:00:00 2001 From: Dhruvil Shah Date: Fri, 26 Feb 2021 00:26:03 -0800 Subject: [PATCH 6/7] minor change --- .../kafka/connect/mirror/MirrorSourceConnector.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 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 015d535fff679..2c92029fb352d 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 @@ -318,12 +318,12 @@ void computeAndCreateTopicPartitions() throws ExecutionException, InterruptedExc Set knownSourceTopics = sourceTopicToPartitionCounts.keySet(); Set knownTargetTopics = targetTopicToPartitionCounts.keySet(); - Map sourceToRemoteTopic = knownSourceTopics.stream() + Map sourceToRemoteTopics = knownSourceTopics.stream() .collect(Collectors.toMap(Function.identity(), sourceTopic -> formatRemoteTopic(sourceTopic))); // compute existing and new source topics Map> partitionedSourceTopics = knownSourceTopics.stream() - .collect(Collectors.partitioningBy(sourceTopic -> knownTargetTopics.contains(sourceToRemoteTopic.get(sourceTopic)), + .collect(Collectors.partitioningBy(sourceTopic -> knownTargetTopics.contains(sourceToRemoteTopics.get(sourceTopic)), Collectors.toSet())); Set existingSourceTopics = partitionedSourceTopics.get(true); Set newSourceTopics = partitionedSourceTopics.get(false); @@ -335,7 +335,7 @@ void computeAndCreateTopicPartitions() throws ExecutionException, InterruptedExc // compute topics with new partitions Map sourceTopicsWithNewPartitions = existingSourceTopics.stream() .filter(sourceTopic -> { - String targetTopic = sourceToRemoteTopic.get(sourceTopic); + String targetTopic = sourceToRemoteTopics.get(sourceTopic); return sourceTopicToPartitionCounts.get(sourceTopic) > targetTopicToPartitionCounts.get(targetTopic); }) .collect(Collectors.toMap(Function.identity(), sourceTopicToPartitionCounts::get)); @@ -343,7 +343,7 @@ void computeAndCreateTopicPartitions() throws ExecutionException, InterruptedExc // create new partitions if (!sourceTopicsWithNewPartitions.isEmpty()) { Map newTargetPartitions = sourceTopicsWithNewPartitions.entrySet().stream() - .collect(Collectors.toMap(sourceTopicAndPartitionCount -> sourceToRemoteTopic.get(sourceTopicAndPartitionCount.getKey()), + .collect(Collectors.toMap(sourceTopicAndPartitionCount -> sourceToRemoteTopics.get(sourceTopicAndPartitionCount.getKey()), sourceTopicAndPartitionCount -> NewPartitions.increaseTo(sourceTopicAndPartitionCount.getValue().intValue()))); createNewPartitions(newTargetPartitions); } From 79d2029b369c7f345d8bf6e8355af6a76140245e Mon Sep 17 00:00:00 2001 From: Dhruvil Shah Date: Fri, 26 Feb 2021 10:23:50 -0800 Subject: [PATCH 7/7] address review comment --- .../kafka/connect/mirror/MirrorSourceConnector.java | 12 +++++------- .../connect/mirror/MirrorSourceConnectorTest.java | 5 ++--- 2 files changed, 7 insertions(+), 10 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 2c92029fb352d..7b844c84d2f56 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 @@ -352,7 +352,7 @@ void computeAndCreateTopicPartitions() throws ExecutionException, InterruptedExc private void createNewTopics(Set newSourceTopics, Map sourceTopicToPartitionCounts) throws ExecutionException, InterruptedException { Map sourceTopicToConfig = describeTopicConfigs(newSourceTopics); - List newTopics = newSourceTopics.stream() + Map newTopics = newSourceTopics.stream() .map(sourceTopic -> { String remoteTopic = formatRemoteTopic(sourceTopic); int partitionCount = sourceTopicToPartitionCounts.get(sourceTopic).intValue(); @@ -360,19 +360,17 @@ private void createNewTopics(Set newSourceTopics, Map sour return new NewTopic(remoteTopic, partitionCount, (short) replicationFactor) .configs(configs); }) - .collect(Collectors.toList()); + .collect(Collectors.toMap(NewTopic::name, Function.identity())); createNewTopics(newTopics); } // visible for testing - void createNewTopics(List newTopics) { - Map newTopicMap = newTopics.stream() - .collect(Collectors.toMap(NewTopic::name, Function.identity())); - targetAdminClient.createTopics(newTopics, new CreateTopicsOptions()).values().forEach((k, v) -> v.whenComplete((x, e) -> { + void createNewTopics(Map newTopics) { + targetAdminClient.createTopics(newTopics.values(), new CreateTopicsOptions()).values().forEach((k, v) -> v.whenComplete((x, e) -> { if (e != null) { log.warn("Could not create topic {}.", k, e); } else { - log.info("Created remote topic {} with {} partitions.", k, newTopicMap.get(k).numPartitions()); + log.info("Created remote topic {} with {} partitions.", k, newTopics.get(k).numPartitions()); } })); } 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 66a7429e6d482..42d7951cd60fc 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 @@ -203,9 +203,8 @@ public void testRefreshTopicPartitions() throws Exception { Map configMap = MirrorSourceConnector.configToMap(topicConfig); assertEquals(2, configMap.size()); - List expectedNewTopics = Arrays.asList( - new NewTopic("source.topic", 1, (short) 0) - .configs(configMap)); + Map expectedNewTopics = new HashMap<>(); + expectedNewTopics.put("source.topic", new NewTopic("source.topic", 1, (short) 0).configs(configMap)); verify(connector, times(2)).computeAndCreateTopicPartitions(); verify(connector, times(2)).createNewTopics(eq(expectedNewTopics));