skip to content

Leadership Failover and Replica State

What a partition does when its leader dies: failover to an in-sync replica, an epoch bump, follower truncation, and transparent client re-discovery. Interviewers walk through this as the 'a broker just died' scenario.

part ofApache Kafkaoverview, primer and where to startread it →
on this pageshow

questions

5

When the broker hosting a partition's leader replica dies, what happens to that partition so producers and consumers can keep working?

level: juniorimportance: must knowfreq 78%

answer

  1. ISR member promoted to leader
  2. controller elects + bumps epoch
  3. NOT_LEADER_OR_FOLLOWER -> metadata refresh
  4. producer retries hide it
  5. acks=all -> no acked data lost

basics

~10 s

Kafka promotes one of the in-sync follower replicas to be the new leader. Clients get an error, refresh metadata, find the new leader's broker, and resume producing/consuming there. No manual action is needed.

solid answer

~40 s

Every Kafka partition has one leader replica (handles all reads/writes) and several followers that replicate it. The set of followers caught up with the leader is the ISR (in-sync replica set). When the leader's broker dies, the cluster controller detects it and elects a new leader from the surviving ISR members, then updates cluster metadata. Clients that try the old leader receive a NOT_LEADER_OR_FOLLOWER error, which triggers a metadata refresh; they learn the new leader and reconnect transparently. The application sees at most a brief pause and some retried requests. As long as at least one in-sync replica survives, no acknowledged data is lost. This whole flow is automatic — it's the core of Kafka's high availability.

go deeper

for a junior

Know the one-liner: a caught-up follower (ISR member) is promoted to leader and clients auto-reconnect.

for a middle

Explain ISR eligibility, the NOT_LEADER_OR_FOLLOWER -> metadata-refresh loop, and how producer retries hide the blip.

for a senior

Tie failover to durability: acks=all + min.insync.replicas guarantee acked data survives; explain leader-epoch bump and unclean-election trade-off.

for a principal

Reason about availability vs durability trade-offs cluster-wide, RTO during mass failover, and when unclean election is an acceptable policy.

## The setup A Kafka **topic** is split into **partitions**. Each partition is replicated onto several brokers for fault tolerance. Among those copies (replicas), exactly one is the **leader**: it is the only replica that accepts produce (write) and consume (read) requests. The others are **followers** — they continuously fetch records from the leader to stay in sync. The **ISR** (in-sync replica set) is the subset of replicas that are sufficiently caught up with the leader (within `replica.lag.time.max.ms`, default 30s). Only ISR members are eligible to become leader under the default safe setting. ## What happens when the leader's broker dies 1. **Detection.** The broker stops sending heartbeats. In KRaft mode the active controller notices the broker's session expiring and marks it fenced/dead. (The controller/quorum machinery itself is out of scope here — a sibling topic owns it.) 2. **Leader election.** For every partition that had its leader on the dead broker, the controller picks a new leader from the surviving ISR. It also bumps the partition's **leader epoch** (a monotonically increasing version number for leadership) so everyone can tell old leadership from new. 3. **Metadata propagation.** The controller writes the new leader/ISR/epoch into cluster metadata and propagates it to all brokers. 4. **Client rediscovery.** A producer or consumer still pointing at the old leader sends a request and gets back an error such as `NOT_LEADER_OR_FOLLOWER` (or a connection failure). The client library treats this as a signal to refresh its metadata (it asks any broker "who leads this partition now?"), learns the new leader, and retransmits to the correct broker. The producer's internal retries (`retries`, default effectively `Integer.MAX_VALUE`, bounded by `delivery.timeout.ms`) make this invisible to application code. ## Durability guarantee With `acks=all` and `min.insync.replicas >= 2`, the leader only acknowledges a write once it's replicated to the required number of in-sync replicas. So any record the producer saw acknowledged already exists on a follower that can become the new leader — no acknowledged data is lost. ## Edge cases - If the **entire ISR is gone**, the partition goes offline; it can only recover by waiting for an ISR member to return, or by enabling **unclean leader election** (`unclean.leader.election.enable=true`), which promotes an out-of-sync replica and risks data loss. - Consumers continue from their committed offsets against the new leader; because acknowledged data is preserved, offsets stay valid. ## Net effect From the application's perspective: a short blip, some retried requests, and then normal operation — all automatic.

  • How does a producer find out it's talking to the wrong (old) leader?
    Its request to the old leader fails — typically a NOT_LEADER_OR_FOLLOWER error or a connection error. The client treats that as 'my metadata is stale', issues a metadata-refresh request to learn the current leader, and retries against the new broker.
  • Can data acknowledged to the producer be lost during failover?
    Not if you use acks=all with min.insync.replicas>=2. Acked records already live on an in-sync follower eligible to be the new leader. Loss only happens with weak acks, or with unclean leader election promoting an out-of-sync replica.

saying these in an interview costs you the question

  • Saying clients need manual reconfiguration or a restart to find the new leader.
  • Claiming failover always loses data (it doesn't, with acks=all + min.insync.replicas).
  • Saying any follower can become leader (only ISR members, unless unclean election is enabled).
  • Confusing the leader (the writable replica) with the controller (the cluster coordinator).

context

open as a page

Walk through how a producer transparently re-discovers a new partition leader after a failover, including the relevant errors and configs.

level: middleimportance: should knowfreq 55%

basics

~20 s

The producer's send to the old leader fails with a retriable error like NOT_LEADER_OR_FOLLOWER. The client marks its metadata stale, asks any broker for fresh metadata, learns the new leader, and retries automatically within delivery.timeout.ms.

open as a page

What is the leader epoch, and why is it bumped on every leadership change?

level: seniorimportance: should knowfreq 52%

basics

~20 s

The leader epoch is a number that increases by one each time a partition gets a new leader. It lets brokers and clients tell stale leadership apart from current, so followers truncate correctly and old leaders can't corrupt the log.

open as a page

During failover, what is unclean leader election, what does it trade off, and when (if ever) would you enable it?

level: seniorimportance: should knowfreq 48%

basics

~20 s

Unclean leader election lets an out-of-sync replica become leader when no in-sync replica survives. It restores availability but can lose acknowledged data. It's off by default; enable it only when uptime matters more than durability.

open as a page

After a failover, how does a former-leader replica that rejoins reconcile its log with the new leader, and what role does the high watermark play?

level: principalimportance: nice to knowfreq 30%

basics

~20 s

When the old leader comes back as a follower, it may hold uncommitted records past the high watermark that the new leader never had. Using leader-epoch lookups it finds the divergence point, truncates the extra records, then re-fetches from the new leader to catch up and rejoin the ISR.

open as a page