--:--
notes/commonplace/kafka/build-kafka-1-records-on-disk.mdx

NOTES / Commonplace / Kafka ·

Build your own Kafka, 1: records on disk

Chapter 0 decided that Kafka stores messages in an append-only log. Now: what exactly goes into the file?

The problem

A record has a key, a value and a timestamp. The naive file format writes them one after another. Three questions break it:

  1. Where does one record end? Keys and values are arbitrary bytes, so no separator character is safe. Each one needs a length in front.
  2. Was it written completely? If the machine dies mid-write, the file ends with half a record. A reader must be able to tell.
  3. How much does the framing cost? Kafka handles millions of small records per second. If every record carried 8-byte offsets, 8-byte timestamps and 4-byte lengths, the framing could be bigger than the data.

Batches

Kafka's answer to the third question is that records never travel or sit on disk alone. A producer sends a record batch, and the broker writes that batch to the file as is, byte for byte. Everything that is the same for all records in the batch is stored once, in a 61-byte header: the offset of the first record, the timestamp of the first record, the producer id, the leader epoch, and a checksum.

Each record then stores only what is different, and only as a delta: its offset minus the batch's base offset (0, 1, 2…) and its timestamp minus the base timestamp (often a few milliseconds). Those deltas are small numbers, so they are written as varints that take one byte when the number is small.

Type some records below and click any byte to see which field it is.

Things to notice:

  • One record with a 7-byte key and an 11-byte value costs 61 + 7 bytes of framing. Ten records in one batch cost 61 + 10 × 7, a little more once the deltas need two bytes. Batching is how Kafka makes small records cheap; the producer's linger.ms and batch.size settings exist to build bigger batches.
  • batchLength (bytes 8 to 11) tells a reader where the batch ends without parsing the records. Chapter 4 uses this to walk a damaged file.
  • A null key is length -1, a single byte 01. A null value is a tombstone, used by compaction in chapter 5.
  • The CRC covers every byte from attributes to the end. Pick a byte in a record and click Flip a bit: the CRC no longer matches and a broker would reject the batch with CORRUPT_MESSAGE.

Varints

A varint writes a number 7 bits per byte, lowest bits first, and sets the top bit of every byte except the last to mean "more bytes follow". So 0 to 127 take one byte, up to 16,383 take two, and so on.

Lengths can be -1, and timestamps can go backwards within a batch, so negative numbers must also be small. Two's complement would make -1 a huge unsigned number. Zigzag encoding fixes that by interleaving signs: 0 → 0, -1 → 1, 1 → 2, -2 → 3, 2 → 4. Then the 7-bit encoding runs on the result. Protocol Buffers uses the same scheme.

Write it

Write writeVarint. Every varint in the demo above goes through your function as soon as all the tests pass, so a mistake shows up as bytes Kafka can't read back.

Break it

  • Flip a bit in the last byte of batchLength. The batch now claims to be one byte longer or shorter than it is, so a reader either runs off the end of the file or checks the CRC over the wrong bytes.
  • Flip a bit in a record's value. Kafka can still parse the batch, but the CRC catches it.
  • Flip a bit in magic. The CRC still matches, because it starts after magic. A real broker checks magic on its own and refuses versions it doesn't know.

In Kafka

  • The batch format is DefaultRecordBatch and each record is DefaultRecord, in clients/src/main/java/org/apache/kafka/common/record/. The comment at the top of DefaultRecordBatch has the same field table as the demo.
  • Varints: ByteUtils.writeVarint and writeVarlong in org.apache.kafka.common.utils.
  • The checksum is CRC-32C (Castagnoli), Crc32C in the same package. Modern CPUs have an instruction for it.
  • The format arrived in Kafka 0.11 with KIP-98, which needed producer ids and sequence numbers in the header for idempotence and transactions (chapters 7 and 14).
  • What the demo leaves out: compression. When attributes names a codec, everything after recordCount is compressed as one block, which is another reason batches matter.