--:--
notes/commonplace/kafka/build-kafka-4-page-cache-and-recovery.mdx

NOTES / Commonplace / Kafka ·

Build your own Kafka, 4: the page cache and crash recovery

So far the log has been a file and appending has been instant. Real disks don't work like that, and Kafka's answer is one of its least intuitive design choices.

The problem

When a program calls write(), the bytes don't go to the disk. They are copied into the operating system's page cache, memory the kernel uses to hold file contents, and the call returns. The kernel writes dirty pages back to the disk later, a few seconds or tens of seconds later, in page-sized pieces (usually 4 KiB). Only fsync() forces them out, and it waits for the device: on a spinning disk that can mean milliseconds, a hundred times slower than the write.

So a broker can choose:

  • fsync every batch: nothing acknowledged is ever lost in a power cut, and throughput collapses.
  • Let the OS decide: writes are memory-speed, and a power cut loses whatever was still in the page cache. Worse, the kernel may have written back some pages of a batch and not others.

Kafka's choice

By default Kafka never calls fsync for durability. log.flush.interval.messages and log.flush.interval.ms are effectively off, and the design documentation argues for leaving the page cache in charge:

  • The page cache is already a cache. Keeping records in the JVM as well would store them twice and fight the garbage collector, so Kafka keeps nothing in its own heap and reads from the page cache too. Recent data, which consumers almost always want, is served from memory.
  • Sequential writes and reads are what disks and the page cache handle best, so the log stays fast even as it grows.
  • When nothing has to be transformed, the broker sends file bytes straight to the network socket with sendfile (Java's FileChannel.transferTo), so data never even enters the JVM: zero copy. TLS turns this off, because bytes must be encrypted in user space.
  • Durability comes from replication instead (chapters 8 to 10): a record acknowledged with acks=all is in the page cache of several brokers, and losing power on all of them at once is far less likely than on one.

The demo shows one segment. Green bytes are on disk; gold bytes are only in the page cache; ticks are page boundaries; white lines separate batches.

Produce several batches, then Cut power. The OS had written some dirty pages back, so the file usually ends in the middle of a batch (red). Then Restart the broker.

Recovery

On startup a broker checks whether the last shutdown was clean (it leaves a .kafka_cleanshutdown marker). If not, it recovers the log: starting from the last offset it knows was flushed (the recovery point, saved in recovery-point-offset-checkpoint), it reads every batch, checks that it is complete and that its CRC matches, and rebuilds the indexes as it goes. At the first bad batch it truncates the file and deletes any later segments.

That walk only works because of two decisions from chapter 1: every batch starts with its length, so the broker can step from batch to batch without parsing records, and every batch carries a CRC, so a half-written batch is detected rather than read as garbage.

Write it

With your version running, cut the power and restart. If it keeps a torn batch, the demo tries to read the file afterwards and shows where it breaks.

Break it

  • Set Flush every to 1 batch. Power cuts now lose nothing, and you have just made every write wait for the disk.
  • Leave flushing off, produce, cut power and restart several times. Count how many records disappear. These were records the broker had appended, and, with acks=1, already acknowledged to the producer. The scenario acks=1 loses a write on the Kafka page shows the cluster-level version of this.

In Kafka

  • LogSegment.recover walks the batches, validates each one (batch.ensureValid() checks the CRC), rebuilds .index and .timeindex, and truncates at the first invalid byte. LogLoader decides which segments need recovery after an unclean shutdown.
  • UnifiedLog.flush writes up to an offset and advances the recovery point; the log.flush.* settings and flush.messages / flush.ms per topic control it.
  • FileRecords.writeTo is where transferTo sends a slice of a segment file to a socket.
  • Jack Vanlightly's Why Apache Kafka doesn't need fsync to be safe explains when skipping fsync is safe (with replication and a correct recovery protocol) and when it is not.