Page Cache and Sequential I/O: The Engine Behind Kafka Throughput
Kafka achieves high throughput primarily because it turns random I/O into sequential I/O and lets the operating system page cache absorb the difference between producer and consumer speeds. When a producer writes to a partition, Kafka appends the record to the active segment file. The write goes into the OS page cache first, not directly to disk. The OS flushes dirty pages to disk lazily in the background. Similarly, when a consumer reads, Kafka serves the data from the page cache if it is still resident, avoiding a disk read entirely. This is why Kafka can sustain hundreds of thousands of messages per second on modest hardware: the hot working set lives in RAM managed by the OS, not by the JVM heap.
The mechanism has several important consequences. First, Kafka avoids the JVM garbage collection cost of caching messages in heap. Instead, it relies on the OS, which is very good at managing memory and can evict pages under pressure without stopping the world. Second, because writes are sequential appends and reads are sequential scans, the OS can prefetch aggressively and the disk head (or SSD controller) stays busy with large contiguous operations. Third, Kafka uses zero-copy transfer (sendfile) when sending data to consumers over the network, which means the bytes go from the page cache to the network socket without being copied into user space. This is a significant CPU saving.
A common misconception is that Kafka is fast because it fsyncs every message. The opposite is true: fsyncing every message would destroy throughput. Kafka's durability model relies on replication, not on immediate disk persistence. With acks=all and replication factor 3, a message is considered committed once all in-sync replicas have written it to their page caches, not necessarily to disk. This is a trade-off: you get high throughput and good durability, but a simultaneous power loss across all replicas can lose recent data. For workloads that require true disk durability, you can set log.flush.interval.messages=1, but expect a massive throughput drop. Another misconception is that page cache is a Kafka feature; it is an OS feature, and Kafka's design is what exploits it well. If you run Kafka in a container with a memory limit that does not account for page cache, you can starve the cache and see throughput collapse. In Kubernetes, for example, page cache is not counted against the container's memory limit in the same way as heap, but cgroup v2 and memory pressure can still evict it.
Kafka writes to the OS page cache first; the OS flushes to disk lazily.
Consumers read from page cache when data is hot, avoiding disk I/O.
Sequential appends and scans let the OS prefetch and the disk stay busy with large contiguous operations.
Zero-copy (sendfile) moves bytes from page cache to socket without copying into user space.
Durability comes from replication and acks, not per-message fsync; fsyncing every message kills throughput.
In containers, page cache pressure can evict hot data and collapse throughput; monitor memory and cgroup limits.
0-2 years experience
2-5 years experience
5-8 years experience
8+ years experience