--:--
kafka/build/3-segments-and-retention.mdx

KAFKA / Build your own Kafka / Chapter 3 ·

Build your own Kafka, 3: segments and retention

Chapter 0 traded "delete on read" for "keep everything". Now the bill arrives: the disk fills up.

The problem

Deleting the oldest records from the front of a file is not something file systems do cheaply. You would have to rewrite the whole file without them, while producers keep appending and consumers keep reading. A single growing file also makes the index from chapter 2 awkward: 4-byte relative offsets and positions only work up to 2 GiB.

Segments

Kafka never has one file per partition. It has a directory of segments, each a .log file plus its .index and .timeindex, named after the first offset it holds, padded to 20 digits:

orders-0/
  00000000000000000000.log    offsets 0 to 8,421
  00000000000000008422.log    offsets 8,422 to 16,950
  00000000000000016951.log    the active segment, still being written

Only the newest, the active segment, is ever written to. When appending a batch would make it bigger than segment.bytes (default 1 GiB), or it is older than segment.ms (default 7 days), the broker rolls: it closes it and starts a new file named after the next offset. Closed segments never change again.

Now deleting old data is trivial: delete a whole file. No rewriting, no locking out readers of other segments. Finding offset N first means finding its segment, the one with the largest base offset ≤ N (the broker keeps them in a sorted map), then the index lookup from chapter 2 inside it.

Try it: produce for a while with segment.bytes at 512, then set retention.ms to 1 h and advance the clock. Segments disappear from the front and the log start offset jumps forward. Your slow consumer is still at offset 0; read next, and the broker answers OFFSET_OUT_OF_RANGE.

Retention

The rules a broker applies every log.retention.check.interval.ms (5 minutes):

  • retention.ms (default 7 days): delete a segment when its newest record is older than this. Using the newest timestamp, not the file's age, means a segment is only deleted once everything in it has expired.
  • retention.bytes (default off, per partition): delete the oldest segments while the partition is still at least this big without them.
  • Only from the front: the log must stay a contiguous range of offsets, so if an old segment is kept, everything after it is kept too.
  • Never the active segment. A partition that gets no writes keeps its last segment until segment.ms rolls it.

Two consequences people trip over: retention is per segment, so data can live up to one segment's worth longer than retention.ms says; and setting retention.bytes on a topic limits each partition, not the topic.

Write it

When your version passes, the demo above runs on it. Try retention.bytes at 2048 with segment.bytes at 1024 and check that your rule leaves at least 2048 bytes.

Break it

  • Set segment.bytes to 2048 and retention.ms to 1 h, produce once, then advance the clock many hours. Nothing is deleted: everything is in the active segment.
  • Let the consumer fall behind, then click Reset to earliest. That is what auto.offset.reset=earliest does after OFFSET_OUT_OF_RANGE: it skips to the log start offset, and whatever it never read is gone without an error.

In Kafka

  • LogSegment is one segment (its .log, .index and .timeindex files); UnifiedLog owns the segments of one partition and decides when to roll. Both live in the storage module, org.apache.kafka.storage.internals.log.
  • Retention runs in UnifiedLog.deleteOldSegments, which applies time and size limits through predicates over the segment list and deletes a prefix, like your function. Deletion is asynchronous: segments are first renamed with a .deleted suffix and removed after file.delete.delay.ms.
  • Time retention used to use the file's modification time, which moved whenever a replica was copied. KIP-33 (0.10.1) switched to the largest record timestamp in the segment and added the time index.
  • The other way to bound a log is compaction: keep the newest record per key instead of the newest records overall. That is chapter 5.

Your code is saved in this browser only.