Kafka achieves massive throughput on standard mechanical disks and SSDs by writing data sequentially to an append-only log, relying on the operating system page cache, and streaming bytes to network sockets with zero-copy system calls. Throughput is multiplied across the cluster by batching and compressing records, distributing topic partitions across multiple brokers, and letting consumers pull data independently at their own speed.
The sequential disk advantage
Spinning disks and solid-state drives both deliver high throughput when handling sequential writes compared to random seeks. Random writes force disk heads to physically move or flash controllers to execute complex block erase cycles, capping throughput at a few megabytes per second. By treating partitions as append-only commit logs where new events simply attach to the end of the active segment file, Kafka achieves hundreds of megabytes per second of sustained write speed on inexpensive storage.
Operating system page cache and zero copy
Instead of managing a complex in-memory cache inside the JVM heap, Kafka delegates caching entirely to the Linux OS page cache. Storing cached data in the JVM would introduce heavy garbage collection pauses, double memory footprint due to Java object overhead, and lose cached contents during broker restarts. The OS page cache keeps hot segments in kernel memory automatically.
When a consumer reads data, Kafka bypasses user-space memory entirely using the Linux sendfile system call, often referred to as zero-copy data transfer. Under standard network transfers, data moves from kernel disk cache to user space memory and back to the kernel socket buffer. With zero copy, the operating system copies data directly from the page cache into the network interface buffer, eliminating CPU cycles and memory bus contention.
Batches, partitions, and the pull model
Producers accumulate individual messages into batches before sending them across the wire. Applying compression codecs like snappy, lz4, or zstd to a full batch yields much higher compression ratios than compressing single messages, conserving network bandwidth and disk space.
Topics are split into partitions across multiple brokers, spreading I/O load horizontally across separate machines and drive controllers. Finally, consumers pull messages rather than having brokers push them. This pull model prevents fast brokers from overwhelming slower downstream consumers and avoids stateful broker buffer queues, allowing each consumer group to read at its own maximum processing capacity.