What do flush.messages and flush.ms control, and why does Kafka discourage relying on broker fsync for durability?
answer
- dirty pages -> fsync forces to disk
- flush.messages / flush.ms default ~never
- durability = replication, acks=all, min.insync.replicas
- fsync breaks batched sequential writeback
- correlated AZ/power loss is the exception
basics
~20 sflush.messages and flush.ms force the broker to fsync log data to disk after N messages or T milliseconds. By default both are effectively unset, so Kafka lets the OS flush dirty pages lazily and relies on replication, not local fsync, for durability — forcing frequent fsync hurts throughput.
solid answer
~50 sWhen a broker appends records, they land in the page cache as **dirty pages**; they are not yet physically on disk. `log.flush.interval.messages` (`flush.messages` per-topic) and `log.flush.interval.ms` (`flush.ms`) let you force an **fsync** after every N messages or T ms, guaranteeing those bytes hit the disk platter. Kafka's deliberate default is to leave these effectively **disabled** (max values), letting the OS writeback (`vm.dirty_*`) flush dirty pages in the background. Durability instead comes from **replication**: with `acks=all`, `replication.factor`, and `min.insync.replicas`, a record is acknowledged only once it is in the page cache of enough in-sync replicas, so a single node losing its un-fsynced pages in a crash does not lose committed data — another replica has it. Forcing frequent fsync serializes I/O, defeats batched sequential writeback, and tanks throughput and latency. The principled stance: tolerate single-node un-fsynced loss and buy durability with replicas across failure domains; only tighten flush settings for special single-broker or correlated-failure (e.g. datacenter power loss) scenarios.
go deeper
Know that data first sits in memory (page cache) and fsync forces it to disk; Kafka usually lets the OS handle it.
Explain flush.messages/flush.ms defaults and that durability comes mainly from replication.
Articulate the fsync-vs-throughput trade-off and the acks/min.insync.replicas model.
Reason about correlated failure domains, when to deviate from defaults, and design durability across AZs/power domains.
## Dirty pages and fsync When the broker writes a record batch, the bytes go into the **page cache** as **dirty pages** — modified pages the kernel will eventually write to the physical disk during **writeback**. Until that happens, a power loss or kernel crash on that machine loses those bytes. The `fsync(2)` syscall forces all dirty pages of a file to the storage device and returns only when they are durable. fsync is **expensive**: it stalls until the device confirms, and it breaks the kernel's ability to batch many writes into efficient sequential flushes. ## The flush settings - **`log.flush.interval.messages`** (topic override `flush.messages`): fsync the log after this many messages have been appended. - **`log.flush.interval.ms`** (topic override `flush.ms`): fsync the log at most this often in time. - `log.flush.scheduler.interval.ms`: how often the flusher thread checks whether a flush is due. By **default these are set to effectively 'never'** (Long.MAX_VALUE), meaning Kafka does **not** proactively fsync on a schedule — it leaves flushing to the OS. ## Why Kafka prefers replication over fsync Kafka's durability model is **replication-based**: - A topic has `replication.factor` copies across brokers. - Producers set **`acks=all`** so the leader waits for all **in-sync replicas (ISR)** to receive the record before acknowledging. - **`min.insync.replicas`** sets how many replicas must be in sync for a write to be accepted; if fewer are available, the producer gets an error rather than a silent under-replicated write. Under this model, a committed record exists in the page cache of multiple brokers, ideally in **different failure domains** (racks/AZs). If one broker crashes and loses its un-fsynced pages, the data still lives on the surviving replicas, which become leader. So durability does **not** require each broker to fsync every write. This lets Kafka keep writes **sequential, batched, and OS-flushed** — the basis of its throughput. ## Why frequent fsync hurts - It serializes I/O and adds latency on the hot path. - It prevents the kernel from coalescing many dirty pages into large contiguous sequential flushes. - It can turn a bandwidth-bound workload into a latency/IOPS-bound one. Hence the docs explicitly recommend leaving flush intervals at defaults and relying on replication, tuning OS writeback (`vm.dirty_background_ratio`, `vm.dirty_ratio`, `vm.dirty_expire_centisecs`) if you want to smooth flush behavior. ## When you might tighten flush - **Single-broker / no replication** deployments where there is no replica to recover from. - **Correlated failures**: a whole-datacenter power loss can crash all replicas simultaneously, so 'another replica has it' fails — some operators fsync or place replicas across power/AZ domains to mitigate. - Regulatory/strict-durability requirements where even a rare un-fsynced loss window is unacceptable. ## Related durability levers (not flush) - `acks`, `min.insync.replicas`, `replication.factor`, `unclean.leader.election.enable=false` (don't elect an out-of-sync replica as leader and lose data), rack/AZ-aware replica placement. ## Edge cases / gotchas - fsync flushes data pages but Kafka also relies on consistent recovery on restart (log recovery, index rebuild) for partially written tails. - Setting `flush.messages=1` (fsync per message) is a classic anti-pattern that cripples throughput. - 'acks=all' guarantees replicas have it in **page cache/memory**, not necessarily fsynced — which is exactly the deliberate trade-off.
- If Kafka doesn't fsync per write, how is committed data not lost when a broker crashes?With acks=all and min.insync.replicas, a committed record is already in the page cache of multiple in-sync replicas across failure domains; a crashed broker's lost un-fsynced pages are recovered from the surviving replicas that become leader.
- When is relying purely on replication insufficient, motivating tighter flush settings?When failures are correlated — e.g. a whole-datacenter power loss crashing all replicas at once, or single-broker deployments with no replica. Then operators may fsync more aggressively or spread replicas across independent power/AZ domains.
- Why is flush.messages=1 considered an anti-pattern?It forces an fsync per message, serializing I/O and preventing batched sequential writeback, which collapses throughput while replication already provides durability.
saying these in an interview costs you the question
- Claiming Kafka fsyncs every message by default.
- Saying acks=all guarantees data is fsynced to disk (it guarantees replicas have it in memory/page cache).
- Recommending flush.messages=1 for 'safety' on a replicated cluster.
- Ignoring correlated-failure scenarios where replication alone is insufficient.