--:--
notes/commonplace/kafka/build-kafka-0-why-a-log.mdx

NOTES / Commonplace / Kafka ·

Build your own Kafka, 0: why a log

This course rebuilds Kafka from nothing. Each chapter starts with a problem, you write the one function that solves it, and the demo next to it switches over to your code. By the last chapter, every part of the cluster on the Kafka page is something you have written yourself.

This chapter has no code. It is about the single decision everything else follows from.

Two ways to hand over messages

Say an order service has to tell three other services about every new order: billing, shipping and analytics. The obvious tool is a queue: the producer puts a message in, a consumer takes it out, and the queue deletes it. That is how classic message brokers work, and it is a good fit when one worker should do each job once.

Now try it with two readers. Use the demo below: send a few messages, then let A and B take turns.

With the queue, each message goes to exactly one reader. B never sees what A took. If analytics had a bug last night and wants to read yesterday's orders again, it can't: they are gone. Adding a fourth service means changing the producer or the broker to copy messages into a new queue.

Switch the demo to Log. Now nothing is deleted when it is read. The broker only appends, and each reader keeps one number, its offset: the position of the next message it will read. Readers don't affect each other at all. B can rewind to 0 and read everything again, while A carries on.

What the log buys

That small change, "don't delete on read, let readers keep a position", has big consequences:

  • Any number of readers. A new service starts at offset 0 (or at the end) without anyone else noticing.
  • Replay. Fix a bug, reset the offset, process the history again.
  • Order. Every reader sees the same messages in the same order. Two services that read the same log end up agreeing.
  • Cheap writes. Appending to the end of a file is the fastest thing a disk does. There are no per-message acknowledgements to track on the broker, only one offset per reader.
  • The broker stays simple. It doesn't need to know who has read what. It stores bytes in order and serves ranges of them.

Jay Kreps, one of Kafka's authors, put it this way in The Log: a log is the simplest possible storage abstraction, an append-only, totally ordered sequence of records ordered by time. Databases already use one internally (the write-ahead log), and so do consensus systems. Kafka makes it the product.

What it costs

A log never deletes on read, so something else has to delete old data, or the disk fills up. And a single file on a single machine can only grow so fast, can only be read so fast, and dies with that machine. The rest of the course is a series of answers to these problems:

ChapterProblemWhat you build
1How are records laid out on disk?The record batch format
2Reading offset 1,000,000 can't mean scanning from the start.A sparse index and binary search
3The file grows forever.Segments and retention
4Writing to disk is slow; not writing is unsafe.The page cache, flushing, crash recovery

Later chapters split the log into partitions, copy it to other machines, elect leaders, coordinate consumers and add transactions.

Try it

  • Send five messages in queue mode, let A take two and B take one. How many messages could a third reader still get?
  • Switch to log mode and do the same. Then rewind B. Which numbers would you need to store to remember everyone's progress?