How does a partition reassignment preserve availability and avoid data loss while replicas are moving?
answer
- assignment = union(old, new) during move
- new replicas join ISR before old removed
- leadership only to ISR member → no acked loss
- risk is load: under-replicated + disk peak
- watch UnderReplicatedPartitions
basics
~20 sKafka adds the new target replicas and waits for them to catch up and join the ISR before removing the old replicas, so the partition always keeps enough in-sync copies. Leadership only moves to a new replica once it is in the ISR, so no acknowledged data is lost.
solid answer
~50 sDuring reassignment the controller computes a transitional replica set that is the union of the old and new assignments. New replicas are added as followers and begin catching up; only once a new replica is fully caught up does it join the ISR. The old replicas are not removed until the new ones are in-sync, so the partition never drops below its intended replication during the move — producers with acks=all still get min.insync.replicas honored. Leadership is only transferred to a replica that is in the ISR, so no acknowledged records are lost. Risks come from the *load* of moving, not correctness: unthrottled moves can saturate I/O and push replicas out of the ISR (under-replicated partitions), so throttling and bounded concurrency matter. A too-aggressive plan across many partitions at once can also exhaust disk on destination brokers — capacity goals (Cruise Control) or careful batching mitigate this.
go deeper
Know that moving replicas keeps the partition available — old copies stay until new ones are ready.
Explain that new replicas must join the ISR before old ones are removed, so no acked data is lost.
Describe the union transitional set, ISR-only leadership, and the under-replication/disk-peak load risks.
Reason about concurrency/throttle policy, destination disk headroom, unclean-election interaction, and monitoring during large drains.
## The core invariant Kafka's durability rests on the **ISR** (in-sync replica set): a record acknowledged with `acks=all` is only acked once it is replicated to all members of the ISR, and the cluster refuses to lose acked data by only ever electing leaders from the ISR (unless unclean leader election is enabled). A reassignment must not break this. ## Transitional replica set (union of old + new) Suppose partition P moves from replicas `[1,2,3]` to `[4,5,6]`. The controller does **not** swap them atomically. It first sets the assignment to the **union** `[1,2,3,4,5,6]`: 1. Brokers 4, 5, 6 are added as **followers** and start fetching from the leader, replaying the log from the start (or from log-start-offset). 2. As each new replica fully catches up to the leader's log end, it **joins the ISR**. 3. Once the new replicas are in-sync, the controller transfers leadership if needed (only to an in-sync replica) and then **removes** the old replicas 1, 2, 3 from the assignment. Because old replicas stay until new ones are in-sync, **at every instant the partition has at least its original number of in-sync replicas** — availability and `min.insync.replicas` guarantees hold throughout. ## Why no data loss - Leadership only ever moves to an **ISR member**, which by definition has all acknowledged records. - Removing a not-yet-leader old replica loses nothing because the data exists on the in-sync new replicas. - The dangerous setting is `unclean.leader.election.enable=true`, which would allow electing an out-of-sync replica — orthogonal to reassignment but a real loss vector if combined with an ill-timed failure. ## Where the real risk lives: load, not correctness - **I/O saturation:** the catch-up fetch is heavy. If it starves normal replication, healthy replicas fall behind and **under-replicated partitions** spike — exactly what throttles (`leader/follower.replication.throttled.rate`) prevent. - **Destination disk exhaustion:** the union state means destination brokers temporarily hold the new replica *while* sources still hold the old — peak disk is higher than the final state. Moving too many large partitions concurrently can fill disks; bound with concurrency limits or Cruise Control capacity (hard) goals. - **Controller/metadata load:** very large plans churn metadata; batch them. ## Operational safeguards - Throttle so moves run as background work. - Limit concurrent reassignments (`--max-...` style caps / Cruise Control concurrency). - Keep a rollback plan (the original assignment from `--generate`). - Use `--cancel` if a move misbehaves. - Monitor `UnderReplicatedPartitions`, `UnderMinIsrPartitionCount`, and replication MB/s during the move.
- Why can destination brokers temporarily need more disk during a reassignment than the final layout requires?During the move the assignment is the union of old and new replicas: destinations hold the new replica while sources still hold the old. Peak disk usage exceeds the steady-state footprint until the old replicas are dropped, so moving many large partitions at once can fill disks.
- Reassignment preserves acked data, but what setting could still cause loss if a failure hits mid-move?unclean.leader.election.enable=true. Reassignment itself only elects in-sync leaders, but if unclean election is enabled and the in-sync replicas are lost during the move, an out-of-sync replica could be elected and lose acknowledged records.
saying these in an interview costs you the question
- Saying reassignment briefly drops the partition below its replication while swapping — it uses the union set so it never does.
- Claiming reassignment itself risks data loss — correctness is preserved; the real risks are I/O saturation and disk peaks.
- Forgetting that destination disk peaks above final size because old+new coexist.