From 233b2c2978de6ea5f4f95cc883dcf0a59ffaa24d Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Mon, 11 Aug 2025 15:43:45 +0800 Subject: [PATCH 1/7] init commit --- .../reassign/PartitionReassignmentState.java | 7 +++++++ .../tools/reassign/VerifyAssignmentResult.java | 8 ++++++++ .../reassign/ReassignPartitionsCommandTest.java | 17 +++++++++++------ 3 files changed, 26 insertions(+), 6 deletions(-) diff --git a/tools/src/main/java/org/apache/kafka/tools/reassign/PartitionReassignmentState.java b/tools/src/main/java/org/apache/kafka/tools/reassign/PartitionReassignmentState.java index 53ae960f863db..644064647768b 100644 --- a/tools/src/main/java/org/apache/kafka/tools/reassign/PartitionReassignmentState.java +++ b/tools/src/main/java/org/apache/kafka/tools/reassign/PartitionReassignmentState.java @@ -54,4 +54,11 @@ public boolean equals(Object o) { public int hashCode() { return Objects.hash(currentReplicas, targetReplicas, done); } + + @Override + public String toString() { + return "{currentReplicas=" + currentReplicas + ", " + + "targetReplicas=" + targetReplicas + ", " + + "done=" + done; + } } diff --git a/tools/src/main/java/org/apache/kafka/tools/reassign/VerifyAssignmentResult.java b/tools/src/main/java/org/apache/kafka/tools/reassign/VerifyAssignmentResult.java index 9dbde0fd1af9a..28a5ca10afce2 100644 --- a/tools/src/main/java/org/apache/kafka/tools/reassign/VerifyAssignmentResult.java +++ b/tools/src/main/java/org/apache/kafka/tools/reassign/VerifyAssignmentResult.java @@ -66,4 +66,12 @@ public boolean equals(Object o) { public int hashCode() { return Objects.hash(partStates, partsOngoing, moveStates, movesOngoing); } + + @Override + public String toString() { + return "partStates=" + partStates + ", " + + "partsOngoing=" + partsOngoing + ", " + + "moveStates=" + moveStates + ", " + + "movesOngoing=" + movesOngoing; + } } diff --git a/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommandTest.java b/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommandTest.java index 41b35179c18fc..9c2cc6badb6fa 100644 --- a/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommandTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommandTest.java @@ -153,13 +153,18 @@ public void testGenerateAssignmentWithBootstrapServer() throws Exception { produceMessages(foo0.topic(), foo0.partition(), 100); try (Admin admin = Admin.create(Map.of(CommonClientConfigs.BOOTSTRAP_SERVERS_CONFIG, clusterInstance.bootstrapServers()))) { - String assignment = "{\"version\":1,\"partitions\":" + - "[{\"topic\":\"foo\",\"partition\":0,\"replicas\":[3,1,2],\"log_dirs\":[\"any\",\"any\",\"any\"]}" + - "]}"; - generateAssignment(admin, assignment, "1,2,3", false); + String topicsToMoveJson = "{\"topics\":\n\t[{\"topic\": \"foo\"}],\n\"version\":1\n}"; + var assignment = generateAssignment(admin, topicsToMoveJson, "1,2,3", false); + Map> proposedAssignments = assignment.getKey(); + String assignmentJson = String.format("{\"version\":1,\"partitions\":" + + "[{\"topic\":\"foo\",\"partition\":0,\"replicas\":%s,\"log_dirs\":[\"any\",\"any\",\"any\"]}" + + "]}", proposedAssignments.get(foo0)); + + runExecuteAssignment(false, assignmentJson, -1L, -1L); + Map finalAssignment = Map.of(foo0, - new PartitionReassignmentState(List.of(0, 1, 2), List.of(3, 1, 2), true)); - waitForVerifyAssignment(admin, assignment, false, + new PartitionReassignmentState(proposedAssignments.get(foo0), proposedAssignments.get(foo0), true)); + waitForVerifyAssignment(admin, assignmentJson, false, new VerifyAssignmentResult(finalAssignment)); } } From 8709a26b5929f28075cd0e07fb27f9b24c966a90 Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Mon, 11 Aug 2025 15:47:43 +0800 Subject: [PATCH 2/7] fix PartitionReassignmentState.toString --- .../apache/kafka/tools/reassign/PartitionReassignmentState.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tools/src/main/java/org/apache/kafka/tools/reassign/PartitionReassignmentState.java b/tools/src/main/java/org/apache/kafka/tools/reassign/PartitionReassignmentState.java index 644064647768b..c36fd33c785d0 100644 --- a/tools/src/main/java/org/apache/kafka/tools/reassign/PartitionReassignmentState.java +++ b/tools/src/main/java/org/apache/kafka/tools/reassign/PartitionReassignmentState.java @@ -57,7 +57,7 @@ public int hashCode() { @Override public String toString() { - return "{currentReplicas=" + currentReplicas + ", " + return "currentReplicas=" + currentReplicas + ", " + "targetReplicas=" + targetReplicas + ", " + "done=" + done; } From 6e181b66a5edfa650dfc656bae4e30ce972afc25 Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Mon, 11 Aug 2025 16:18:59 +0800 Subject: [PATCH 3/7] make record --- .../reassign/PartitionReassignmentState.java | 44 ++------------- .../reassign/ReassignPartitionsCommand.java | 8 +-- .../reassign/VerifyAssignmentResult.java | 56 ++++--------------- .../ReassignPartitionsCommandTest.java | 14 ++--- 4 files changed, 26 insertions(+), 96 deletions(-) diff --git a/tools/src/main/java/org/apache/kafka/tools/reassign/PartitionReassignmentState.java b/tools/src/main/java/org/apache/kafka/tools/reassign/PartitionReassignmentState.java index c36fd33c785d0..bcdcecca8f124 100644 --- a/tools/src/main/java/org/apache/kafka/tools/reassign/PartitionReassignmentState.java +++ b/tools/src/main/java/org/apache/kafka/tools/reassign/PartitionReassignmentState.java @@ -18,47 +18,13 @@ package org.apache.kafka.tools.reassign; import java.util.List; -import java.util.Objects; /** * The state of a partition reassignment. The current replicas and target replicas * may overlap. */ -final class PartitionReassignmentState { - public final List currentReplicas; - - public final List targetReplicas; - - public final boolean done; - - /** - * @param currentReplicas The current replicas. - * @param targetReplicas The target replicas. - * @param done True if the reassignment is done. - */ - public PartitionReassignmentState(List currentReplicas, List targetReplicas, boolean done) { - this.currentReplicas = currentReplicas; - this.targetReplicas = targetReplicas; - this.done = done; - } - - @Override - public boolean equals(Object o) { - if (this == o) return true; - if (o == null || getClass() != o.getClass()) return false; - PartitionReassignmentState state = (PartitionReassignmentState) o; - return done == state.done && Objects.equals(currentReplicas, state.currentReplicas) && Objects.equals(targetReplicas, state.targetReplicas); - } - - @Override - public int hashCode() { - return Objects.hash(currentReplicas, targetReplicas, done); - } - - @Override - public String toString() { - return "currentReplicas=" + currentReplicas + ", " - + "targetReplicas=" + targetReplicas + ", " - + "done=" + done; - } -} +public record PartitionReassignmentState( + List currentReplicas, + List targetReplicas, + boolean done +) { } diff --git a/tools/src/main/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommand.java b/tools/src/main/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommand.java index 334f0738ca363..591c127ae2f3b 100644 --- a/tools/src/main/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommand.java +++ b/tools/src/main/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommand.java @@ -279,12 +279,12 @@ static String partitionReassignmentStatesToString(Map { PartitionReassignmentState state = states.get(topicPartition); - if (state.done) { - if (state.currentReplicas.equals(state.targetReplicas)) { + if (state.done()) { + if (state.currentReplicas().equals(state.targetReplicas())) { bld.add(String.format("Reassignment of partition %s is completed.", topicPartition)); } else { - String currentReplicaStr = state.currentReplicas.stream().map(String::valueOf).collect(Collectors.joining(",")); - String targetReplicaStr = state.targetReplicas.stream().map(String::valueOf).collect(Collectors.joining(",")); + String currentReplicaStr = state.currentReplicas().stream().map(String::valueOf).collect(Collectors.joining(",")); + String targetReplicaStr = state.targetReplicas().stream().map(String::valueOf).collect(Collectors.joining(",")); bld.add("There is no active reassignment of partition " + topicPartition + ", " + "but replica set is " + currentReplicaStr + " rather than " + diff --git a/tools/src/main/java/org/apache/kafka/tools/reassign/VerifyAssignmentResult.java b/tools/src/main/java/org/apache/kafka/tools/reassign/VerifyAssignmentResult.java index 28a5ca10afce2..0232eb3d1c1aa 100644 --- a/tools/src/main/java/org/apache/kafka/tools/reassign/VerifyAssignmentResult.java +++ b/tools/src/main/java/org/apache/kafka/tools/reassign/VerifyAssignmentResult.java @@ -21,57 +21,21 @@ import org.apache.kafka.common.TopicPartitionReplica; import java.util.Map; -import java.util.Objects; /** * A result returned from verifyAssignment. + * @param partStates A map from partitions to reassignment states. + * @param partsOngoing True if there are any ongoing partition reassignments. + * @param moveStates A map from log directories to movement states. + * @param movesOngoing True if there are any ongoing moves that we know about. */ -public final class VerifyAssignmentResult { - public final Map partStates; - public final boolean partsOngoing; - public final Map moveStates; - public final boolean movesOngoing; - +public record VerifyAssignmentResult( + Map partStates, + boolean partsOngoing, + Map moveStates, + boolean movesOngoing +) { public VerifyAssignmentResult(Map partStates) { this(partStates, false, Map.of(), false); } - - /** - * @param partStates A map from partitions to reassignment states. - * @param partsOngoing True if there are any ongoing partition reassignments. - * @param moveStates A map from log directories to movement states. - * @param movesOngoing True if there are any ongoing moves that we know about. - */ - public VerifyAssignmentResult( - Map partStates, - boolean partsOngoing, - Map moveStates, - boolean movesOngoing - ) { - this.partStates = partStates; - this.partsOngoing = partsOngoing; - this.moveStates = moveStates; - this.movesOngoing = movesOngoing; - } - - @Override - public boolean equals(Object o) { - if (this == o) return true; - if (o == null || getClass() != o.getClass()) return false; - VerifyAssignmentResult that = (VerifyAssignmentResult) o; - return partsOngoing == that.partsOngoing && movesOngoing == that.movesOngoing && Objects.equals(partStates, that.partStates) && Objects.equals(moveStates, that.moveStates); - } - - @Override - public int hashCode() { - return Objects.hash(partStates, partsOngoing, moveStates, movesOngoing); - } - - @Override - public String toString() { - return "partStates=" + partStates + ", " - + "partsOngoing=" + partsOngoing + ", " - + "moveStates=" + moveStates + ", " - + "movesOngoing=" + movesOngoing; - } } diff --git a/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommandTest.java b/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommandTest.java index 9c2cc6badb6fa..fa2cb136a3d5d 100644 --- a/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommandTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommandTest.java @@ -242,15 +242,15 @@ public void testThrottledReassignment() throws Exception { // Check the reassignment status. VerifyAssignmentResult result = runVerifyAssignment(admin, assignment, true); - if (!result.partsOngoing) { + if (!result.partsOngoing()) { return true; } else { assertFalse( - result.partStates.values().stream().allMatch(state -> state.done), + result.partStates().values().stream().allMatch(PartitionReassignmentState::done), "Expected at least one partition reassignment to be ongoing when result = " + result ); - assertEquals(List.of(0, 3, 2), result.partStates.get(new TopicPartition("foo", 0)).targetReplicas); - assertEquals(List.of(3, 2, 1), result.partStates.get(new TopicPartition("baz", 2)).targetReplicas); + assertEquals(List.of(0, 3, 2), result.partStates().get(new TopicPartition("foo", 0)).targetReplicas()); + assertEquals(List.of(3, 2, 1), result.partStates().get(new TopicPartition("baz", 2)).targetReplicas()); waitForInterBrokerThrottle(admin, List.of(0, 1, 2, 3), interBrokerThrottle); return false; } @@ -389,7 +389,7 @@ public void testCancellationWithAddingReplicaInIsr() throws Exception { /** * Test moving partitions between directories. */ - @ClusterTest(types = {Type.KRAFT}) + @ClusterTest(types = {Type.KRAFT, Type.CO_KRAFT}) public void testLogDirReassignment() throws Exception { createTopics(); TopicPartition topicPartition = new TopicPartition("foo", 0); @@ -545,7 +545,7 @@ private void executeAndVerifyReassignment() throws InterruptedException { finalAssignment.put(bar0, new PartitionReassignmentState(List.of(3, 2, 0), List.of(3, 2, 0), true)); VerifyAssignmentResult verifyAssignmentResult = runVerifyAssignment(admin, assignment, false); - assertFalse(verifyAssignmentResult.movesOngoing); + assertFalse(verifyAssignmentResult.movesOngoing()); // Wait for the assignment to complete waitForVerifyAssignment(admin, assignment, false, @@ -791,7 +791,7 @@ private void testCancellationAction(boolean useBootstrapServer) throws Interrupt // This time, the broker throttles were removed. waitForBrokerLevelThrottles(admin, unthrottledBrokerConfigs); // Verify that there are no ongoing reassignments. - assertFalse(runVerifyAssignment(admin, assignment, false).partsOngoing); + assertFalse(runVerifyAssignment(admin, assignment, false).partsOngoing()); } // Verify that the partition is removed from cancelled replicas verifyReplicaDeleted(new TopicPartitionReplica(foo0.topic(), foo0.partition(), 3)); From 240ce9c3fbcfc93b9d4705c96815c81637abcda9 Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Mon, 11 Aug 2025 16:22:58 +0800 Subject: [PATCH 4/7] revert --- .../kafka/tools/reassign/ReassignPartitionsCommandTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommandTest.java b/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommandTest.java index fa2cb136a3d5d..054cd04bb7e53 100644 --- a/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommandTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommandTest.java @@ -389,7 +389,7 @@ public void testCancellationWithAddingReplicaInIsr() throws Exception { /** * Test moving partitions between directories. */ - @ClusterTest(types = {Type.KRAFT, Type.CO_KRAFT}) + @ClusterTest(types = {Type.KRAFT}) public void testLogDirReassignment() throws Exception { createTopics(); TopicPartition topicPartition = new TopicPartition("foo", 0); From 18116c5eb2f7e1c4a6539871bec2829dd0751618 Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Mon, 11 Aug 2025 16:29:15 +0800 Subject: [PATCH 5/7] remove public identify --- .../apache/kafka/tools/reassign/PartitionReassignmentState.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tools/src/main/java/org/apache/kafka/tools/reassign/PartitionReassignmentState.java b/tools/src/main/java/org/apache/kafka/tools/reassign/PartitionReassignmentState.java index bcdcecca8f124..9d3ae1defc2d5 100644 --- a/tools/src/main/java/org/apache/kafka/tools/reassign/PartitionReassignmentState.java +++ b/tools/src/main/java/org/apache/kafka/tools/reassign/PartitionReassignmentState.java @@ -23,7 +23,7 @@ * The state of a partition reassignment. The current replicas and target replicas * may overlap. */ -public record PartitionReassignmentState( +record PartitionReassignmentState( List currentReplicas, List targetReplicas, boolean done From 727efd9c604488e2e978893f802c6bba4b8f84b5 Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Mon, 11 Aug 2025 16:45:38 +0800 Subject: [PATCH 6/7] address comments --- .../kafka/tools/reassign/PartitionReassignmentState.java | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/tools/src/main/java/org/apache/kafka/tools/reassign/PartitionReassignmentState.java b/tools/src/main/java/org/apache/kafka/tools/reassign/PartitionReassignmentState.java index 9d3ae1defc2d5..a81b430c8b3e3 100644 --- a/tools/src/main/java/org/apache/kafka/tools/reassign/PartitionReassignmentState.java +++ b/tools/src/main/java/org/apache/kafka/tools/reassign/PartitionReassignmentState.java @@ -22,6 +22,10 @@ /** * The state of a partition reassignment. The current replicas and target replicas * may overlap. + * + * @param currentReplicas The current replicas. + * @param targetReplicas The target replicas. + * @param done True if the reassignment is done. */ record PartitionReassignmentState( List currentReplicas, From afe5c66540aee7bea829aa7496f16e15358e58d6 Mon Sep 17 00:00:00 2001 From: TaiJu Wu Date: Mon, 11 Aug 2025 17:40:16 +0800 Subject: [PATCH 7/7] address comments --- .../ReassignPartitionsCommandTest.java | 25 ++++++++++++++++--- 1 file changed, 21 insertions(+), 4 deletions(-) diff --git a/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommandTest.java b/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommandTest.java index 054cd04bb7e53..f1726721bfddf 100644 --- a/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommandTest.java +++ b/tools/src/test/java/org/apache/kafka/tools/reassign/ReassignPartitionsCommandTest.java @@ -153,12 +153,29 @@ public void testGenerateAssignmentWithBootstrapServer() throws Exception { produceMessages(foo0.topic(), foo0.partition(), 100); try (Admin admin = Admin.create(Map.of(CommonClientConfigs.BOOTSTRAP_SERVERS_CONFIG, clusterInstance.bootstrapServers()))) { - String topicsToMoveJson = "{\"topics\":\n\t[{\"topic\": \"foo\"}],\n\"version\":1\n}"; + String topicsToMoveJson = """ + { + "topics": [ + { "topic": "foo" } + ], + "version": 1 + } + """; var assignment = generateAssignment(admin, topicsToMoveJson, "1,2,3", false); Map> proposedAssignments = assignment.getKey(); - String assignmentJson = String.format("{\"version\":1,\"partitions\":" + - "[{\"topic\":\"foo\",\"partition\":0,\"replicas\":%s,\"log_dirs\":[\"any\",\"any\",\"any\"]}" + - "]}", proposedAssignments.get(foo0)); + String assignmentJson = String.format(""" + { + "version": 1, + "partitions": [ + { + "topic": "foo", + "partition": 0, + "replicas": %s, + "log_dirs": ["any", "any", "any"] + } + ] + } + """, proposedAssignments.get(foo0)); runExecuteAssignment(false, assignmentJson, -1L, -1L);