Design a leadership-balancing strategy for a large, latency-sensitive cluster during rolling restarts. What do you tune, and what are the failure modes?
answer
- Disable auto-rebalance during the roll
- controlled.shutdown for graceful failover
- Wait for ISR, then scoped PREFERRED election
- Stagger to avoid metadata thundering herd
- Never UNCLEAN for routine balance; throttle reassignments
basics
~20 sOften disable auto.leader.rebalance.enable and trigger preferred election deliberately after each restart batch, scoped via JSON and staggered, so leadership churn is controlled. Watch for client metadata storms, imbalance from out-of-sync preferred replicas, and unclean-election durability risks.
solid answer
~40 sFor a large latency-sensitive cluster I'd avoid uncontrolled churn: set auto.leader.rebalance.enable=false (or widen leader.imbalance.per.broker.percentage and lengthen the check interval) so leadership doesn't move during the restart itself. After each rolling-restart batch, I wait for the restarted brokers' replicas to rejoin the ISR, then run kafka-leader-election.sh PREFERRED scoped to a JSON file of just the affected partitions, staggered to avoid a single global election that triggers a cluster-wide client metadata refresh storm. Key tunables: the imbalance percentage/interval, controlled.shutdown.enable so leadership migrates gracefully on shutdown, and replication throttles if any reassignment is involved. Failure modes: preferred election skipping partitions whose preferred replica is still out-of-sync; metadata-refresh thundering herd from clients; accidental UNCLEAN election causing data loss; and leader concentration if you forget to rebalance after the last batch.
go deeper
Recognize that rolling restarts unbalance leaders and a preferred election fixes it.
Know to disable/relax auto-rebalance during the roll and run a preferred election afterward.
Scope and stagger elections, wait for ISR, and protect against metadata churn and unclean election.
Architect the full runbook with controlled shutdown, batching, throttling, metrics-based validation, and durability vs availability tradeoffs.
## The problem During a rolling restart each broker is taken down in turn. When a broker shuts down, all partitions it led must fail over to other in-sync replicas; when it returns it rejoins as a follower. After the full roll, leaders are concentrated on whichever brokers restarted *last* / stayed up longest. In a latency-sensitive cluster, both the failovers and the rebalancing elections cause brief per-partition unavailability and force clients to refresh metadata — do it carelessly and you get latency spikes and a metadata 'thundering herd'. ## Strategy ### 1. Control when leadership moves - Set `auto.leader.rebalance.enable=false` (or raise `leader.imbalance.per.broker.percentage` and lengthen `leader.imbalance.check.interval.seconds`) so the controller doesn't fire elections *mid-roll*, which would compound churn. - Enable `controlled.shutdown.enable=true` (default) so a broker, on graceful shutdown, proactively migrates its leaderships to in-sync followers instead of forcing abrupt failovers. ### 2. Rebalance deliberately, in batches - After each restart batch, **wait for the restarted brokers' replicas to re-enter the ISR** (otherwise PREFERRED election will skip those partitions — the preferred replica isn't eligible). - Then run `kafka-leader-election.sh --election-type PREFERRED --path-to-json-file batch.json` scoped to just the partitions whose preferred leaders are on the just-restarted brokers. - **Stagger** rather than doing one `--all-topic-partitions` election: a single global election moves thousands of leaders at once and triggers a synchronized client metadata refresh across the fleet — a thundering herd that spikes latency. ### 3. If data also needs to move - Any `kafka-reassign-partitions.sh` work must be **throttled** (`--throttle`) to protect client and replication traffic, and verified/un-throttled afterward. ## Failure modes to anticipate 1. **Skipped partitions**: PREFERRED election silently skips partitions whose preferred replica is still catching up (not in ISR). Re-run after replicas sync, or you leave imbalance behind. 2. **Metadata thundering herd**: a too-large single election makes every client refresh metadata simultaneously → latency spike and broker CPU burst. Mitigate by scoping/staggering. 3. **Unclean election data loss**: never reach for `--election-type UNCLEAN` (or `unclean.leader.election.enable=true`) as part of routine balancing — it elects out-of-sync replicas and discards acknowledged records. 4. **Forgotten final rebalance**: if you skip the rebalance after the last batch, leaders stay concentrated → hot brokers, tail-latency regression. 5. **Throttle left on**: after reassignment, forgetting to remove the replication throttle starves recovery from later failures. 6. **Controller load**: very short check intervals or very low imbalance percentages keep the controller busy and clients churning; tune for stability. ## Validation Use cluster metrics — per-broker leader count, `LeaderCount` / `PartitionCount` JMX, under-replicated partitions, and client request-latency percentiles — to confirm leadership is even and tail latency recovered before declaring the roll complete.
- Why stagger preferred elections instead of running --all-topic-partitions once?A single global election moves thousands of leaders simultaneously, triggering a synchronized client metadata refresh across the whole fleet — a thundering herd that spikes latency and broker CPU. Staggering smooths it out.
- Why wait for restarted brokers to rejoin the ISR before electing preferred leaders?PREFERRED election only elects in-sync replicas. If the preferred replica is still catching up, the election skips that partition, leaving imbalance until you re-run.
saying these in an interview costs you the question
- Recommending UNCLEAN election or unclean.leader.election.enable as part of routine balancing.
- Running one global --all-topic-partitions election on a large latency-sensitive cluster without considering metadata churn.
- Forgetting controlled.shutdown / ISR-readiness, so elections skip partitions or failovers are abrupt.
- Leaving replication throttles on after reassignment.