03 / 05

What happens when a partition leader fails?

Difficulty: 5/10
ISR and leader election

The controller detects the failure, elects a new leader from the ISR, bumps the leader epoch, publishes new metadata, and clients refresh and retry

The sequence is: detection, election, propagation, client recovery. Detection: in KRaft the active controller notices the broker has stopped heartbeating (broker.session.timeout.ms) or the broker shuts down in a controlled way, in which case it hands off leadership first. Election: for each partition that lost its leader, the controller picks a new leader from the ISR, normally the first live replica in the replica-list order that is in the ISR, increments the leader epoch, and records the change in the metadata log. Propagation: brokers learn the new leadership by reading metadata updates. Recovery: clients that still send requests to the old leader get NOT_LEADER_OR_FOLLOWER errors (or time out), refresh their metadata, find the new leader and retry. Producers with retries and idempotence resend batches safely; consumers resume from their last committed offset.

Why it is safe: every ISR member holds all records up to the high watermark, so no acknowledged record (with acks=all) is lost. The leader epoch lets a rejoining former leader detect that its log diverged and truncate to the right point instead of keeping records the new leader never had. If no ISR member is alive, the partition goes offline: with unclean.leader.election.enable=false (the default) it stays unavailable until an ISR member returns, rather than electing a stale replica and losing acknowledged data. Afterwards, leadership can be imbalanced because the old broker's replicas return as followers; preferred leader election (automatic via auto.leader.rebalance.enable, or manual) moves leaders back to the first replica in the list.

javascript
  1. 1

    Trade-off: planned shutdown (controlled shutdown) moves leaders before the broker stops, with almost no client-visible gap. A crash relies on timeout-based detection, so failover time includes the detection delay plus client metadata refresh.

  2. 2

    Trade-off: unclean election restores availability at the cost of acknowledged data. Keep it off for critical topics; consider it only for data you can regenerate.

  3. 3

    Common mistake: assuming failover is instant and invisible. Clients see errors or latency until they refresh metadata, so retries and sensible delivery.timeout.ms are essential.

  4. 4

    Common mistake: assuming acks=1 writes survive a leader crash. The record may exist only on the dead leader, and the new leader never had it.

  5. 5

    Common mistake: forgetting preferred leader rebalancing, leaving the cluster skewed after restarts with some brokers carrying most leaders.

  6. 6

    Eligible leader replicas: KIP-966 adds the concept of eligible leader replicas (ELR) so that replicas which were in the ISR when it shrank below min.insync.replicas can still be considered for safe election. It was introduced as early access in Kafka 4.0, so check its maturity in your release before relying on it.

  7. 7

    Version note: in ZooKeeper-based clusters the controller detected failure through an expired ZooKeeper session (zookeeper.session.timeout.ms), and failover time depended on that and on controller load. KRaft improves metadata propagation and recovery, but exact timings still depend on your settings.

Share

Share via WhatsApp, X, Facebook, LinkedIn or copy link. Open Graph preview enabled.