Build your own Kafka, 2: offsets and the index
The file from chapter 1 is a run of batches back to back. Every fetch request says "give me records starting at offset N". The broker has to find the byte position of offset N.
The problem
Offsets are not byte positions. Batches have different sizes, so offset 217 could be anywhere. Without help, the only way is to start at byte 0, read each batch's header, jump batchLength bytes ahead, and repeat until a batch contains 217. On a 1 GiB file that is millions of header reads for every fetch.
The design question
Before reading on: if you could keep a second file next to the log to make this fast, what would it contain?
- An entry for every record: offset → position. Lookups are instant, but the index is as long as the log has records, and every append writes to two files.
- An entry every N bytes of log. The index is tiny, but the broker lands somewhere before the record and must scan forward a bit.
Kafka picks the second, a sparse index. Every time index.interval.bytes (default 4096) more bytes have been appended, it adds an 8-byte entry: the batch's last offset (stored relative to the segment's base offset, in 4 bytes) and the batch's byte position (4 bytes). Finding offset N is then:
- Find the last index entry with offset ≤ N.
- Jump to its position in the log.
- Scan forward batch by batch until one contains N. At most about
index.interval.bytesof reading.
Change the interval below and watch the trade-off: dark cells in the index are the entries the lookup read, gold cells in the log are the batches it scanned, and the outlined cell is where it landed.
With every batch, the index is as big as the log has batches and the scan is a single batch. With no index, the scan walks hundreds of batches. With 4096 bytes the index has a handful of entries and the scan stays short. Since the index is tiny, Kafka keeps it in memory: it is a memory-mapped file, preallocated to segment.index.bytes (10 MiB) and trimmed when the segment closes.
Write it
The index is sorted by offset, so step 1 is a search for "the largest entry ≤ target". A loop over all entries is correct, but this runs on every fetch against indexes with hundreds of thousands of entries, and one of the tests counts how many entries you read.
Break it
Once your version runs the demo, imagine an off-by-one: returning the first entry greater than the target instead. The scan would then start after the record and stop at the wrong batch. The demo checks for this and says so, but Kafka itself trusts its index completely, which is why a broker that finds a corrupt index file on startup deletes and rebuilds it instead of using it.
In Kafka
OffsetIndexand its parentAbstractIndexin the storage module (org.apache.kafka.storage.internals.log).lookupcallslargestLowerBoundSlotFor, a binary search.- A refinement you can now appreciate: consumers almost always read near the end of the log. A plain binary search over a huge memory-mapped index starts in the middle, touching pages that may have been evicted from memory. Since KAFKA-6432 the search first looks at the warm section, the last 8 KiB of entries, and only falls back to the whole index when the target is older.
- There is a second index per segment, the
.timeindex: timestamp → offset, 12 bytes per entry. It answers "where was the log at 9:00?" foroffsetsForTimesand time-based retention. - The index stores offsets relative to the segment's base offset, so 4 bytes are enough. That only works because segments are bounded, which is the next chapter.
Your code is saved in this browser only.