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
  1. Append: length-prefixed records with offsets assigned by the log.
  2. Segments and a sparse index: seek, then scan.
  3. Recovery stops at the last complete record.
  4. Partitions by key: order is per key, and work runs in parallel across keys.
  5. Consumer groups: assignment, rebalance, commit after processing.
  6. Replication: the ISR, the high watermark, and failover without losing committed records.
real Kafkatiny-kafka
record batches with CRC, compression and headers; zero-copy sendfile to consumersone JSON record per entry, read through the process
.index and .timeindex files, memory-mappedan in-memory sparse index rebuilt on open
a binary protocol over TCP with fetch long-pollingan in-process API
group coordinator, heartbeats, cooperative rebalancingsynchronous join and leave with range assignment
leader epochs, follower truncation, KRaft metadata quorumthe HW, ISR and failover sketch
log compaction, transactions, idempotent producersexercises
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).