From 9082a19e2fa0eeeb5675f73d1a392da512a2194c Mon Sep 17 00:00:00 2001 From: Jorge Esteban Quilcate Otoya Date: Tue, 18 Aug 2026 13:26:15 +0300 Subject: [PATCH] fix(inkless:switch): stage reassignments of switched partitions changePartitionReassignment applied the target replica set in one step for any diskless topic, setting targetIsr = target.replicas(). For born-diskless that is safe: all data is in object storage, so a target replica is current on arrival. A partition switched from classic breaks that premise. Records below classicToDisklessStartOffset live only in the replicas' local logs, so a replica newly added by a reassignment entered ISR holding none of the classic prefix -- not a stale prefix, none -- and was immediately electable. Nothing gated the reassignment path on the switch state. Restrict the one-step path to born-diskless partitions. A switched partition now takes the staged path: the target enters as addingReplicas and earns ISR through AlterPartition once its catch-up fetcher reaches the seal, which applyLocalFollowersDelta already arms for a newly added replica of a switched topic. Operator-visible consequence: reassigning a switched topic is pending rather than instant, and listPartitionReassignments reports it in progress until the new replica catches up. That is the point -- the wait is the evidence. This closes the grow half only. A pure RF shrink has no adding replicas, so completeReassignmentIfNeeded finishes the reassignment in the same operation and replicas collapses to the target before maybePopulateTargetElr runs, leaving ELR outside the replica set exactly as before (measured: elr=[1,2] replicas=[0] for a [0,1,2] -> [0] shrink at minISR=2). That half is not diskless-specific -- classic topics escape it only because RF is fixed without KIP-860 allow-RF-change -- and needs Rule B restricted to the target replica set instead, so it is left alone here. Tests cover the staged grow through to completion; cancellation, which unwinds adding/removing state on a diskless partition for the first time; and the consolidating variant, which stages the same way because its follower reaches the seal on the classic fetcher before isReadyForConsolidation hands off, so the leader still observes the fetch state that expands ISR. --- .../controller/ReplicationControlManager.java | 13 +- .../ReplicationControlManagerInklessTest.java | 169 ++++++++++++++++++ 2 files changed, 180 insertions(+), 2 deletions(-) diff --git a/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java b/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java index ef3a5bc6ca..9d95237d51 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java @@ -2750,6 +2750,8 @@ void generateLeaderAndIsrUpdates(String context, // replica set with rack-aware placement, and because the data lives in object storage a // replica-set resize is immediate and safe (no inter-broker catch-up). Honor the // KIP-860 allow-RF-change flag for managed diskless topics, exactly like classic topics. + // The resize is only immediate for born-diskless partitions; a switched one still has a + // classic prefix to replicate, so it stages the change (see changePartitionReassignment). boolean rfChangeAllowedForTopic = !isDisklessTopic(topic.name()) || isDisklessManagedReplicasEnabled; boolean effectiveRFChange = allowRFChange && rfChangeAllowedForTopic; @@ -2889,7 +2891,14 @@ Optional changePartitionReassignment(TopicIdPartition tp, new PartitionReassignmentReplicas(currentAssignment, targetAssignment); String topicName = topics.get(tp.topicId()).name; - boolean isDiskless = isDisklessTopic(topicName); + // Only born-diskless partitions take the one-step path below. A partition switched from classic + // still has records under classicToDisklessStartOffset in the replicas' local logs, so a target + // replica is not interchangeable and has to earn ISR through the staged path. + // -1 means "no committed seal", which also covers a switch aborted via AlterDisklessSwitch -- + // that partition keeps a classic prefix, but every other seal reader already routes it as + // born-diskless, so the prefix is unreachable regardless of ISR membership. + boolean isBornDiskless = isDisklessTopic(topicName) && + part.classicToDisklessStartOffset == PartitionRegistration.NO_CLASSIC_TO_DISKLESS_START_OFFSET; PartitionChangeBuilder builder = new PartitionChangeBuilder( part, @@ -2901,7 +2910,7 @@ Optional changePartitionReassignment(TopicIdPartition tp, ); builder.setEligibleLeaderReplicasEnabled(featureControl.isElrFeatureEnabled()); - if (isDiskless) { + if (isBornDiskless) { // Diskless: data is in object storage, no replica sync needed. // Apply target replicas directly — skip the staged adding/removing process // (no addingReplicas/removingReplicas). This is safe because: diff --git a/metadata/src/test/java/org/apache/kafka/controller/ReplicationControlManagerInklessTest.java b/metadata/src/test/java/org/apache/kafka/controller/ReplicationControlManagerInklessTest.java index 24db2f7ab8..dab1c3d312 100644 --- a/metadata/src/test/java/org/apache/kafka/controller/ReplicationControlManagerInklessTest.java +++ b/metadata/src/test/java/org/apache/kafka/controller/ReplicationControlManagerInklessTest.java @@ -34,6 +34,7 @@ import org.apache.kafka.common.message.AlterPartitionReassignmentsResponseData; import org.apache.kafka.common.message.AlterPartitionReassignmentsResponseData.ReassignablePartitionResponse; import org.apache.kafka.common.message.AlterPartitionReassignmentsResponseData.ReassignableTopicResponse; +import org.apache.kafka.common.message.AlterPartitionRequestData.BrokerState; import org.apache.kafka.common.message.CreatePartitionsRequestData.CreatePartitionsAssignment; import org.apache.kafka.common.message.CreatePartitionsRequestData.CreatePartitionsTopic; import org.apache.kafka.common.message.CreatePartitionsResponseData.CreatePartitionsTopicResult; @@ -62,6 +63,7 @@ import org.apache.kafka.metadata.Replicas; import org.apache.kafka.server.common.ApiMessageAndVersion; import org.apache.kafka.server.common.MetadataVersion; +import org.apache.kafka.server.common.TopicIdPartition; import org.apache.kafka.server.config.ServerConfigs; import org.junit.jupiter.api.Nested; @@ -1797,6 +1799,173 @@ public void testDecreaseReplicationFactorForManagedDisklessTopic() { assertEquals(List.of(0), Replicas.toList(partition.isr)); } + @Test + public void testReassignSwitchedPartitionStagesTargetInsteadOfGrantingIsr() { + ReplicationControlTestContext ctx = new ReplicationControlTestContext.Builder() + .setMetadataVersion(MetadataVersion.latestTesting()) + .setDisklessStorageSystemEnabled(true) + .setDisklessManagedReplicasEnabled(true) + .build(); + + ReplicationControlManager replication = ctx.replicationControl; + ctx.registerBrokers(0, 1, 2); + ctx.unfenceBrokers(0, 1, 2); + + // Classic topic switched to diskless with a committed seal, so records below offset 100 + // exist only in the local logs of brokers 0 and 1. + String topic = "switched"; + Uuid topicId = ctx.createTestTopic(topic, new int[][] {new int[] {0, 1}}, Map.of(), (short) 0) + .topicId(); + ctx.alterTopicConfig(topic, DISKLESS_ENABLE_CONFIG, "true"); + setClassicToDisklessStartOffset(ctx, topicId, 100L); + + // Swap broker 1 for broker 2. Broker 2 holds none of the classic prefix. + ControllerResult alterResult = + replication.alterPartitionReassignments( + new AlterPartitionReassignmentsRequestData().setTopics(List.of( + new ReassignableTopic().setName(topic).setPartitions(List.of( + new ReassignablePartition().setPartitionIndex(0) + .setReplicas(List.of(0, 2))))))); + ctx.replay(alterResult.records()); + + PartitionRegistration partition = replication.getPartition(topicId, 0); + assertFalse(Replicas.contains(partition.isr, 2), + "Broker 2 must not be granted ISR by the reassignment: it holds none of the classic " + + "prefix below the seal and has to earn ISR through AlterPartition"); + assertTrue(Replicas.contains(partition.addingReplicas, 2), + "Broker 2 should be staged as an adding replica"); + assertTrue(Replicas.contains(partition.removingReplicas, 1), + "Broker 1 should be staged as a removing replica, not dropped immediately"); + assertNotEquals(NONE_REASSIGNING, replication.listPartitionReassignments(List.of( + new ListPartitionReassignmentsTopics().setName(topic) + .setPartitionIndexes(List.of(0))), Long.MAX_VALUE), + "The reassignment stays in progress until the new replica catches up"); + assertEquals(100L, partition.classicToDisklessStartOffset, + "The seal should be untouched by the reassignment"); + + // Once broker 2 has caught up and the leader reports it in sync, the reassignment must + // complete -- otherwise routing switched partitions through the staged path would leave + // them reassigning forever. + ctx.alterPartition(new TopicIdPartition(topicId, 0), 0, + List.of( + new BrokerState().setBrokerId(0).setBrokerEpoch(defaultBrokerEpoch(0)), + new BrokerState().setBrokerId(1).setBrokerEpoch(defaultBrokerEpoch(1)), + new BrokerState().setBrokerId(2).setBrokerEpoch(defaultBrokerEpoch(2))), + LeaderRecoveryState.RECOVERED); + + PartitionRegistration completed = replication.getPartition(topicId, 0); + assertEquals(List.of(0, 2), Replicas.toList(completed.replicas), + "Reassignment should have completed to the target replica set"); + assertEquals(List.of(), Replicas.toList(completed.addingReplicas)); + assertEquals(List.of(), Replicas.toList(completed.removingReplicas)); + assertEquals(NONE_REASSIGNING, replication.listPartitionReassignments(List.of( + new ListPartitionReassignmentsTopics().setName(topic) + .setPartitionIndexes(List.of(0))), Long.MAX_VALUE)); + } + + @Test + public void testCancelReassignmentOfSwitchedPartitionRevertsStagedState() { + ReplicationControlTestContext ctx = new ReplicationControlTestContext.Builder() + .setMetadataVersion(MetadataVersion.latestTesting()) + .setDisklessStorageSystemEnabled(true) + .setDisklessManagedReplicasEnabled(true) + .build(); + + ReplicationControlManager replication = ctx.replicationControl; + ctx.registerBrokers(0, 1, 2); + ctx.unfenceBrokers(0, 1, 2); + + String topic = "switched"; + Uuid topicId = ctx.createTestTopic(topic, new int[][] {new int[] {0, 1}}, Map.of(), (short) 0) + .topicId(); + ctx.alterTopicConfig(topic, DISKLESS_ENABLE_CONFIG, "true"); + setClassicToDisklessStartOffset(ctx, topicId, 100L); + + ControllerResult alterResult = + replication.alterPartitionReassignments( + new AlterPartitionReassignmentsRequestData().setTopics(List.of( + new ReassignableTopic().setName(topic).setPartitions(List.of( + new ReassignablePartition().setPartitionIndex(0) + .setReplicas(List.of(0, 2))))))); + ctx.replay(alterResult.records()); + assertTrue(Replicas.contains(replication.getPartition(topicId, 0).addingReplicas, 2), + "precondition: the reassignment must be staged, not applied in one step"); + + // Switched partitions are the first diskless partitions to carry real adding/removing + // state, so cancellation has to unwind it through the standard revert path. + ControllerResult cancelResult = + replication.alterPartitionReassignments( + new AlterPartitionReassignmentsRequestData().setTopics(List.of( + new ReassignableTopic().setName(topic).setPartitions(List.of( + new ReassignablePartition().setPartitionIndex(0) + .setReplicas(null)))))); + ctx.replay(cancelResult.records()); + + PartitionRegistration reverted = replication.getPartition(topicId, 0); + assertEquals(List.of(0, 1), Replicas.toList(reverted.replicas), + "Cancelling should revert to the original replica set"); + assertEquals(List.of(), Replicas.toList(reverted.addingReplicas)); + assertEquals(List.of(), Replicas.toList(reverted.removingReplicas)); + assertFalse(Replicas.contains(reverted.isr, 2), + "Broker 2 must not be left in ISR after cancellation"); + assertEquals(NONE_REASSIGNING, replication.listPartitionReassignments(List.of( + new ListPartitionReassignmentsTopics().setName(topic) + .setPartitionIndexes(List.of(0))), Long.MAX_VALUE)); + assertEquals(100L, reverted.classicToDisklessStartOffset, + "The seal should be untouched by cancellation"); + } + + @Test + public void testReassignConsolidatingSwitchedPartitionAlsoStages() { + // A consolidating switched partition must stage like any other switched one. Its follower + // reaches the seal on the classic fetcher before isReadyForConsolidation hands off to the + // consolidation fetcher, so the leader still observes the fetch state that expands ISR. + ReplicationControlTestContext ctx = new ReplicationControlTestContext.Builder() + .setMetadataVersion(MetadataVersion.latestTesting()) + .setDisklessStorageSystemEnabled(true) + .setDisklessManagedReplicasEnabled(true) + .setDisklessRemoteStorageConsolidationEnabled(true) + .build(); + + ReplicationControlManager replication = ctx.replicationControl; + ctx.registerBrokers(0, 1, 2); + ctx.unfenceBrokers(0, 1, 2); + + String topic = "switched-consolidating"; + Uuid topicId = ctx.createTestTopic(topic, new int[][] {new int[] {0, 1}}, Map.of(), (short) 0) + .topicId(); + ctx.alterTopicConfig(topic, DISKLESS_ENABLE_CONFIG, "true"); + ctx.alterTopicConfig(topic, REMOTE_LOG_STORAGE_ENABLE_CONFIG, "true"); + setClassicToDisklessStartOffset(ctx, topicId, 100L); + + ControllerResult alterResult = + replication.alterPartitionReassignments( + new AlterPartitionReassignmentsRequestData().setTopics(List.of( + new ReassignableTopic().setName(topic).setPartitions(List.of( + new ReassignablePartition().setPartitionIndex(0) + .setReplicas(List.of(0, 2))))))); + ctx.replay(alterResult.records()); + + PartitionRegistration staged = replication.getPartition(topicId, 0); + assertFalse(Replicas.contains(staged.isr, 2), + "Consolidation must not restore the one-step ISR grant"); + assertTrue(Replicas.contains(staged.addingReplicas, 2)); + + ctx.alterPartition(new TopicIdPartition(topicId, 0), 0, + List.of( + new BrokerState().setBrokerId(0).setBrokerEpoch(defaultBrokerEpoch(0)), + new BrokerState().setBrokerId(1).setBrokerEpoch(defaultBrokerEpoch(1)), + new BrokerState().setBrokerId(2).setBrokerEpoch(defaultBrokerEpoch(2))), + LeaderRecoveryState.RECOVERED); + + PartitionRegistration completed = replication.getPartition(topicId, 0); + assertEquals(List.of(0, 2), Replicas.toList(completed.replicas)); + assertEquals(List.of(), Replicas.toList(completed.addingReplicas)); + assertEquals(NONE_REASSIGNING, replication.listPartitionReassignments(List.of( + new ListPartitionReassignmentsTopics().setName(topic) + .setPartitionIndexes(List.of(0))), Long.MAX_VALUE)); + } + @Test public void testReassignDisklessPartitionsToFencedBrokerIncludesInIsr() { MetadataVersion metadataVersion = MetadataVersion.latestTesting();