Published on

Build a Mini Kafka Broker to Actually Understand Kafka

Authors
  • Mehdi Akiki avatar
    Name
    Mehdi Akiki
    Twitter

The right way to learn Kafka is to build a tiny log broker, not a fake "queue with ack/delete." That is the core mistake beginners make.

Kafka is a distributed append-only log, closer to a database write-ahead log than a traditional message queue. Consumer progress is one integer per partition: the offset of the next record to read. Because Kafka stores committed offsets separately, consumers can resume, rewind, and reprocess old data instead of having messages disappear after an ack.

The mental model

One topic = one partition (for this toy version):

filepurpose
log.jsonlappend-only record log
offsets.jsoncommitted offset per consumer group
  • produce = append a record and assign it the next offset
  • consume = fetch records starting at an offset
  • commit = store next_offset so the group can resume after a crash

That is the whole thing.

FastAPI implementation

from __future__ import annotations

from collections import defaultdict
from pathlib import Path
from threading import Lock
from typing import Any
import json
import time

from fastapi import FastAPI, HTTPException, Query
from pydantic import BaseModel

app = FastAPI(title="Mini Kafka")

DATA_DIR = Path("data")
DATA_DIR.mkdir(exist_ok=True)
locks = defaultdict(Lock)


class ProduceIn(BaseModel):
    key: str | None = None
    value: Any


class CommitIn(BaseModel):
    group: str
    next_offset: int


def paths(topic: str) -> tuple[Path, Path, Path]:
    d = DATA_DIR / topic
    return d, d / "log.jsonl", d / "offsets.json"


def read_json(path: Path, default: Any) -> Any:
    if not path.exists():
        return default
    return json.loads(path.read_text(encoding="utf-8"))


def write_json(path: Path, data: Any) -> None:
    tmp = path.with_suffix(path.suffix + ".tmp")
    tmp.write_text(json.dumps(data, indent=2), encoding="utf-8")
    tmp.replace(path)


def ensure_topic(topic: str) -> None:
    d, log_path, offsets_path = paths(topic)
    d.mkdir(parents=True, exist_ok=True)
    if not log_path.exists():
        log_path.touch()
    if not offsets_path.exists():
        write_json(offsets_path, {})


def high_watermark(topic: str) -> int:
    """Next offset that would be assigned. This toy scans the log — real Kafka does not."""
    ensure_topic(topic)
    _, log_path, _ = paths(topic)
    next_offset = 0
    with log_path.open("r", encoding="utf-8") as f:
        for line in f:
            if line.strip():
                next_offset = json.loads(line)["offset"] + 1
    return next_offset


def committed_offset(topic: str, group: str) -> int:
    ensure_topic(topic)
    _, _, offsets_path = paths(topic)
    return int(read_json(offsets_path, {}).get(group, 0))


@app.post("/topics/{topic}")
def create_topic(topic: str):
    ensure_topic(topic)
    return {"topic": topic, "status": "created"}


@app.post("/topics/{topic}/produce")
def produce(topic: str, body: ProduceIn):
    with locks[topic]:
        ensure_topic(topic)
        _, log_path, _ = paths(topic)
        offset = high_watermark(topic)
        record = {
            "offset": offset,
            "ts_ms": int(time.time() * 1000),
            "key": body.key,
            "value": body.value,
        }
        with log_path.open("a", encoding="utf-8") as f:
            f.write(json.dumps(record, separators=(",", ":")) + "\n")
            f.flush()
        return {"topic": topic, "offset": offset, "record": record}


@app.get("/topics/{topic}/consume")
def consume(
    topic: str,
    group: str | None = None,
    offset: int | None = None,
    limit: int = Query(default=10, ge=1, le=100),
):
    ensure_topic(topic)
    _, log_path, _ = paths(topic)
    hw = high_watermark(topic)
    start = offset if offset is not None else (committed_offset(topic, group) if group else 0)

    if start < 0 or start > hw:
        raise HTTPException(400, detail=f"offset out of range: start={start}, high_watermark={hw}")

    messages: list[dict] = []
    next_offset = start

    with log_path.open("r", encoding="utf-8") as f:
        for line in f:
            if not line.strip():
                continue
            rec = json.loads(line)
            if rec["offset"] < start:
                continue
            messages.append(rec)
            next_offset = rec["offset"] + 1
            if len(messages) >= limit:
                break

    return {
        "topic": topic,
        "group": group,
        "start_offset": start,
        "next_offset": next_offset,
        "high_watermark": hw,
        "messages": messages,
    }


@app.post("/topics/{topic}/commit")
def commit(topic: str, body: CommitIn):
    with locks[topic]:
        ensure_topic(topic)
        _, _, offsets_path = paths(topic)
        hw = high_watermark(topic)
        if body.next_offset < 0 or body.next_offset > hw:
            raise HTTPException(400, detail=f"next_offset must be between 0 and {hw}")
        offsets = read_json(offsets_path, {})
        offsets[body.group] = body.next_offset
        write_json(offsets_path, offsets)
        return {"topic": topic, "group": body.group, "committed_offset": body.next_offset}


@app.get("/topics/{topic}/groups/{group}")
def group_state(topic: str, group: str):
    ensure_topic(topic)
    hw = high_watermark(topic)
    committed = committed_offset(topic, group)
    return {
        "topic": topic,
        "group": group,
        "committed_offset": committed,
        "high_watermark": hw,
        "lag": hw - committed,
    }

Run it

pip install fastapi uvicorn pydantic
uvicorn main:app --reload

Try it

# Create topic
curl -X POST http://127.0.0.1:8000/topics/orders

# Produce
curl -X POST http://127.0.0.1:8000/topics/orders/produce \
  -H "Content-Type: application/json" \
  -d '{"key":"user-1","value":{"amount":42}}'

# Consume (from committed offset for group, or 0 if none)
curl "http://127.0.0.1:8000/topics/orders/consume?group=analytics&limit=10"

# Commit the next offset to read
curl -X POST http://127.0.0.1:8000/topics/orders/commit \
  -H "Content-Type: application/json" \
  -d '{"group":"analytics","next_offset":1}'

# Check lag
curl "http://127.0.0.1:8000/topics/orders/groups/analytics"

How the offset actually works

The commit stores next_offset, not "last seen offset." That is the key detail.

If a consumer reads records 5, 6, 7 and finishes processing them, it commits 8. The next fetch starts at 8.

In real Kafka the committed offset is stored durably in the internal __consumer_offsets topic and served from cache by the group coordinator. In this toy broker it is a JSON file. The model is identical.

Why this toy scans, and real Kafka does not

This broker scans log.jsonl on every high_watermark() call because it is for learning. Real Kafka maintains segment files plus a sparse offset index that maps offsets to file positions. A fetch starting at offset X jumps directly to the right byte in the right segment without reading anything before it.

So "how does Kafka know where to place the cursor?" becomes:

  1. consumer says "give me records starting at offset X"
  2. broker looks up X in the index, finds the file position
  3. broker reads from there
  4. consumer commits the new next offset later

The cursor is a logical position number, not a mutable queue head.

What this toy is missing

Real Kafka also has:

  • multiple partitions
  • leader/follower replication
  • segment files + sparse indexes
  • batching and compression
  • long polling
  • consumer group rebalancing
  • retention and log compaction
  • checksums and crash recovery

Why Kafka is fast

Four cheap operations, stacked together.

1. It appends at the end of a file

Not "insert somewhere." Not "update a row." Not "rebalance a tree."

Just:

[msg0][msg1][msg2]  →  [msg0][msg1][msg2][msg3][msg4][msg5]

The broker never needs to search for where to put data.

2. It batches messages

Instead of:

write(msg1)
write(msg2)
write(msg3)
...1000 times

Kafka does:

write(big_chunk_of_many_messages)  ← once

Every write has syscall overhead, kernel work, and flush behavior. 1000 small writes are far more expensive than 1 write of ~1 MB. Batching is one of the main reasons Kafka gets high throughput.

3. Consumers often read from RAM, not disk

When Kafka writes to a file, Linux keeps recent pages in the OS page cache (RAM). A consumer asking for data shortly after it was produced usually gets it straight from memory:

producer → broker → OS page cache (RAM) → consumer

not:

producer → disk → broker → disk again → consumer

Kafka writes like a log and the OS turns a lot of that into memory-speed reads. "Kafka is fast because disk is fast" misses this.

4. It tracks almost no per-message state

A conventional queue marks every message: delivered, acknowledged, deleted. Kafka does none of that. For each consumer group, the state is just:

partition 0 → next offset = 120845
partition 1 → next offset = 9981
partition 2 → next offset = 551002

The message stays in the log. The consumer tracks one number per partition. That is much cheaper bookkeeping.

The full picture in one example

10,000 sensor events arrive:

  1. broker collects events into a batch
  2. appends batch to end of log file (one write)
  3. OS keeps recent pages in RAM
  4. consumer fetches starting at offset 500000
  5. broker streams bytes directly from page cache
  6. consumer commits 510000

No deletes per message. No row updates. No index maintenance. No broker-side ack bookkeeping. Just appends and offset numbers.

That is why it is fast.

I build and scale reliable production systems. Open to full-time and freelance work with U.S.-based teams that value ownership and execution.

Got something in mind?

Book a Discovery Call