Walk through how a producer transparently re-discovers a new partition leader after a failover, including the relevant errors and configs.
answer
- cached partition->leader map goes stale
- NOT_LEADER_OR_FOLLOWER / LEADER_NOT_AVAILABLE
- Metadata request refreshes the map
- delivery.timeout.ms bounds total retry
- idempotence dedupes resends
basics
~20 sThe 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.
solid answer
~40 sA producer caches a metadata map of partition -> leader broker. After failover that map is briefly stale. When it sends to the old leader it gets back a retriable error — typically NOT_LEADER_OR_FOLLOWER (or LEADER_NOT_AVAILABLE / a connection failure). The client treats these as a signal to refresh: it issues a Metadata request to any broker, which returns the current leader and leader epoch. The producer reroutes the batch to the new leader and retries. This is bounded by delivery.timeout.ms (default 120s) and retry.backoff.ms; retries default to effectively unlimited within that window. Metadata is also proactively refreshed every metadata.max.age.ms (default 5m). With enable.idempotence=true (the default), retries don't create duplicates. So application code calling send() never sees the failover unless the whole window expires.
go deeper
Know that the client auto-refreshes metadata and retries; the app doesn't reconnect manually.
Name the error codes and the error->refresh->retry loop, plus delivery.timeout.ms.
Tie in idempotence for dedupe, leader-epoch fencing, and proactive vs reactive refresh.
Reason about tuning the retry budget vs latency SLOs and failure amplification across many partitions during a broker loss.
## The producer's view of the cluster A Kafka producer keeps a **metadata cache**: for each topic-partition it tracks which broker is the current **leader** (plus the leader epoch and the broker's host:port). All produce requests for a partition are sent directly to that leader broker — the producer does the routing itself. ## What goes stale on failover When the leader's broker dies and a new leader is elected, the producer's cached map still points at the dead/old broker for a short time. The next send exposes the staleness. ## The error -> refresh -> retry loop 1. **Send fails.** Sending to the old leader yields one of: - a **connection error** (broker is down), or - a **retriable Kafka error code** returned by a broker that no longer leads that partition: `NOT_LEADER_OR_FOLLOWER` (the modern code; older name `NotLeaderForPartition`), `LEADER_NOT_AVAILABLE` (election still in progress), or `FENCED_LEADER_EPOCH` (the client's epoch is stale). 2. **Mark metadata stale.** The producer flags that topic's metadata for refresh. 3. **Refresh.** It sends a **Metadata** request to any reachable broker (bootstrap or known). The response lists the current leader, replicas, ISR, and leader epoch for each partition. 4. **Reroute and retry.** The batched records are re-queued to the new leader and re-sent. Retries are governed by: - `retries` (default `Integer.MAX_VALUE`), - `retry.backoff.ms` (default 100ms) between attempts, - and the hard ceiling **`delivery.timeout.ms`** (default 120000ms) — total time from `send()` to success/failure. If failover finishes within this window, the application never sees an error. 5. **Proactive refresh.** Independently, the producer refreshes metadata every `metadata.max.age.ms` (default 300000ms), so it self-heals even without an error. ## Why no duplicates Because retries resend batches, naive retrying could duplicate records. Kafka's **idempotent producer** (`enable.idempotence=true`, the default since 3.0) stamps each batch with a producer ID + sequence number so the broker dedupes resends. So transparent re-discovery is also exactly-once-per-partition safe within a session. ## Consumer side (parallel mechanism) Consumers do the same: a fetch to the old leader returns NOT_LEADER_OR_FOLLOWER, the consumer refreshes metadata, and re-fetches from the new leader at its committed offset. Because acknowledged data survived (with acks=all), no records are skipped. ## What the app should NOT do It should not catch these retriable errors and reconnect manually, nor hardcode broker addresses per partition — the client library owns rediscovery. The app only needs sane `delivery.timeout.ms` and acks settings.
- Which config bounds how long the producer keeps retrying during a failover before giving up?delivery.timeout.ms (default 120000ms). It's the total wall-clock budget from send() to terminal success or failure, covering all retries and backoffs. retries and retry.backoff.ms operate inside that ceiling.
- Why don't producer retries during failover create duplicate records?Because enable.idempotence=true (default) tags batches with a producer ID and per-partition sequence numbers; the broker rejects/deduplicates a resent batch it has already persisted, giving exactly-once delivery per partition within the producer session.
saying these in an interview costs you the question
- Saying the application must catch the error and manually reconnect or restart the producer.
- Confusing metadata.max.age.ms (proactive refresh) with delivery.timeout.ms (retry ceiling).
- Claiming retries inherently cause duplicates (idempotence prevents that).
- Saying the producer always reconnects to a fixed broker for a partition — leadership and thus the target broker changes.