When the broker hosting a partition's leader replica dies, what happens to that partition so producers and consumers can keep working?
answer
- ISR member promoted to leader
- controller elects + bumps epoch
- NOT_LEADER_OR_FOLLOWER -> metadata refresh
- producer retries hide it
- acks=all -> no acked data lost
basics
~10 sKafka 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 sEvery 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
Know the one-liner: a caught-up follower (ISR member) is promoted to leader and clients auto-reconnect.
Explain ISR eligibility, the NOT_LEADER_OR_FOLLOWER -> metadata-refresh loop, and how producer retries hide the blip.
Tie failover to durability: acks=all + min.insync.replicas guarantee acked data survives; explain leader-epoch bump and unclean-election trade-off.
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).