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..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 @@ -18,40 +18,17 @@ 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. + * + * @param currentReplicas The current replicas. + * @param targetReplicas The target replicas. + * @param done True if the reassignment is done. */ -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); - } -} +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 9dbde0fd1af9a..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,49 +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); - } } 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..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,13 +153,35 @@ 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": [ + { "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)); + + 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)); } } @@ -237,15 +259,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; } @@ -540,7 +562,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, @@ -786,7 +808,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));