Kafka for One Machine: An Embedded Event Log With Offsets, Compaction and No Broker
A desktop app wants an undo history that survives a crash. An edge gateway buffers sensor readings while its uplink is down and has to replay them, in order, once the link is back. An AI agent writes down every action it takes so a bad run can be replayed later. All three are the same data structure: an append-only sequence of records, plus readers that remember where they stopped. Kafka is that structure, and also a cluster, partitions and the people who tune them.
The idea worth building keeps Kafka’s three useful ideas and drops the broker: an append-only log, consumers that remember their offset, and compaction by key, as a library over a directory of files. One writer process, many readers, no network protocol. The scope is one machine, and that limit is the product.
The Closest Precedents Are Servers
Lighter brokers exist. NATS JetStream adds persistent streams to the NATS server, and Redpanda is a Kafka-compatible broker that’s lighter to run than Kafka. Both are still processes you deploy, secure and monitor. Event-sourcing databases exist too: KurrentDB, formerly EventStoreDB, is built around domain events and projections, and it’s a server as well. The nearest embedded precedent is Chronicle Queue, a memory-mapped queue for Java on a single machine. What’s missing is an embeddable library with a documented file format and a CLI that can read what it wrote.
The honest competitor is a SQLite table with an integer primary key. It handles more cases than people admit: append, read from id N, keep a cursor per consumer. Where it falls short is housekeeping, because dropping a month of old events means a big DELETE and then a VACUUM to give the space back, while a segment-file log unlinks a file. If you don’t need that, use the table. Whether events suit your design at all is a separate decision, the one the event-driven architecture post is about.
A Directory of Segments
Each segment is a plain file named by its first offset, so finding the right one is a directory listing and a binary search.
events/
00000000000000000000.log closed segment, named by its first offset
00000000000000000000.idx sparse index: (offset, byte position)
00000000000004182917.log active segment, the only one being appended to
consumers.json committed offsets, replaced by write, fsync, rename
LOCK exclusive lock held by the single writer
record: length u32 | crc32c u32 | offset u64 | ts_ms i64 | schema u16 | key_len u16 | key | value
Inside a segment, records are length-prefixed and checksummed with CRC32C. The index is sparse, one entry every few kilobytes, so a lookup jumps near the target and scans a few records. It’s also derived data: delete it and the library rebuilds it by scanning the segment, so the log is the only thing you ever have to trust. Every record stores its own offset instead of deriving it from its position, because compaction leaves holes. After compaction, offsets 5, 9 and 12 can sit side by side, and a reader asking for 6 must be handed 9.
Appends go through one writer with group commit: records pile up, one fsync covers the batch, and callers are acknowledged afterwards. The writer then publishes a durable_offset, and readers only see records up to it. That’s the single-machine version of Kafka’s high watermark, and skipping it is a quiet disaster. Without it a consumer can read a record still sitting in the page cache, commit its offset, and then the power goes. After restart the log ends before the consumer’s cursor, new events reuse the old offsets, and the consumer skips them without an error.
Recovery truncates at the first bad record. On open, the library scans the active segment, checks each CRC and cuts at the first incomplete or corrupt record, even when later records look intact. The OS doesn’t promise to write pages back in order, so a crash can leave a hole of zeros followed by valid-looking data. Only fsynced records were ever acknowledged, so cutting at the hole loses nothing a producer was promised.
log = Log.open("events/", writer=True) # a second writer gets an error
off = log.append(key=b"doc:42", value=payload)
log.flush() # one fsync; durable_offset catches up
for rec in Log.open("events/").read(from_offset=cursor + 1, follow=True):
with db: # the consumer's own SQLite transaction
apply(rec, db)
db.execute("UPDATE cursor SET off = ?", (rec.offset,))
Retention by time or size removes whole closed segments, one unlink each. Compaction rewrites closed segments to keep the newest record per key, and an empty value works as a tombstone that stays around for a grace period so slow readers still see the delete.
Exactly-Once Is a Transaction Problem
The log delivers at least once, because a consumer can crash after doing the work and before saving its cursor. Exactly-once effects come from making the cursor update and the work one atomic step, and that only works when both live in the same transactional store. A consumer that writes to its own SQLite database keeps its cursor in the same transaction, as in the sketch above, and consumers.json is the fallback for consumers with no database. A consumer that calls an external API can’t do that. Derive an idempotency key from the log name and offset and let the other side deduplicate.
The producer side has the mirror problem. The log is a separate file, so an append can’t join your application’s database transaction. The answer is an outbox table in that database, drained into the log by a relay that tolerates duplicates. A sync outbox for a local-first app is this log with a cursor per peer, and what makes that hard is conflict rules, which the database sync post covers.
One writer is a deliberate limit, enforced by an exclusive lock on the LOCK file, and a second writer gets an error at open. Two processes appending to the same file have to agree on offset assignment, torn writes and index updates for every record. The only ways out are a lock around each append or routing appends through one process, and the second is networking by another name. Advisory file locks can also misbehave on network filesystems, so the rule is the same as for SQLite: local disk only. Many readers are fine, since they only read the durable prefix.
Compaction while readers are mid-segment is the next trap. The new segment is written beside the old one and renamed into place. A reader’s cursor is an offset, never a byte position, so a swapped segment can’t confuse it; it reopens by name and asks for the first record at or above its next offset. On Linux a replaced file stays readable by anyone who already has it open, so deleting the old segment is safe. Windows is stricter about deleting open files, so a library that wants to run there has to refcount open segments.
Events outlive the code that wrote them. The header carries a schema id so a reader two years from now can tell what it’s holding, and payloads should be a format with named or numbered fields, like JSON or protobuf, rather than a struct dump. Compaction won’t rescue a log full of unreadable payloads, because it only drops records and never rewrites them.
Skip mmap in 0.1. Plain write and pread through the page cache are a safer start and avoid a nasty failure: a mapped file that shrinks while a reader touches the missing pages kills the reader with SIGBUS instead of returning an error you can handle.
What Version 0.1 Refuses
Version 0.1 is a library for one language: single writer, many readers, named consumers with committed offsets, time and size retention, key compaction, and a CLI to tail the log, list segments and verify checksums. The CLI earns its place the first time someone asks what the agent did at 3 a.m. Agents are the newest customer: a recorder that belongs in one portable file is an event log with its format agreed in advance.
The refusals come next. No replication and no partitions across machines, no network protocol, and no consumer groups with rebalancing, because the problem they solve, many machines sharing a topic, doesn’t exist on one. No competing consumers with retries either: a queue forgets what it delivered and a log remembers, so if each message should be handled once by one of several workers, you want the job queue. And no queries beyond offset, key and time. Asking what a record looked like last Tuesday is a job for a change-history database, with the log as its input at best.
Most of what’s useful about Kafka fits in a directory.