Part 2 · 2 chapters · ~20 min
M2: Kafka
A partition as an append-only file: length-prefixed records, segments named by base offset, a sparse index, torn-tail recovery; keyed partitions, consumer groups with assignment, rebalance and committed offsets; and replication with the ISR and a high watermark.
4
How Kafka works, from the file up
one data structure, everything else on top
- Append: length-prefixed records with offsets assigned by the log.
- Segments and a sparse index: seek, then scan.
- Recovery stops at the last complete record.
- Partitions by key: order is per key, and work runs in parallel across keys.
- Consumer groups: assignment, rebalance, commit after processing.
- Replication: the ISR, the high watermark, and failover without losing committed records.
| real Kafka | tiny-kafka |
|---|---|
| record batches with CRC, compression and headers; zero-copy sendfile to consumers | one JSON record per entry, read through the process |
| .index and .timeindex files, memory-mapped | an in-memory sparse index rebuilt on open |
| a binary protocol over TCP with fetch long-polling | an in-process API |
| group coordinator, heartbeats, cooperative rebalancing | synchronous join and leave with range assignment |
| leader epochs, follower truncation, KRaft metadata quorum | the HW, ISR and failover sketch |
| log compaction, transactions, idempotent producers | exercises |
TINY-KAFKA: A LOG ON DISK
records appended to segment files, found again through a sparse index, split by key, consumed by groups, replicated to a high watermark
swipe the figure sideways, or tap expand for full screen
1/6
append
Append: a record becomes [8-byte offset][4-byte length][JSON body] at the end of the active segment file. Offsets are assigned by the log, increasing by one; the producer gets the offset back. Nothing is ever modified in place.
5
Building it: tiny-kafka
Repo: repos/kafka. Four files: the log, the broker, a consumer, and the replication sketch.
code
// src/log.ts: recovery scans to the last complete record; a torn tail is simply not counted
private recover() {
let pos = 0; const head = Buffer.alloc(HEADER);
while (pos + HEADER <= this.size) {
readSync(this.fd, head, 0, HEADER, pos);
const offset = Number(head.readBigInt64BE(0)); const len = head.readUInt32BE(8);
if (pos + HEADER + len > this.size) break; // torn: stop here
if ((offset - this.base) % INDEX_EVERY === 0) this.index.push([offset, pos]);
this.next = offset + 1;
pos += HEADER + len;
}
this.size = pos; // the next append overwrites the torn bytes
}code
// src/broker.ts: range assignment; every join or leave is a rebalance
private rebalance(g: Group, topic: string) {
const n = this.partitions(topic).length;
const members = [...g.members].sort();
const per = Math.floor(n / members.length), extra = n % members.length;
let p = 0;
members.forEach((m, k) => { const count = per + (k < extra ? 1 : 0); for (let c = 0; c < count; c++) g.assignment.get(m)!.push(p++); });
g.generation++;
}code
// src/replication.ts: committed means "on every in-sync replica"
get highWatermark(): number {
const leos = [this.leader, ...this.followers].filter(r => this.isr.has(r.id)).map(r => r.leo);
return Math.min(...leos);
}
read(from: number, max = 100) { return this.leader.log.read(from, max).filter(r => r.offset < this.highWatermark); }Run it. In
repos/kafka: npm test. Then read test/kafka.test.ts top to bottom. It produces 50 records into tiny segments to force rolling, reopens the log, rebalances a group, restarts a broker and resumes from committed offsets, and fails over a partition whose slowest follower has left the ISR.exercises
1. Log compaction: rewrite closed segments keeping the last record per key. 2. Idempotent producer: producer id plus a sequence per partition, with duplicates dropped at the broker. 3. Leader epochs: a returning old leader truncates its divergent tail to the new leader's epoch boundary. 4. A fetch API over HTTP with long polling (use the M10 server).