diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java index 14af7d9dce57d..46d32e5cc26b3 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java @@ -4945,7 +4945,7 @@ public StartMirrorTopicsResult startMirrorTopics(String mirrorName, Set validateRegexPatterns(options.excludePatterns()); // Fetch source metadata so the controller can create topics synchronously - Map topicMetadata; + Map topicMetadata; if (!topics.isEmpty()) { try { topicMetadata = fetchSourceTopicMetadata(mirrorName, topics); @@ -4966,9 +4966,9 @@ StartMirrorTopicsRequest.Builder createRequest(int timeoutMs) { StartMirrorTopicsRequestData data = new StartMirrorTopicsRequestData(); data.setMirrorName(mirrorName); topics.forEach(t -> { - StartMirrorTopicsRequestData.TopicData existing = topicMetadata.get(t); + StartMirrorTopicsRequestData.TopicMetadata existing = topicMetadata.get(t); data.topics().add(existing != null ? existing - : new StartMirrorTopicsRequestData.TopicData().setTopicName(t)); + : new StartMirrorTopicsRequestData.TopicMetadata().setTopicName(t)); }); data.setIncludePatterns(options.includePatterns()); data.setExcludePatterns(options.excludePatterns()); @@ -5002,7 +5002,7 @@ void handleFailure(Throwable throwable) { return new StartMirrorTopicsResult(future); } - private Map fetchSourceTopicMetadata( + private Map fetchSourceTopicMetadata( String mirrorName, Set topics) throws Exception { ConfigResource mirrorResource = new ConfigResource(ConfigResource.Type.CLUSTER_MIRROR, mirrorName); var configResult = describeConfigs(List.of(mirrorResource)).all().get(); @@ -5018,9 +5018,9 @@ private Map fetchSourceTopicMeta try (Admin sourceAdmin = Admin.create(sourceProps)) { var descriptions = sourceAdmin.describeTopics(topics).allTopicNames().get(); - Map metadata = new HashMap<>(); + Map metadata = new HashMap<>(); descriptions.forEach((name, desc) -> - metadata.put(name, new StartMirrorTopicsRequestData.TopicData() + metadata.put(name, new StartMirrorTopicsRequestData.TopicMetadata() .setTopicName(name) .setTopicId(desc.topicId()) .setNumPartitions(desc.partitions().size()))); @@ -5042,7 +5042,7 @@ public StopMirrorTopicsResult stopMirrorTopics(String mirrorName, Set to StopMirrorTopicsRequest.Builder createRequest(int timeoutMs) { StopMirrorTopicsRequestData data = new StopMirrorTopicsRequestData(); data.setMirrorName(mirrorName); - topics.forEach(t -> data.topics().add(new StopMirrorTopicsRequestData.TopicData().setTopicName(t))); + topics.forEach(t -> data.topics().add(new StopMirrorTopicsRequestData.TopicMetadata().setTopicName(t))); data.setPatterns(options.patterns()); return new StopMirrorTopicsRequest.Builder(data); } diff --git a/clients/src/main/java/org/apache/kafka/common/requests/PauseMirrorTopicsRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/PauseMirrorTopicsRequest.java index f344a332f924d..b1878c0e0898f 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/PauseMirrorTopicsRequest.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/PauseMirrorTopicsRequest.java @@ -40,7 +40,7 @@ public Builder(String mirrorName, Set topics) { ApiKeys.PAUSE_MIRROR_TOPICS.latestVersion()); PauseMirrorTopicsRequestData data = new PauseMirrorTopicsRequestData(); data.setMirrorName(mirrorName); - topics.forEach(topic -> data.topics().add(new PauseMirrorTopicsRequestData.TopicData().setTopicName(topic))); + topics.forEach(topic -> data.topics().add(new PauseMirrorTopicsRequestData.TopicMetadata().setTopicName(topic))); this.data = data; } diff --git a/clients/src/main/java/org/apache/kafka/common/requests/ResumeMirrorTopicsRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/ResumeMirrorTopicsRequest.java index 09a873a5d20a6..6c6f5cae75d24 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/ResumeMirrorTopicsRequest.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/ResumeMirrorTopicsRequest.java @@ -40,7 +40,7 @@ public Builder(String mirrorName, Set topics) { ApiKeys.RESUME_MIRROR_TOPICS.latestVersion()); ResumeMirrorTopicsRequestData data = new ResumeMirrorTopicsRequestData(); data.setMirrorName(mirrorName); - topics.forEach(topic -> data.topics().add(new ResumeMirrorTopicsRequestData.TopicData().setTopicName(topic))); + topics.forEach(topic -> data.topics().add(new ResumeMirrorTopicsRequestData.TopicMetadata().setTopicName(topic))); this.data = data; } diff --git a/clients/src/main/java/org/apache/kafka/common/requests/StartMirrorTopicsRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/StartMirrorTopicsRequest.java index 1fa45e4b22529..3890dc89346e8 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/StartMirrorTopicsRequest.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/StartMirrorTopicsRequest.java @@ -40,7 +40,7 @@ public Builder(String mirrorName, Set topics) { ApiKeys.START_MIRROR_TOPICS.latestVersion()); StartMirrorTopicsRequestData data = new StartMirrorTopicsRequestData(); data.setMirrorName(mirrorName); - topics.forEach(topic -> data.topics().add(new StartMirrorTopicsRequestData.TopicData().setTopicName(topic))); + topics.forEach(topic -> data.topics().add(new StartMirrorTopicsRequestData.TopicMetadata().setTopicName(topic))); this.data = data; } diff --git a/clients/src/main/java/org/apache/kafka/common/requests/StopMirrorTopicsRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/StopMirrorTopicsRequest.java index 59f0cced52a7c..048f612bd36c8 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/StopMirrorTopicsRequest.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/StopMirrorTopicsRequest.java @@ -40,7 +40,7 @@ public Builder(String mirrorName, Set topics) { ApiKeys.STOP_MIRROR_TOPICS.latestVersion()); StopMirrorTopicsRequestData data = new StopMirrorTopicsRequestData(); data.setMirrorName(mirrorName); - topics.forEach(topic -> data.topics().add(new StopMirrorTopicsRequestData.TopicData().setTopicName(topic))); + topics.forEach(topic -> data.topics().add(new StopMirrorTopicsRequestData.TopicMetadata().setTopicName(topic))); this.data = data; } diff --git a/clients/src/main/resources/common/message/PauseMirrorTopicsRequest.json b/clients/src/main/resources/common/message/PauseMirrorTopicsRequest.json index b277321487758..f1a3423b61752 100644 --- a/clients/src/main/resources/common/message/PauseMirrorTopicsRequest.json +++ b/clients/src/main/resources/common/message/PauseMirrorTopicsRequest.json @@ -23,8 +23,8 @@ "flexibleVersions": "0+", "fields": [ { "name": "MirrorName", "type": "string", "versions": "0+", "entityType": "mirrorName", - "about": "The mirror name to pause the topics for." }, - { "name": "Topics", "type": "[]TopicData", "versions": "0+", "about": "The data for the topics.", + "about": "The cluster mirror name." }, + { "name": "Topics", "type": "[]TopicMetadata", "versions": "0+", "about": "The data for the topics.", "fields": [ { "name": "TopicId", "type": "uuid", "versions": "0+", "about": "The unique topic ID."}, { "name": "TopicName", "type": "string", "versions": "0+", "mapKey": true, "entityType": "topicName", diff --git a/clients/src/main/resources/common/message/ReadMirrorStatesRequest.json b/clients/src/main/resources/common/message/ReadMirrorStatesRequest.json index 5735a89208484..4b5880ecc51e6 100644 --- a/clients/src/main/resources/common/message/ReadMirrorStatesRequest.json +++ b/clients/src/main/resources/common/message/ReadMirrorStatesRequest.json @@ -24,7 +24,7 @@ "fields": [ { "name": "MirrorName", "type": "string", "versions": "0+", "entityType": "mirrorName", "about": "The cluster mirror name." }, - { "name": "Topics", "type": "[]TopicData", "versions": "0+", + { "name": "Topics", "type": "[]TopicMetadata", "versions": "0+", "about": "The data for the topics.", "fields": [ { "name": "Name", "type": "string", "versions": "0+", "entityType": "topicName", "about": "The topic name." }, diff --git a/clients/src/main/resources/common/message/ResumeMirrorTopicsRequest.json b/clients/src/main/resources/common/message/ResumeMirrorTopicsRequest.json index d6cfeba6ab97a..2dbaef74b315a 100644 --- a/clients/src/main/resources/common/message/ResumeMirrorTopicsRequest.json +++ b/clients/src/main/resources/common/message/ResumeMirrorTopicsRequest.json @@ -24,7 +24,7 @@ "fields": [ { "name": "MirrorName", "type": "string", "versions": "0+", "entityType": "mirrorName", "about": "The cluster mirror name." }, - { "name": "Topics", "type": "[]TopicData", "versions": "0+", "about": "The data for the topics.", + { "name": "Topics", "type": "[]TopicMetadata", "versions": "0+", "about": "The data for the topics.", "fields": [ { "name": "TopicId", "type": "uuid", "versions": "0+", "about": "The unique topic ID."}, { "name": "TopicName", "type": "string", "versions": "0+", "mapKey": true, "entityType": "topicName", diff --git a/clients/src/main/resources/common/message/StartMirrorTopicsRequest.json b/clients/src/main/resources/common/message/StartMirrorTopicsRequest.json index 0345b4bd43ff6..f64c5fb28cd89 100644 --- a/clients/src/main/resources/common/message/StartMirrorTopicsRequest.json +++ b/clients/src/main/resources/common/message/StartMirrorTopicsRequest.json @@ -24,7 +24,7 @@ "fields": [ { "name": "MirrorName", "type": "string", "versions": "0+", "entityType": "mirrorName", "about": "The cluster mirror name." }, - { "name": "Topics", "type": "[]TopicData", "versions": "0+", "about": "The data for the topics.", + { "name": "Topics", "type": "[]TopicMetadata", "versions": "0+", "about": "The data for the topics.", "fields": [ { "name": "TopicId", "type": "uuid", "versions": "0+", "about": "The unique topic ID."}, { "name": "TopicName", "type": "string", "versions": "0+", "mapKey": true, "entityType": "topicName", diff --git a/clients/src/main/resources/common/message/StopMirrorTopicsRequest.json b/clients/src/main/resources/common/message/StopMirrorTopicsRequest.json index 9cdc3ddaad062..d89cb5fad63af 100644 --- a/clients/src/main/resources/common/message/StopMirrorTopicsRequest.json +++ b/clients/src/main/resources/common/message/StopMirrorTopicsRequest.json @@ -24,7 +24,7 @@ "fields": [ { "name": "MirrorName", "type": "string", "versions": "0+", "entityType": "mirrorName", "about": "The cluster mirror name." }, - { "name": "Topics", "type": "[]TopicData", "versions": "0+", "about": "The data for the topics.", + { "name": "Topics", "type": "[]TopicMetadata", "versions": "0+", "about": "The data for the topics.", "fields": [ { "name": "TopicId", "type": "uuid", "versions": "0+", "about": "The unique topic ID."}, { "name": "TopicName", "type": "string", "versions": "0+", "mapKey": true, "entityType": "topicName", diff --git a/clients/src/main/resources/common/message/WriteMirrorStatesRequest.json b/clients/src/main/resources/common/message/WriteMirrorStatesRequest.json index eb41f01472132..208b52bbfa4a8 100644 --- a/clients/src/main/resources/common/message/WriteMirrorStatesRequest.json +++ b/clients/src/main/resources/common/message/WriteMirrorStatesRequest.json @@ -23,8 +23,8 @@ "flexibleVersions": "0+", "fields": [ { "name": "MirrorName", "type": "string", "versions": "0+", "entityType": "mirrorName", - "about": "The mirror name." }, - { "name": "Topics", "type": "[]TopicData", "versions": "0+", + "about": "The cluster mirror name." }, + { "name": "Topics", "type": "[]TopicMetadata", "versions": "0+", "about": "The data for the topics.", "fields": [ { "name": "Name", "type": "string", "versions": "0+", "entityType": "topicName", "about": "The topic name." }, diff --git a/clients/src/test/java/org/apache/kafka/common/requests/RequestResponseTest.java b/clients/src/test/java/org/apache/kafka/common/requests/RequestResponseTest.java index 844e7a3b7c80b..4d5004d742492 100644 --- a/clients/src/test/java/org/apache/kafka/common/requests/RequestResponseTest.java +++ b/clients/src/test/java/org/apache/kafka/common/requests/RequestResponseTest.java @@ -1223,8 +1223,8 @@ private AbstractResponse getResponse(ApiKeys apikey, short version) { } private StartMirrorTopicsRequest createStartMirrorTopicsRequest(short version) { - StartMirrorTopicsRequestData.TopicDataCollection topics = new StartMirrorTopicsRequestData.TopicDataCollection(); - topics.add(new StartMirrorTopicsRequestData.TopicData() + StartMirrorTopicsRequestData.TopicMetadataCollection topics = new StartMirrorTopicsRequestData.TopicMetadataCollection(); + topics.add(new StartMirrorTopicsRequestData.TopicMetadata() .setTopicId(Uuid.randomUuid()) .setTopicName("topic") ); @@ -1242,8 +1242,8 @@ private StartMirrorTopicsResponse createStartMirrorTopicsResponse() { } private StopMirrorTopicsRequest createStopMirrorTopicsRequest(short version) { - StopMirrorTopicsRequestData.TopicDataCollection topics = new StopMirrorTopicsRequestData.TopicDataCollection(); - topics.add(new StopMirrorTopicsRequestData.TopicData() + StopMirrorTopicsRequestData.TopicMetadataCollection topics = new StopMirrorTopicsRequestData.TopicMetadataCollection(); + topics.add(new StopMirrorTopicsRequestData.TopicMetadata() .setTopicId(Uuid.randomUuid()) .setTopicName("topic") ); @@ -1303,7 +1303,7 @@ public DescribeClusterMirrorsResponse createDescribeClusterMirrorsResponse() { public ReadMirrorStatesRequest createReadMirrorStatesRequest(short version) { ReadMirrorStatesRequestData data = new ReadMirrorStatesRequestData() .setMirrorName("mirror") - .setTopics(List.of(new ReadMirrorStatesRequestData.TopicData() + .setTopics(List.of(new ReadMirrorStatesRequestData.TopicMetadata() .setName("topic") .setPartitions(List.of(new ReadMirrorStatesRequestData.PartitionData().setPartitionIndex(0))) )); @@ -1325,7 +1325,7 @@ public ReadMirrorStatesResponse createReadMirrorStatesResponse() { public WriteMirrorStatesRequest createWriteMirrorStatesRequest(short version) { WriteMirrorStatesRequestData data = new WriteMirrorStatesRequestData() .setMirrorName("mirror") - .setTopics(List.of(new WriteMirrorStatesRequestData.TopicData() + .setTopics(List.of(new WriteMirrorStatesRequestData.TopicMetadata() .setName("topic") .setPartitions(List.of(new WriteMirrorStatesRequestData.PartitionData() .setPartitionIndex(0) @@ -1346,8 +1346,8 @@ public WriteMirrorStatesResponse createWriteMirrorStatesResponse() { } public PauseMirrorTopicsRequest createPauseMirrorTopicsRequest(short version) { - PauseMirrorTopicsRequestData.TopicDataCollection topics = new PauseMirrorTopicsRequestData.TopicDataCollection(); - topics.add(new PauseMirrorTopicsRequestData.TopicData() + PauseMirrorTopicsRequestData.TopicMetadataCollection topics = new PauseMirrorTopicsRequestData.TopicMetadataCollection(); + topics.add(new PauseMirrorTopicsRequestData.TopicMetadata() .setTopicId(Uuid.randomUuid()) .setTopicName("topic") ); @@ -1366,8 +1366,8 @@ public PauseMirrorTopicsResponse createPauseMirrorTopicsResponse() { } public ResumeMirrorTopicsRequest createResumeMirrorTopicsRequest(short version) { - ResumeMirrorTopicsRequestData.TopicDataCollection topics = new ResumeMirrorTopicsRequestData.TopicDataCollection(); - topics.add(new ResumeMirrorTopicsRequestData.TopicData() + ResumeMirrorTopicsRequestData.TopicMetadataCollection topics = new ResumeMirrorTopicsRequestData.TopicMetadataCollection(); + topics.add(new ResumeMirrorTopicsRequestData.TopicMetadata() .setTopicId(Uuid.randomUuid()) .setTopicName("topic") ); diff --git a/core/src/main/java/kafka/server/mirror/MirrorMetadataManager.java b/core/src/main/java/kafka/server/mirror/MirrorMetadataManager.java index 188a8306536bc..825bb2b652dbe 100644 --- a/core/src/main/java/kafka/server/mirror/MirrorMetadataManager.java +++ b/core/src/main/java/kafka/server/mirror/MirrorMetadataManager.java @@ -1250,7 +1250,7 @@ private void discoverTopicsByPattern(String mirrorName, ClusterMirrorConfig mirr Set configuredTopics = getConfiguredTopics(mirrorName, true); final Pattern topicsExcludePattern = mirrorConfig.topicsExcludePattern(); - List newTopics; + List newTopics; try { Set allSourceTopics = srcAdmin.listTopics() .names().get(brokerConfig.requestTimeoutMs(), TimeUnit.MILLISECONDS); @@ -1271,7 +1271,7 @@ private void discoverTopicsByPattern(String mirrorName, ClusterMirrorConfig mirr cacheSourceLeaders(mirrorName, descriptions.values()); newTopics = descriptions.values().stream() - .map(td -> new StartMirrorTopicsRequestData.TopicData() + .map(td -> new StartMirrorTopicsRequestData.TopicMetadata() .setTopicName(td.name()) .setTopicId(td.topicId()) .setNumPartitions(td.partitions().size())) @@ -1286,7 +1286,7 @@ private void discoverTopicsByPattern(String mirrorName, ClusterMirrorConfig mirr } log.info("Discovered {} new topic(s) matching mirror.topics.include pattern for mirror {}: {}", - newTopics.size(), mirrorName, newTopics.stream().map(StartMirrorTopicsRequestData.TopicData::topicName).toList()); + newTopics.size(), mirrorName, newTopics.stream().map(StartMirrorTopicsRequestData.TopicMetadata::topicName).toList()); StartMirrorTopicsRequestData data = new StartMirrorTopicsRequestData(); data.setMirrorName(mirrorName); @@ -1459,10 +1459,10 @@ void writeStatesToRemoteCoordinator(String mirrorName, // Send one batched request per coordinator node nodeToTopicPartitions.forEach((node, topicPartitionsMap) -> { WriteMirrorStatesRequestData data = new WriteMirrorStatesRequestData().setMirrorName(mirrorName); - List topicDataList = new ArrayList<>(); + List topicDataList = new ArrayList<>(); topicPartitionsMap.forEach((topic, partitionDataList) -> - topicDataList.add(new WriteMirrorStatesRequestData.TopicData() + topicDataList.add(new WriteMirrorStatesRequestData.TopicMetadata() .setName(topic) .setPartitions(partitionDataList))); @@ -1517,10 +1517,10 @@ void readStatesFromRemoteCoordinator(String mirrorName, // Send one batched request per coordinator node nodeToTopicPartitions.forEach((node, topicPartitionsMap) -> { ReadMirrorStatesRequestData data = new ReadMirrorStatesRequestData().setMirrorName(mirrorName); - List topicDataList = new ArrayList<>(); + List topicDataList = new ArrayList<>(); topicPartitionsMap.forEach((topic, partitionDataList) -> - topicDataList.add(new ReadMirrorStatesRequestData.TopicData() + topicDataList.add(new ReadMirrorStatesRequestData.TopicMetadata() .setName(topic) .setPartitions(partitionDataList)));