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..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 @@ -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,40 +307,82 @@ private void createOffsetSyncsTopic() { MirrorUtils.createSinglePartitionCompactedTopic(config.offsetSyncsTopic(), config.offsetSyncsTopicReplicationFactor(), config.sourceAdminConfig()); } - // visible for testing - void computeAndCreateTopicPartitions() - throws InterruptedException, ExecutionException { - Map partitionCounts = 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)) - .collect(Collectors.toList()); - Map newPartitions = partitionCounts.entrySet().stream() - .filter(x -> knownTargetTopics.contains(x.getKey())) - .collect(Collectors.toMap(Entry::getKey, x -> NewPartitions.increaseTo(x.getValue().intValue()))); - createTopicPartitions(partitionCounts, newTopics, newPartitions); + 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(); + 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(sourceToRemoteTopics.get(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 = sourceToRemoteTopics.get(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 -> sourceToRemoteTopics.get(sourceTopicAndPartitionCount.getKey()), + sourceTopicAndPartitionCount -> NewPartitions.increaseTo(sourceTopicAndPartitionCount.getValue().intValue()))); + createNewPartitions(newTargetPartitions); + } + } + + private void createNewTopics(Set newSourceTopics, Map sourceTopicToPartitionCounts) + throws ExecutionException, InterruptedException { + Map sourceTopicToConfig = describeTopicConfigs(newSourceTopics); + Map 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.toMap(NewTopic::name, Function.identity())); + createNewTopics(newTopics); } // visible for testing - void createTopicPartitions(Map partitionCounts, List newTopics, - Map newPartitions) { - 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, partitionCounts.get(k)); + log.info("Created remote topic {} with {} partitions.", k, newTopics.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()); } })); } @@ -359,6 +402,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::name, ConfigEntry::value)); + } + @SuppressWarnings("deprecation") // use deprecated alterConfigs API for broker compatibility back to 0.11.0 private void updateTopicConfigs(Map topicConfigs) @@ -390,7 +438,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..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 @@ -183,10 +183,16 @@ public void testRefreshTopicPartitions() throws Exception { connector.initialize(mock(ConnectorContext.class)); connector = spy(connector); + Config topicConfig = new Config(Arrays.asList( + 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(); - 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 @@ -194,13 +200,15 @@ 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()); + + 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)).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(); @@ -217,11 +225,19 @@ public void testRefreshTopicPartitionsTopicOnTargetFirst() throws Exception { connector.initialize(mock(ConnectorContext.class)); connector = spy(connector); + Config topicConfig = new Config(Arrays.asList( + new ConfigEntry("cleanup.policy", "compact"), + new ConfigEntry("segment.bytes", "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(); - doNothing().when(connector).createTopicPartitions(any(), any(), any()); + doReturn(configs).when(connector).describeTopicConfigs(Collections.singleton("source.topic")); + doReturn(Collections.emptyMap()).when(connector).describeTopicConfigs(Collections.emptySet()); + doNothing().when(connector).createNewTopics(any()); + doNothing().when(connector).createNewPartitions(any()); // partitions appearing on the target cluster should not cause reconfiguration connector.refreshTopicPartitions(); @@ -234,6 +250,5 @@ public void testRefreshTopicPartitionsTopicOnTargetFirst() throws Exception { // when partitions are added to the source cluster, reconfiguration is triggered connector.refreshTopicPartitions(); verify(connector, times(1)).computeAndCreateTopicPartitions(); - } }