Why does Kafka rely on replication for durability instead of fsync-ing every message to disk before acknowledging it?
answer
- ack on replication, not on fsync
- acks=all → ISR has it in page cache
- fsync per message = throughput death
- one flushed copy still dies with the disk
- RF=3, min.insync.replicas=2
basics
~20 sKafka 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.
solid answer
~40 sKafka treats durability as a replication problem, not a disk-flush problem. A producer with acks=all is acknowledged when all in-sync replicas (ISR) have copied the record into their OS page cache and the leader has confirmed it — Kafka does NOT force an fsync per message. The data physically reaches disk later, asynchronously, when the OS flushes dirty pages or the optional flush.messages/flush.ms thresholds fire. The reasoning: an fsync on every message would cripple throughput (each fsync is a synchronous disk round-trip), and even a flushed single copy can be lost if that one broker's disk dies. Replicating to N brokers on independent hardware gives stronger guarantees than flushing one copy, because a record survives as long as at least one ISR member survives. min.insync.replicas controls how many copies must exist before acks=all succeeds.
go deeper
Know that Kafka confirms a write once enough copies exist on other brokers, not once it's written to disk.
Explain acks=all + ISR + page cache, and why per-message fsync would kill throughput while a single flushed copy is still fragile.
Tie together acks=all, min.insync.replicas, replication factor, and the asynchronous OS flush; reason about the residual correlated-power-loss window.
Frame durability as a replication design decision vs. disk-flush, justify the throughput/latency trade-off, and set fleet-wide defaults across racks/AZs.
## The core idea **Durability** means a write, once acknowledged, won't be lost. There are two ways to achieve it: (1) **persist to disk** — force the bytes onto stable storage with an `fsync` system call; or (2) **replicate** — copy the data to multiple independent machines. Kafka deliberately chooses (2) as its primary mechanism. ## What actually happens on a produce When a producer sends a record, the leader broker appends it to the active log segment. That `write()` call lands the bytes in the **OS page cache** — RAM managed by the kernel — not yet on the physical disk. The kernel marks those pages "dirty" and writes them to disk later, asynchronously. With `acks=all` (a.k.a. `acks=-1`), the leader waits until every **in-sync replica (ISR)** — followers that are caught up — has fetched and appended the record (also into *their* page cache). Only then does the leader acknowledge the producer. Crucially, **none of this involves a forced fsync**; acknowledgement is based on replication, not on the data hitting a physical platter/SSD. ## Why not fsync every message? - **Throughput collapse.** `fsync` is a synchronous, high-latency operation (it must wait for the storage device). Doing it per message (or even per small batch) would drop throughput by orders of magnitude. - **A single flushed copy is still fragile.** If you fsync to one broker and that broker's disk fails (or the file is corrupted), the data is gone. Flushing doesn't protect against disk/host loss. - **Replication is stronger AND faster.** Network round-trips to a few peers in the same datacenter are far cheaper than disk fsyncs, and the record survives as long as one ISR member survives. So replication gives you *better* durability at *lower* latency. ## The knobs involved - `acks=all` on the producer — wait for all ISR members. - `min.insync.replicas` (broker/topic) — the minimum ISR size for an `acks=all` write to succeed; if fewer replicas are in sync, the broker rejects the write with `NotEnoughReplicas`. Typical safe config: replication factor 3, `min.insync.replicas=2`, `acks=all` — tolerates one broker loss with no data loss. - `flush.messages` / `flush.ms` — *optional* knobs that force fsync after N messages or T milliseconds. Almost always left at defaults (effectively "let the OS decide") because replication already provides durability. ## The trade-off / residual risk Because acknowledged data may still be only in page cache on all replicas, a **correlated power loss** — every ISR broker losing power simultaneously before the OS flushes — can lose acknowledged records. In practice replicas are spread across racks/availability zones with independent power, making correlated loss rare, and that residual risk is accepted in exchange for huge throughput gains.
- When acks=all returns, is the data guaranteed to be on disk?No. It's guaranteed to be in the page cache (RAM) of all in-sync replicas. The physical flush to disk happens later, asynchronously, unless flush.messages/flush.ms force it sooner.
- What's the canonical no-data-loss configuration?Replication factor 3, min.insync.replicas=2, producer acks=all (plus enable.idempotence=true). This survives the loss of one broker without losing acknowledged writes.
saying these in an interview costs you the question
- Claiming acks=all means the data is fsynced to disk before acknowledgement.
- Saying Kafka fsyncs every message — it relies on the OS background flush by default.
- Believing flushing one replica's copy to disk is as safe as replicating to multiple brokers.