“How does Kafka make sure a message is not lost?” is the durability question every Kafka interview reaches. Interviewers expect you to talk about replication factor, leaders and followers, the in-sync replica set, acks, min.insync.replicas and unclean leader election, and then they test the edges: “with acks=all, can you still lose data?” and “what happens when two of three brokers are down?” Those edges are exactly where real clusters have lost acknowledged orders.
This guide builds the replication model from the leader outwards and shows how each producer and topic setting changes what “acknowledged” means. It targets Kafka 3.8/4.0 in KRaft mode. Commands and code are illustrative: they follow the documented tools and configs but were not run against a live cluster.
Before you start
You should know that a topic is split into partitions stored on brokers, and that a producer sends records to a partition. It helps to understand the general idea of a leader-based replicated system, where one node accepts writes and others copy them. Basic Java is enough for the producer snippets.
The short answer
Each partition has a replication factor, usually 3: one leader replica takes all writes and the followers continuously fetch from it. The in-sync replica set (ISR) is the leader plus the followers that are caught up, and the high watermark is the highest offset copied to every ISR member; consumers only read below it. With acks=all the leader acknowledges a write only once every ISR member has it, and min.insync.replicas (typically 2 with RF 3) refuses writes when the ISR is too small, so every acknowledged record is on at least two brokers. Unclean leader election, off by default, would let an out-of-sync replica lead and lose data.
How it works
A produce request goes to the partition leader, which appends the batch to its log and advances its log end offset (LEO). Followers send fetch requests to the leader, exactly like consumers, starting from their own LEO. Each fetch tells the leader how far that follower has replicated, which is how the leader knows when a record is safely copied.
The leader computes the high watermark (HW) as the smallest LEO among ISR members. Records below the HW are committed: they exist on every in-sync replica, so any of them could become leader without losing those records. Consumers fetch only up to the HW. A follower that has not caught up to the leader’s end within replica.lag.time.max.ms (30 seconds) is removed from the ISR; when it catches up again it is added back. In KRaft mode the leader reports ISR changes to the controller with an AlterPartition request, and the controller records them in the metadata log.
When a leader fails, the controller picks a new leader from the ISR. Because every ISR member has everything below the HW, no committed record is lost. Records the old leader had appended above the HW, never fully replicated, may be gone; a replica that returns uses leader epochs to find where its log diverged and truncates the uncommitted tail.
The producer’s acks setting decides when it hears “success”:
acks |
Leader replies when | Can lose acknowledged data when |
|---|---|---|
0 |
never; producer does not wait | always possible, silently |
1 |
leader has appended | leader dies before followers copy |
all (default since 3.0) |
every current ISR member has the record | ISR has shrunk to the leader and it then dies |
The last row is why min.insync.replicas exists. It is a topic or broker setting, and the broker default is 1. With acks=all and min.insync.replicas=2, a leader rejects writes with NotEnoughReplicasException when the ISR has fewer than two members.
Step-by-step walkthrough
Step 1: Create topics with replication that survives a failure
Replication factor 3 is the common baseline: it survives one broker failure while still keeping two copies. Set the minimum ISR on the topic, and spread replicas across racks or availability zones by setting broker.rack on every broker so one zone failure cannot take all copies.
kafka-topics.sh --bootstrap-server localhost:9092 --create --topic payments \
--partitions 12 --replication-factor 3 --config min.insync.replicas=2Also check broker defaults, because auto-created topics use default.replication.factor, which is 1 unless you change it. Many teams disable auto.create.topics.enable for this reason.
Step 2: Configure the producer for durable writes
The current Java defaults are already durable, but stating them protects you from a later change that quietly weakens them, such as someone setting acks=1 for speed, which also disables idempotence.
Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ProducerConfig.ACKS_CONFIG, "all");
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
props.put(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG, 120_000); // total time to keep retryingStep 3: Treat the callback as part of the durability contract
acks=all only helps if you notice failures. When the ISR is below the minimum, the producer retries the retriable error until delivery.timeout.ms expires and then fails the send. If nobody checks the callback, that record is lost even though every setting looked durable.
producer.send(record, (metadata, exception) -> {
if (exception != null) {
// e.g. TimeoutException after NotEnoughReplicasException retries
failedSends.increment();
outbox.markForRetry(record.key());
}
});Step 4: Decide what happens when every in-sync replica is gone
If the leader and all ISR members are unavailable, the partition goes offline until one of them returns. unclean.leader.election.enable=true would let an out-of-sync follower lead instead, restoring availability but losing every record it never copied. Keep it false (the default) for business data. In a real outage you can trigger a one-off unclean election for a specific partition, which is a deliberate, recorded decision rather than a standing risk:
kafka-leader-election.sh --bootstrap-server localhost:9092 \
--election-type unclean --topic metrics --partition 3Step 5: Monitor the ISR, not just broker uptime
Brokers can all be “up” while partitions are one failure away from losing data. Alert on the broker metrics UnderReplicatedPartitions, UnderMinIsrPartitionCount and IsrShrinksPerSec (under kafka.server:type=ReplicaManager), and on OfflinePartitionsCount from the active controller.
Worked scenario
An orders topic had replication factor 3 and producers used acks=all, so the team believed every acknowledged order was on three brokers. Nobody set min.insync.replicas, so it was the broker default of 1.
During a network problem between availability zones, both followers fell behind for more than 30 seconds and were removed from the ISR. The leader was now the only in-sync replica, and acks=all meant “acknowledged by the leader alone.” For twelve minutes, orders were accepted with one copy. Then the leader’s disk failed. With unclean election disabled, the partition went offline because no in-sync replica was left. To restore service, an engineer forced an unclean election, a follower became leader from where it had stopped, and twelve minutes of acknowledged orders were gone.
The fix was small and immediate:
kafka-configs.sh --bootstrap-server localhost:9092 --alter --entity-type topics \
--entity-name orders --add-config min.insync.replicas=2With that setting, the same network problem would have made producers fail with NotEnoughReplicasException after the ISR shrank, which is loud and recoverable through retries or an outbox, instead of accepting writes that could not survive one more failure. The team also set broker.rack per zone and added alerts on ISR shrinks.
Common mistake
- “
acks=allmeans all replicas.” It means all replicas currently in the ISR, which can be just the leader. - “
min.insync.replicasapplies to every producer.” It is enforced only foracks=allwrites;acks=1producers still succeed. - “Set
min.insync.replicasequal to the replication factor for maximum safety.” Then one broker restart blocks all writes. RF 3 with a minimum of 2 tolerates one failure. - “Consumers can read anything the leader wrote.” Only below the high watermark.
- “Replication factor 3 means you can lose two brokers and keep writing.” With a minimum ISR of 2, losing two makes the partition read-only for
acks=allproducers. That is intended.
Verify the behavior
On a local three-broker KRaft cluster (for example with Docker Compose), create the topic from Step 1, then check its replicas:
kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic payments
# Partition: 0 Leader: 1 Replicas: 1,2,3 Isr: 1,2,3 (one line per partition)Stop two brokers and produce with acks=all:
kafka-console-producer.sh --bootstrap-server localhost:9092 --topic payments \
--producer-property acks=all
# warnings that the produce request failed with NOT_ENOUGH_REPLICAS, retried until the delivery timeout
kafka-topics.sh --bootstrap-server localhost:9092 --describe --under-min-isr-partitionsRepeat with --producer-property acks=1 and the write succeeds, which demonstrates that the minimum ISR applies only to acks=all. Restart the brokers and --describe shows them rejoining the Isr column once they catch up.
Follow-up questions
Does one slow follower slow down acks=all? Yes, until it is removed from the ISR after replica.lag.time.max.ms. Then writes no longer wait for it.
What is NotEnoughReplicasAfterAppendException? The ISR shrank after the leader appended the record but before replication finished. The record is in the leader’s log but was not acknowledged, so the producer retries and idempotence prevents a duplicate.
What is the preferred leader? The first replica in the assignment. With auto.leader.rebalance.enable (on by default), leadership moves back to it after a failover so load stays balanced.
Can consumers read from followers? Yes, since Kafka 2.4: set replica.selector.class to the rack-aware selector on brokers and client.rack on consumers to fetch from the nearest replica and cut cross-zone traffic.
Interview exercise
A topic has replication factor 3 and min.insync.replicas=2. Producer A uses acks=all, producer B uses acks=1. Broker 2 crashes, then broker 3 crashes, then broker 1 (the leader) loses its disk permanently. For each stage, what happens to A, B and consumers, and what data can be lost?
Answer and reasoning
After broker 2 crashes, the ISR is brokers 1 and 3. Both producers succeed, and A’s records are on two brokers. After broker 3 crashes, the ISR is just broker 1: A’s writes fail with NotEnoughReplicasException and are retried, while B’s writes succeed on the leader alone. Consumers still read up to the high watermark, which no longer waits for the crashed brokers. When broker 1’s disk is lost, no in-sync replica remains, so the partition is offline. A lost nothing it was told succeeded, because nothing was acknowledged after the ISR fell below two. B’s records accepted after broker 3 crashed existed only on broker 1 and are gone. Recovering availability requires an unclean election from broker 2 or 3, which confirms that loss.
Continue learning
Practise with the Apache Kafka interview questions and the Kafka MCQs. Related notes: Kafka idempotent producer and batching and replication in system design. Primary sources: the Apache Kafka documentation, Confluent’s replication and committed messages notes and the producer configuration reference.