- Published on
Build a Mini Kafka Broker to Actually Understand Kafka
- Authors

- Name
- Mehdi Akiki
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):
| file | purpose |
|---|---|
log.jsonl | append-only record log |
offsets.json | committed offset per consumer group |
- produce = append a record and assign it the next offset
- consume = fetch records starting at an offset
- commit = store
next_offsetso 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:
- consumer says "give me records starting at offset X"
- broker looks up X in the index, finds the file position
- broker reads from there
- 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:
- broker collects events into a batch
- appends batch to end of log file (one write)
- OS keeps recent pages in RAM
- consumer fetches starting at offset 500000
- broker streams bytes directly from page cache
- 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