From 601ee9f12d48e774871686f17d2d009acf851027 Mon Sep 17 00:00:00 2001 From: Federico Valeri Date: Wed, 17 Jun 2026 18:19:13 +0200 Subject: [PATCH 1/3] Address AS32 The use of StartMirrorTopicsRequestData.TopicData in an external interface looks odd. Signed-off-by: Federico Valeri --- .../scala/kafka/server/ControllerApis.scala | 8 +++--- .../ConfigurationControlManager.java | 13 +++++---- .../apache/kafka/controller/Controller.java | 27 ++++++++++--------- .../kafka/controller/QuorumController.java | 21 +++++++-------- .../kafka/common/test/MockController.java | 3 +-- 5 files changed, 36 insertions(+), 36 deletions(-) diff --git a/core/src/main/scala/kafka/server/ControllerApis.scala b/core/src/main/scala/kafka/server/ControllerApis.scala index 8c274807fd072..fb284d1307c22 100644 --- a/core/src/main/scala/kafka/server/ControllerApis.scala +++ b/core/src/main/scala/kafka/server/ControllerApis.scala @@ -285,16 +285,18 @@ class ControllerApis( throw new ClusterMirrorAuthorizationException(s"Request $request needs ALTER permission on ClusterMirror:$mirrorName.") if (!ClusterMirrorUtils.isClusterMirroringEnabled(apiVersionManager.features.finalizedFeatures)) throw new UnsupportedVersionException("Cluster mirroring requires mirror.version >= 1.") - val topics = startMirrorTopicsRequest.data().topics() - val unauthorizedTopics = topics.asScala.map(_.topicName()).filterNot(topic => + val wireTopics = startMirrorTopicsRequest.data().topics() + val unauthorizedTopics = wireTopics.asScala.map(_.topicName()).filterNot(topic => authHelper.authorize(request.context, ALTER_CONFIGS, TOPIC, topic, logIfDenied = false)).toSet if (unauthorizedTopics.nonEmpty) throw new TopicAuthorizationException(unauthorizedTopics.asJava) + val topics = wireTopics.asScala.map(t => + new Controller.MirrorTopicMetadata(t.topicName(), t.topicId(), t.numPartitions())).toList.asJava val includePatterns = startMirrorTopicsRequest.data().includePatterns() val excludePatterns = startMirrorTopicsRequest.data().excludePatterns() val context = new ControllerRequestContext(request.context.header.data, request.context.principal, OptionalLong.empty()) - controller.startMirrorTopics(context, mirrorName, new java.util.ArrayList(topics), + controller.startMirrorTopics(context, mirrorName, topics, includePatterns, excludePatterns) .handle[Unit] { (response, exception) => if (exception != null) { diff --git a/metadata/src/main/java/org/apache/kafka/controller/ConfigurationControlManager.java b/metadata/src/main/java/org/apache/kafka/controller/ConfigurationControlManager.java index 9139dc61dcc22..9fd8800327ded 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ConfigurationControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ConfigurationControlManager.java @@ -31,7 +31,6 @@ import org.apache.kafka.common.message.DeleteClusterMirrorResponseData; import org.apache.kafka.common.message.PauseMirrorTopicsResponseData; import org.apache.kafka.common.message.ResumeMirrorTopicsResponseData; -import org.apache.kafka.common.message.StartMirrorTopicsRequestData; import org.apache.kafka.common.message.StartMirrorTopicsResponseData; import org.apache.kafka.common.message.StopMirrorTopicsResponseData; import org.apache.kafka.common.metadata.ClearElrRecord; @@ -399,7 +398,7 @@ ControllerResult resumeMirrorTopics(String mirro ControllerResult startMirrorTopics( String mirrorName, - List topics, + List topics, List includePatterns, List excludePatterns, ReplicationControlManager replicationControl) { @@ -407,7 +406,7 @@ ControllerResult startMirrorTopics( StartMirrorTopicsResponseData data = new StartMirrorTopicsResponseData(); data.setMirrorName(mirrorName); - Set topicNames = topics.stream().map(StartMirrorTopicsRequestData.TopicData::topicName).collect(Collectors.toSet()); + Set topicNames = topics.stream().map(Controller.MirrorTopicMetadata::name).collect(Collectors.toSet()); ApiError patternError = updatePatternsAndStopExcluded(mirrorName, records, topicNames, replicationControl, (includeSet, excludeSet) -> { // auto-include explicitly started topics so validation works fine for (String topicName : topicNames) { @@ -429,15 +428,15 @@ ControllerResult startMirrorTopics( } List topicResList = new ArrayList<>(); - for (StartMirrorTopicsRequestData.TopicData topic : topics) { + for (Controller.MirrorTopicMetadata topic : topics) { StartMirrorTopicsResponseData.TopicResult topicRes = new StartMirrorTopicsResponseData.TopicResult(); - String topicName = topic.topicName(); + String topicName = topic.name(); // When topic ID and partition count are missing, the caller only provided the topic // name without metatata (see StartMirrorTopicsOptions). The topic was already added // to mirror.topics.include above, so discoverTopicsByPattern will fetch the metadata // from the source and create it at the next metadata refresh. - if (topic.topicId().equals(Uuid.ZERO_UUID) && topic.numPartitions() <= 0) { + if (topic.id().equals(Uuid.ZERO_UUID) && topic.numPartitions() <= 0) { log.warn("Topic {} for mirror {} has no topic ID or partition info and will be" + " created at the next metadata refresh", topicName, mirrorName); topicRes.setName(topicName); @@ -446,7 +445,7 @@ ControllerResult startMirrorTopics( } // For pre-2.8 sources (ZERO_UUID) with partition info, assign a random UUID - Uuid topicId = topic.topicId().equals(Uuid.ZERO_UUID) ? Uuid.randomUuid() : topic.topicId(); + Uuid topicId = topic.id().equals(Uuid.ZERO_UUID) ? Uuid.randomUuid() : topic.id(); if (topic.numPartitions() > 0) { ApiError createError = replicationControl.createMirrorTopic( diff --git a/metadata/src/main/java/org/apache/kafka/controller/Controller.java b/metadata/src/main/java/org/apache/kafka/controller/Controller.java index 9723aea8781e3..bbe3e595d39a3 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/Controller.java +++ b/metadata/src/main/java/org/apache/kafka/controller/Controller.java @@ -52,7 +52,6 @@ import org.apache.kafka.common.message.RenewDelegationTokenRequestData; import org.apache.kafka.common.message.RenewDelegationTokenResponseData; import org.apache.kafka.common.message.ResumeMirrorTopicsResponseData; -import org.apache.kafka.common.message.StartMirrorTopicsRequestData; import org.apache.kafka.common.message.StartMirrorTopicsResponseData; import org.apache.kafka.common.message.StopMirrorTopicsResponseData; import org.apache.kafka.common.message.UpdateFeaturesRequestData; @@ -152,12 +151,6 @@ CompletableFuture createTopics( Set describable ); - CompletableFuture createClusterMirror( - ControllerRequestContext context, - String mirrorName, - Map> configChanges - ); - /** * Unregister a broker. * @@ -172,10 +165,16 @@ CompletableFuture unregisterBroker( int brokerId ); + CompletableFuture createClusterMirror( + ControllerRequestContext context, + String mirrorName, + Map> configChanges + ); + CompletableFuture startMirrorTopics( ControllerRequestContext context, String mirrorName, - List topics, + List topics, List includePatterns, List excludePatterns ); @@ -199,17 +198,17 @@ CompletableFuture resumeMirrorTopics( Set topics ); - CompletableFuture bumpLeaderEpoch( - ControllerRequestContext context, - Map> partitionLeaderEpochs - ); - CompletableFuture deleteClusterMirror( ControllerRequestContext context, String mirrorName, long brokerMetadataOffset ); + CompletableFuture bumpLeaderEpoch( + ControllerRequestContext context, + Map> partitionLeaderEpochs + ); + /** * Find the ids for topic names. * @@ -493,4 +492,6 @@ default boolean isActive() { * Blocks until we have shut down and freed all resources. */ void close() throws InterruptedException; + + record MirrorTopicMetadata(String name, Uuid id, int numPartitions) { } } diff --git a/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java b/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java index bc817e4703189..e42be80908a0d 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java +++ b/metadata/src/main/java/org/apache/kafka/controller/QuorumController.java @@ -61,7 +61,6 @@ import org.apache.kafka.common.message.RenewDelegationTokenRequestData; import org.apache.kafka.common.message.RenewDelegationTokenResponseData; import org.apache.kafka.common.message.ResumeMirrorTopicsResponseData; -import org.apache.kafka.common.message.StartMirrorTopicsRequestData; import org.apache.kafka.common.message.StartMirrorTopicsResponseData; import org.apache.kafka.common.message.StopMirrorTopicsResponseData; import org.apache.kafka.common.message.UpdateFeaturesRequestData; @@ -1789,20 +1788,11 @@ public CompletableFuture createClusterMirror( }); } - @Override - public CompletableFuture bumpLeaderEpoch( - ControllerRequestContext context, - Map> partitionLeaderEpochs - ) { - return appendWriteEvent("bumpLeaderEpochs", context.deadlineNs(), - () -> replicationControl.bumpLeaderEpochs(partitionLeaderEpochs)); - } - @Override public CompletableFuture startMirrorTopics( ControllerRequestContext context, String mirrorName, - List topics, + List topics, List includePatterns, List excludePatterns ) { @@ -1852,6 +1842,15 @@ public CompletableFuture deleteClusterMirror( () -> configurationControl.deleteClusterMirror(mirrorName, brokerMetadataOffset, replicationControl)); } + @Override + public CompletableFuture bumpLeaderEpoch( + ControllerRequestContext context, + Map> partitionLeaderEpochs + ) { + return appendWriteEvent("bumpLeaderEpochs", context.deadlineNs(), + () -> replicationControl.bumpLeaderEpochs(partitionLeaderEpochs)); + } + @Override public CompletableFuture unregisterBroker( ControllerRequestContext context, diff --git a/test-common/test-common-runtime/src/main/java/org/apache/kafka/common/test/MockController.java b/test-common/test-common-runtime/src/main/java/org/apache/kafka/common/test/MockController.java index feff69c755962..79931cec7d425 100644 --- a/test-common/test-common-runtime/src/main/java/org/apache/kafka/common/test/MockController.java +++ b/test-common/test-common-runtime/src/main/java/org/apache/kafka/common/test/MockController.java @@ -58,7 +58,6 @@ import org.apache.kafka.common.message.RenewDelegationTokenRequestData; import org.apache.kafka.common.message.RenewDelegationTokenResponseData; import org.apache.kafka.common.message.ResumeMirrorTopicsResponseData; -import org.apache.kafka.common.message.StartMirrorTopicsRequestData; import org.apache.kafka.common.message.StartMirrorTopicsResponseData; import org.apache.kafka.common.message.StopMirrorTopicsResponseData; import org.apache.kafka.common.message.UpdateFeaturesRequestData; @@ -134,7 +133,7 @@ public CompletableFuture stopMirrorTopics( public CompletableFuture startMirrorTopics( ControllerRequestContext context, String mirrorName, - List topics, + List topics, List includePatterns, List excludePatterns ) { From cc429384144db44472ea91a941c325a6cdec94ee Mon Sep 17 00:00:00 2001 From: Federico Valeri Date: Wed, 17 Jun 2026 18:28:49 +0200 Subject: [PATCH 2/3] Also address AS33 Signed-off-by: Federico Valeri --- .../common/message/PauseMirrorTopicsResponse.json | 4 +++- .../common/message/ResumeMirrorTopicsResponse.json | 4 +++- .../common/message/StartMirrorTopicsResponse.json | 4 +++- .../common/message/StopMirrorTopicsResponse.json | 4 +++- .../controller/ConfigurationControlManager.java | 13 +++++++++---- 5 files changed, 21 insertions(+), 8 deletions(-) diff --git a/clients/src/main/resources/common/message/PauseMirrorTopicsResponse.json b/clients/src/main/resources/common/message/PauseMirrorTopicsResponse.json index 9cb0a1579f2b4..8ebfc17101e55 100644 --- a/clients/src/main/resources/common/message/PauseMirrorTopicsResponse.json +++ b/clients/src/main/resources/common/message/PauseMirrorTopicsResponse.json @@ -34,7 +34,9 @@ { "name": "Name", "type": "string", "versions": "0+", "entityType": "topicName", "about": "The topic name." }, { "name": "ErrorCode", "type": "int16", "versions": "0+", - "about": "The error code, or 0 if there was no error." } + "about": "The error code, or 0 if there was no error." }, + { "name": "ErrorMessage", "type": "string", "versions": "0+", "nullableVersions": "0+", "default": "null", + "about": "The error message, or null if there was no error." } ]} ] } diff --git a/clients/src/main/resources/common/message/ResumeMirrorTopicsResponse.json b/clients/src/main/resources/common/message/ResumeMirrorTopicsResponse.json index 3b5d0a9fbad36..6d1e5c84ec562 100644 --- a/clients/src/main/resources/common/message/ResumeMirrorTopicsResponse.json +++ b/clients/src/main/resources/common/message/ResumeMirrorTopicsResponse.json @@ -34,7 +34,9 @@ { "name": "Name", "type": "string", "versions": "0+", "entityType": "topicName", "about": "The topic name." }, { "name": "ErrorCode", "type": "int16", "versions": "0+", - "about": "The error code, or 0 if there was no error." } + "about": "The error code, or 0 if there was no error." }, + { "name": "ErrorMessage", "type": "string", "versions": "0+", "nullableVersions": "0+", "default": "null", + "about": "The error message, or null if there was no error." } ]} ] } diff --git a/clients/src/main/resources/common/message/StartMirrorTopicsResponse.json b/clients/src/main/resources/common/message/StartMirrorTopicsResponse.json index cc953a4b31ff0..13b283d6e6e08 100644 --- a/clients/src/main/resources/common/message/StartMirrorTopicsResponse.json +++ b/clients/src/main/resources/common/message/StartMirrorTopicsResponse.json @@ -34,7 +34,9 @@ { "name": "Name", "type": "string", "versions": "0+", "entityType": "topicName", "about": "The topic name." }, { "name": "ErrorCode", "type": "int16", "versions": "0+", - "about": "The error code, or 0 if there was no error." } + "about": "The error code, or 0 if there was no error." }, + { "name": "ErrorMessage", "type": "string", "versions": "0+", "nullableVersions": "0+", "default": "null", + "about": "The error message, or null if there was no error." } ]} ] } diff --git a/clients/src/main/resources/common/message/StopMirrorTopicsResponse.json b/clients/src/main/resources/common/message/StopMirrorTopicsResponse.json index be3b1dfd02cf9..f029024ba5f92 100644 --- a/clients/src/main/resources/common/message/StopMirrorTopicsResponse.json +++ b/clients/src/main/resources/common/message/StopMirrorTopicsResponse.json @@ -34,7 +34,9 @@ { "name": "Name", "type": "string", "versions": "0+", "entityType": "topicName", "about": "The topic name." }, { "name": "ErrorCode", "type": "int16", "versions": "0+", - "about": "The error code, or 0 if there was no error." } + "about": "The error code, or 0 if there was no error." }, + { "name": "ErrorMessage", "type": "string", "versions": "0+", "nullableVersions": "0+", "default": "null", + "about": "The error message, or null if there was no error." } ]} ] } diff --git a/metadata/src/main/java/org/apache/kafka/controller/ConfigurationControlManager.java b/metadata/src/main/java/org/apache/kafka/controller/ConfigurationControlManager.java index 9fd8800327ded..d8d3968572ae6 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ConfigurationControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ConfigurationControlManager.java @@ -328,9 +328,10 @@ ControllerResult pauseMirrorTopics(String mirrorN continue; } - // Don't allow to pause a topic when stopping/stopped if (currMirrorStateChange == MirrorPartitionState.STOPPED.value()) { - topicRes.setErrorCode(Errors.MIRROR_TOPIC_BEING_STOPPED.code()).setName(topic); + topicRes.setErrorCode(Errors.MIRROR_TOPIC_BEING_STOPPED.code()).setName(topic) + .setErrorMessage("Topic '" + topic + "' is in " + + MirrorPartitionState.fromValue((byte) currMirrorStateChange) + " state"); topicResList.add(topicRes); continue; } @@ -379,7 +380,9 @@ ControllerResult resumeMirrorTopics(String mirro } if (currMirrorStateChange != MirrorPartitionState.PAUSED.value()) { - topicRes.setErrorCode(Errors.MIRROR_TOPIC_NOT_PAUSED.code()).setName(topic); + topicRes.setErrorCode(Errors.MIRROR_TOPIC_NOT_PAUSED.code()).setName(topic) + .setErrorMessage("Topic '" + topic + "' is in " + + MirrorPartitionState.fromValue((byte) currMirrorStateChange) + " state"); topicResList.add(topicRes); continue; } @@ -463,7 +466,9 @@ ControllerResult startMirrorTopics( String currMirrorNameValue = topicInfo.mirrorName(); int currMirrorStateChange = topicInfo.mirrorState(); if (currMirrorNameValue != null && !currMirrorNameValue.isBlank() && currMirrorStateChange != MirrorPartitionState.STOPPED.value()) { - topicRes.setErrorCode(Errors.TOPIC_ALREADY_IN_CLUSTER_MIRROR.code()).setName(topicName); + topicRes.setErrorCode(Errors.TOPIC_ALREADY_IN_CLUSTER_MIRROR.code()).setName(topicName) + .setErrorMessage("Topic '" + topicName + "' is already in mirror '" + currMirrorNameValue + + "' in " + MirrorPartitionState.fromValue((byte) currMirrorStateChange) + " state"); topicResList.add(topicRes); continue; } From d2fd01b703889f56b3b4bf69d99848c1d971b9bc Mon Sep 17 00:00:00 2001 From: Federico Valeri Date: Wed, 17 Jun 2026 18:52:12 +0200 Subject: [PATCH 3/3] Rename TopicData to TopicMetadata Signed-off-by: Federico Valeri --- .../kafka/clients/admin/KafkaAdminClient.java | 14 ++++++------- .../requests/PauseMirrorTopicsRequest.java | 2 +- .../requests/ResumeMirrorTopicsRequest.java | 2 +- .../requests/StartMirrorTopicsRequest.java | 2 +- .../requests/StopMirrorTopicsRequest.java | 2 +- .../message/PauseMirrorTopicsRequest.json | 4 ++-- .../message/ReadMirrorStatesRequest.json | 2 +- .../message/ResumeMirrorTopicsRequest.json | 2 +- .../message/StartMirrorTopicsRequest.json | 2 +- .../message/StopMirrorTopicsRequest.json | 2 +- .../message/WriteMirrorStatesRequest.json | 4 ++-- .../common/requests/RequestResponseTest.java | 20 +++++++++---------- .../server/mirror/MirrorMetadataManager.java | 14 ++++++------- 13 files changed, 36 insertions(+), 36 deletions(-) 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)));