skip to content

Replication and Durability

How partitions survive broker loss: leader/follower fetching, the ISR, high watermark and leader epoch, and election choices. Interviewers push here to see whether you can reason about the durability-versus-availability dial.

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

questions

page 1 of 2

What are the three main Kafka settings that work together to control write durability, and what does each do at a high level?

level: juniorimportance: must knowfreq 78%

answer

  1. RF = how many copies
  2. min.insync.replicas = how many must be caught up
  3. acks = how many producer waits for
  4. RF=3, MISR=2, acks=all baseline
  5. MISR only bites with acks=all

basics

~20 s

Replication factor (how many copies of each partition), min.insync.replicas (how many copies must be caught up to accept a write), and acks (how many copies the producer waits for: 0, 1, or all). Together they decide how many failures you can survive without losing data.

solid answer

~40 s

Three knobs interact. Replication factor (RF), a topic property, sets how many brokers hold a copy of each partition; RF=3 means a leader plus two followers. min.insync.replicas (a topic/broker config) sets the minimum number of in-sync replicas (ISR) that must be present for an acks=all write to be accepted; if fewer are in sync, the broker rejects the write with NotEnoughReplicas. acks (a producer config) controls how many replicas the producer waits to acknowledge: acks=0 (fire and forget), acks=1 (leader only), acks=all (every member of the ISR). The classic durable baseline is RF=3, min.insync.replicas=2, acks=all: it tolerates one broker loss and still guarantees no acknowledged write is lost.

go deeper

for a junior

Memorize the three knobs and the RF=3/MISR=2/acks=all baseline and what each word means.

for a middle

Explain ISR dynamics and that min.insync.replicas only applies to acks=all.

for a senior

Articulate the exact guarantee (acknowledged write survives N broker losses) and the durability-vs-availability tension when MISR approaches RF.

for a principal

Reason about org-wide defaults, cost of RF, and how these defaults compose with idempotence/transactions for end-to-end guarantees.

## What governs durability Kafka stores each topic partition as an ordered log replicated across brokers for fault tolerance. Three settings govern the durability of a write: - **Replication factor (RF)** — a per-topic setting (`--replication-factor` at creation, or the broker default `default.replication.factor`). RF=N means each partition has N copies on N different brokers: one **leader** that handles all reads and writes, and N-1 **followers** that continuously fetch from the leader to stay current. RF determines the maximum number of broker failures the data can physically survive (you can lose up to N-1 copies and still have one). - **In-Sync Replicas (ISR)** — the subset of replicas (including the leader) that are fully caught up to the leader's log. A follower drops out of the ISR if it falls behind by more than `replica.lag.time.max.ms` (default 30s). The ISR shrinks and grows dynamically. - **min.insync.replicas** — a topic/broker config naming the minimum ISR size required to accept a write *when the producer uses acks=all*. If the current ISR is smaller than this, the leader rejects the produce request with `NotEnoughReplicasException` / `NotEnoughReplicasAfterAppendException`. This protects you from the dangerous case where the ISR has collapsed to just the leader: without it, an acks=all write to a lone leader could be lost if that leader then dies. - **acks** — a producer config: `acks=0` means the producer never waits (highest throughput, can silently lose data); `acks=1` waits only for the leader to write to its log (lost if the leader dies before a follower replicates); `acks=all` (a.k.a. `acks=-1`) waits for every member of the current ISR to acknowledge. ## How they combine Durability comes from `acks=all` AND `min.insync.replicas >= 2`. With acks=all, the producer is told a write succeeded only after all ISR members have it; with min.insync.replicas=2, the write is only accepted while at least two replicas are in sync, so an acknowledged record always lives on at least two brokers. The canonical durable config is **RF=3, min.insync.replicas=2, acks=all**: it survives one broker failure with zero acknowledged-data loss and stays available for writes. Setting min.insync.replicas equal to RF (e.g. 3) maximizes durability but means *any* single replica falling out of the ISR halts writes — a **durability/availability trade-off**. ## Edge cases - `min.insync.replicas` only matters with acks=all; with acks=1 it is ignored. - RF without acks=all gives you read availability but not write durability guarantees.

  • Does min.insync.replicas have any effect when the producer uses acks=1?
    No. min.insync.replicas is only enforced for acks=all (acks=-1) writes. With acks=1 the leader acknowledges on its own write regardless of ISR size, so the setting is ignored.
  • Why is RF=3 with min.insync.replicas=2 preferred over RF=2 with min.insync.replicas=2?
    With RF=2/MISR=2 you have no headroom: losing one broker drops the ISR to 1, below min.insync.replicas, so all writes stop immediately. RF=3/MISR=2 tolerates one failure and stays writable.

saying these in an interview costs you the question

  • Saying min.insync.replicas guarantees durability even with acks=1 (it does nothing there).
  • Confusing RF with min.insync.replicas — RF is total copies, MISR is the required caught-up subset.
  • Claiming acks=all waits for ALL replicas (it waits for all members of the current ISR, which may be smaller than RF).

context

open as a page

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%

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.

open as a page

After a Kafka write returns successfully, where does the data physically live, and why does that distinction matter?

level: juniorimportance: must knowfreq 60%

basics

~20 s

After a successful write the data is in the broker's RAM (the OS page cache), and with acks=all also in the RAM of other replicas. It isn't necessarily on the physical disk yet — the operating system writes it to disk a bit later.

open as a page

What is the difference between the log-end-offset (LEO) and the high watermark (HW) on a Kafka partition leader, and why does it matter to consumers?

level: juniorimportance: must knowfreq 70%

basics

~20 s

LEO is the offset just past the last message written to the log. HW is the highest offset that has been replicated to all in-sync replicas. Consumers can only read up to (below) the HW, so they never see un-replicated records.

open as a page

What is the In-Sync Replicas (ISR) set in Apache Kafka, and how does it relate to the assigned replicas of a partition?

level: juniorimportance: must knowfreq 78%

basics

~20 s

The ISR is the subset of a partition's replicas (leader plus followers) that are fully caught up with the leader's log. Assigned replicas (AR) are all replicas; ISR is the healthy, up-to-date ones eligible to serve and be elected leader.

open as a page

In Kafka's replication model, which replica handles producer writes and reads, and what role do the other replicas play?

level: juniorimportance: must knowfreq 75%

basics

~10 s

Each partition has one leader replica that handles all produce and consume requests. The other replicas are followers; they only copy the leader's data and stand by to take over if the leader fails.

open as a page

What does the broker-side config min.insync.replicas do, and how does it relate to acks=all?

level: juniorimportance: must knowfreq 70%

basics

~10 s

min.insync.replicas sets the minimum number of in-sync replicas that must acknowledge a write before the broker accepts it. It only takes effect when the producer uses acks=all; otherwise the broker ignores it.

open as a page

What is the 'preferred leader' for a Kafka partition, and why does it matter?

level: juniorimportance: must knowfreq 70%

basics

~20 s

The preferred leader is the first broker listed in a partition's replica assignment (the assigned-replica list). Kafka tries to make it the leader so leadership is spread evenly across brokers and no single broker is overloaded.

open as a page

What is the broker.rack configuration in Kafka, and what does setting it actually do?

level: juniorimportance: must knowfreq 60%

basics

~20 s

broker.rack is a per-broker string tagging which rack or availability zone a broker is in. Kafka uses these tags to spread a partition's replicas across different racks, so one rack failing doesn't take down all copies.

open as a page

What is unclean leader election in Kafka, and what does the unclean.leader.election.enable setting control?

level: juniorimportance: must knowfreq 70%

basics

~20 s

Unclean leader election lets Kafka promote a replica that is NOT fully up to date (out of the in-sync set) to leader when no in-sync replica is available. The unclean.leader.election.enable flag turns this on or off. It trades possible data loss for availability.

open as a page

Explain why setting min.insync.replicas equal to the replication factor hurts availability, and how to pick the value to hit a specific durability/availability target.

level: middleimportance: must knowfreq 70%

basics

~20 s

If min.insync.replicas equals RF, every replica must be in sync to accept writes, so losing or lagging just one broker stops all writes. Setting it to RF-1 (e.g. 2 with RF=3) keeps you durable while tolerating one failure.

open as a page

How do acks=all and min.insync.replicas work together to provide durability, and how would you configure them with a replication factor of 3?

level: middleimportance: must knowfreq 65%

basics

~20 s

acks=all makes the producer wait until all in-sync replicas have the record. min.insync.replicas sets how many replicas must be in sync for that write to be allowed. With replication factor 3, set min.insync.replicas=2 and acks=all so you can lose one broker without losing data and without halting writes.

open as a page

Why does Kafka rely on replication for durability instead of fsync-ing every message to disk before acknowledging it?

level: middleimportance: must knowfreq 70%

basics

~20 s

Kafka acknowledges writes once enough replicas have the data in memory, not once it's flushed to disk. Replication across machines is faster and safer than a per-message fsync, which would be slow and still vulnerable to a single disk failure.

open as a page

How does replica.lag.time.max.ms control ISR shrink and expansion, and why did Kafka switch from a message-count-based lag check?

level: middleimportance: must knowfreq 70%

basics

~20 s

A follower stays in the ISR if it fetched up to the leader's latest offset within replica.lag.time.max.ms (default 30s). Miss that window and the leader removes it; catch back up and it rejoins. The old count-based check (lag in messages) misjudged bursty traffic, so it was replaced by time.

open as a page

What is the ISR, and how does Kafka decide whether a follower is in-sync or should be removed from it?

level: middleimportance: must knowfreq 58%

basics

~20 s

The ISR (in-sync replica set) is the leader plus all followers caught up with it. A follower is dropped from the ISR if it hasn't fetched up to the leader's log end within replica.lag.time.max.ms, and re-added once it catches back up.

open as a page

How does a follower replica actually stay in sync with the leader? Describe the ReplicaFetcherThread and its fetch loop.

level: middleimportance: must knowfreq 60%

basics

~20 s

Each follower broker runs ReplicaFetcherThreads. A thread sends Fetch requests to the leader asking for records starting at the follower's current log-end offset, appends what it gets to its local log, and advances its offset — repeating continuously.

open as a page

Walk through what happens when the ISR shrinks below min.insync.replicas while producers are writing with acks=all.

level: middleimportance: must knowfreq 60%

basics

~10 s

The leader rejects new acks=all writes with NotEnoughReplicasException, so producers fail and retry. The partition becomes read-still-available but write-blocked until enough followers rejoin the ISR.

open as a page

How does auto.leader.rebalance.enable work, and what do leader.imbalance.check.interval.seconds and leader.imbalance.per.broker.percentage control?

level: middleimportance: must knowfreq 60%

basics

~20 s

auto.leader.rebalance.enable (default true) lets the controller automatically run preferred leader election. The check interval sets how often it inspects imbalance, and the per-broker percentage is the imbalance threshold that triggers a rebalance for a broker.

open as a page

What does unclean.leader.election.enable do, and how does it change the durability/availability trade-off when combined with the other replication settings?

level: seniorimportance: must knowfreq 62%

basics

~20 s

It decides whether a partition can elect a leader from a replica that was NOT in the in-sync set when all in-sync replicas are down. Enabled = stay available but possibly lose data; disabled (default) = stay consistent but the partition goes offline until an in-sync replica returns.

open as a page

Before KIP-101, Kafka followers truncated their logs to the high watermark on becoming a follower of a new leader. Describe the data-loss / log-divergence scenario this caused.

level: seniorimportance: must knowfreq 50%

basics

~20 s

On rejoining, a follower truncated everything above its high watermark, assuming records above the HW were uncommitted. But because the follower's HW lags, this could throw away records that were actually committed, or let two replicas keep different records at the same offset — causing data loss or log divergence.

open as a page

Explain how leader epochs and the OffsetsForLeaderEpoch request (KIP-101, completed by KIP-279) let a follower truncate to the correct divergence point instead of to the high watermark.

level: seniorimportance: must knowfreq 45%

basics

~20 s

Each leader is assigned a monotonically increasing leader epoch, and the log records which epoch produced each offset range (the leader-epoch cache). On becoming a follower, the broker asks the leader, via OffsetsForLeaderEpoch, for the end offset of its last known epoch, and truncates exactly there — the true point where the logs diverge — instead of guessing with the HW.

open as a page

Why does a record produced to the leader not become visible to consumers immediately, and how do followers and the high watermark control commit/visibility?

level: seniorimportance: must knowfreq 55%

basics

~20 s

A record appended on the leader is only 'committed' once all in-sync followers have replicated it. The high watermark marks the highest committed offset; consumers can only read up to it, so un-replicated tail records stay invisible.

open as a page

Why is RF=3 with min.insync.replicas=2 the standard durable configuration, rather than RF=3/min.isr=3 or RF=2/min.isr=2?

level: seniorimportance: must knowfreq 65%

basics

~20 s

RF=3/min.isr=2 lets you survive one broker failure while still accepting writes and keeping two durable copies. RF=3/min.isr=3 blocks writes on any single failure; RF=2/min.isr=2 blocks writes on any single failure and has no spare copy.

open as a page

Explain fetch-from-follower (KIP-392): how does rack-aware consumer fetching work, what configs enable it, and why would you use it?

level: seniorimportance: must knowfreq 50%

basics

~20 s

KIP-392 lets a consumer read from a follower replica in its own rack/AZ instead of always from the leader. You set replica.selector.class on the broker and client.rack on the consumer. It cuts cross-AZ network cost and latency.

open as a page

Walk through what happens to a Kafka topic (RF=3 across 3 AZs, min.insync.replicas=2, acks=all) when one entire AZ fails. What survives, and what would break this guarantee?

level: seniorimportance: must knowfreq 48%

basics

~20 s

Each partition keeps 2 of its 3 replicas because they're in different AZs. Leaders that were in the dead AZ fail over to surviving replicas, ISR drops to 2 which still meets min.insync.replicas=2, so acks=all producers and consumers keep working with no data loss.

open as a page

Walk through exactly what data loss and log divergence happens when an unclean leader election promotes an out-of-sync replica.

level: seniorimportance: must knowfreq 60%

basics

~20 s

The lagging replica becomes leader at its own (shorter) log end. Records the old leader had committed but this replica never fetched are gone. When the old replica returns, it must truncate its longer log down to the new leader's offset to re-join, throwing away those records.

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

Walk through exactly how the high watermark propagates from the leader to the followers, and explain why a follower's HW always lags the leader's HW by at least one fetch round-trip.

level: middleimportance: should knowfreq 45%

basics

~20 s

Followers fetch from the leader. Each Fetch response carries the leader's current HW. The follower applies it after appending the fetched data. Because the HW arrives in the next response, a follower's HW always trails the leader's by one fetch round-trip.

open as a page

Where is min.insync.replicas configured, and how do you set up the full durability contract for a topic in practice?

level: middleimportance: should knowfreq 45%

basics

~10 s

min.insync.replicas is a broker default that can be overridden per topic. The full contract needs three parts: replication factor at topic creation, min.insync.replicas as a topic config, and acks=all on the producer.

open as a page

How do you manually trigger preferred leader election with kafka-leader-election.sh, and how does PREFERRED differ from UNCLEAN election?

level: middleimportance: should knowfreq 45%

basics

~10 s

Run kafka-leader-election.sh with --election-type PREFERRED to move leadership back to preferred replicas. PREFERRED only elects in-sync replicas (safe). UNCLEAN can elect an out-of-sync replica to restore availability, risking data loss.

open as a page

showing 1–30 of 50