Correcting leadership skew copies no bytes — so what does a leadership-only rebalance actually cost, and when should it run?
answer
- assignment change, not a migration
- no bandwidth, but interruptions
- the cost scales with units handed over
- home node present and caught up first
basics
~20 sA leadership-only rebalance changes which node serves each unit and moves no records, so it costs no bandwidth. It costs a brief interruption per unit as service hands over and clients refresh, which means the real decision is how many units to hand over at once, and when.
solid answer
~50 sThe correction hands leadership back to each unit's **home node** — the node first in its plan — and copies nothing, because the home node already holds a copy. That makes it unlike a data move: there is no bandwidth to budget and no hours-long tail. What it does cost is one short interruption per unit of ownership: writes to that unit fail or stall for the moment of handover, and clients discover the new serving node when their owner map misses. Those costs scale with the number of units handed over, and they land at the same instant if you correct a large skew in one sweep. So the preconditions are that each home node is present and its copy is caught up, and the operational choice is to stagger the sweep and to avoid running it while the home node is the component under suspicion.
go deeper
Recall that restoring leadership does not move records: the home node already holds a copy. The cost is a short pause for each unit as service changes hands, not hours of copying.
Explain the per-unit sequence — confirm the home node is caught up, hand service over, clients fail once and refresh their owner map — and why that makes the cost proportional to the number of units, not to their size.
Show the operational judgment: batch the handovers, verify the home nodes are present and current, watch the error rate of clients that cannot retry, and refuse to run it toward a node that is still the cause of the drift.
Set the standard: what preconditions a correction must meet, what batch size is acceptable against the client population's retry behaviour, and whether this operation is permitted during an incident or only outside one.
## What the correction is A **leadership-only rebalance** restores leadership for each **unit of ownership** — a partition, or on queue-shaped brokers a queue — to its **home node**, the node listed first in that unit's plan. It changes an assignment and nothing else. The home node already holds a copy of the unit; that is what made it the home node. So there is no byte of record data to transfer, which makes this a fundamentally different operation from moving the underlying data between nodes. That difference drives everything else about how it is scheduled. A data migration is rate-limited work measured in hours that competes with production for disks and links. This is a burst of small, fast state changes. ## What happens per unit 1. The home node is confirmed present and its copy is current enough to serve. 2. Service for that unit is handed from the current serving node to the home node, and the change is recorded. 3. Writers and readers holding a now-out-of-date **owner map** fail their next attempt against the old node, refresh, and reconnect to the new one. 4. The unit resumes, usually within the time a client takes to notice and retry. Each of those is short. Several hundred of them started at once is not: the interruptions are individually brief and collectively a spike. ## What it costs - **A brief per-unit interruption.** Between the old node stopping and the new node starting, writes for that one unit have nowhere to go. Well-behaved clients retry; clients with no retry, a short deadline or a fixed attempt count surface it as an error. - **A refresh burst.** Every client touching a handed-over unit refreshes its owner map at roughly the same moment, producing a spike of metadata requests on top of the reconnects. - **A visible latency bump.** The interruptions land in the tail: p99 moves even when nothing has failed. - **Nothing on the storage path.** No copy traffic, no disk contention, no bandwidth ceiling to set. This is the half people wrongly budget for. | | Leadership-only rebalance | A move of the underlying data | |---|---|---| | Records transferred | None | All of the unit's records | | Dominant cost | Brief interruption per unit | Sustained traffic on disks and links | | Duration | Seconds to minutes | Hours, typically | | Reversible | Immediately, by handing back | Only by another move | ## Choosing the moment Because there is no bandwidth cost, the instinct to hide the operation in the quietest hour is only half right — quiet hours also mean fewer clients affected, but the operation is short enough to run in daylight with people watching. What genuinely gates it: - **Every home node is present and caught up.** Handing a unit to a node whose copy is behind is either refused or, if forced, restores service at the cost of records — a different decision entirely, and one this correction should never be quietly turned into. - **Stagger, don't sweep.** Correcting a large skew in one action concentrates hundreds of interruptions and refreshes into one second. Handing back in batches spreads both. - **Not while the home node is the suspect.** If the drift exists because a node keeps failing, restoring leadership to it hands the traffic straight back to the component that caused the outage — and the next failure re-creates the skew plus another interruption. - **Not blindly during an incident.** The saturated node may be the symptom you are treating, but adding a wave of retries to a cluster already erroring is how a degradation becomes an outage. ## Where designs differ Some platforms run this restoration on their own, periodically, so an operator only ever sees the result. Some expose it purely as an operator action. Some make the handover nearly invisible to clients because ownership changes are pushed rather than discovered on failure, which shrinks the per-unit cost to almost nothing. And in designs where consumers compete for work from a shared queue, there is no single serving node per unit and this correction has no meaning. State the model you are describing rather than assuming the audience shares yours. ## What a strong answer sounds like "It is an assignment change, so budget interruptions, not bandwidth. Check that the home nodes are back and caught up, hand back in batches rather than in one sweep, watch the tail and the error rate of the clients that cannot retry, and do not hand traffic back to a node that is still the reason the skew exists."
- If the home node's copy is behind, should you hand leadership back anyway to relieve the saturated node?No. Handing a unit to a copy that is not caught up is a different decision with a different price — restoring service at the cost of records already acknowledged — and it should never be smuggled in as load balancing. Wait for the copy to become current, or relieve the hot node another way.
- Why does a single large correction hurt more than the same number of handovers spread out?Because the per-unit cost is an interruption, and interruptions add up when they coincide. Hundreds of units handed over together mean hundreds of clients failing, refreshing their owner map and reconnecting in the same second, which shows up as an error spike and a metadata burst that neither node was sized for.
- Can the correction be undone if it makes things worse?Effectively yes, and that is one of its virtues. Since nothing was copied, leadership can be handed back the other way at the same cost as handing it over — another brief interruption per unit. Compare that with a data migration, which can only be undone by moving the records back again.
saying these in an interview costs you the question
- It is a data migration, so it needs a long maintenance slot
- It copies no bytes, so no client notices anything
- It needs a bandwidth ceiling set before it runs
- Hand leadership to the home node whether or not it is caught up
- Run it during an incident to spread the load immediately
- Correct the whole skew in one action; smaller batches waste time