Distributed Message Broker
Pure Java distributed commit log broker with binary wire framing.
A lightweight, event-driven message queue broker in pure Java built from first principles using raw TCP ServerSockets, an append-only commit log engine, and custom binary wire framing.
Industrial-scale distributed event streaming brokers like Apache Kafka are massive distributed systems with hundreds of thousands of lines of code and extensive external dependencies. Understanding how log-structured storage engines, binary wire framing protocols, and offset indices function under the hood requires stripping away distributed coordination layers and reconstructing the core primitives from first principles.
The architectural requirements were strict:
- 01.Zero external web frameworks or third-party serialization libraries—relying purely on standard Java runtime libraries.
- 02.High-throughput sequential disk writes that bypass random-access bottlenecks.
- 03.Fast offset lookup without loading whole commit logs into main memory.
- 04.Reliable crash recovery that can rebuild partition watermarks after unexpected server termination.
I designed the broker around three foundational distributed systems principles:
- 01.Append-Only Commit Log Storage Engine: Message frames are serialized into packed binary arrays with 4-byte magic bytes, CRC32 checksums, timestamps, and payload byte arrays. These frames are appended sequentially to disk using
FileChannel.write(), ensuring O(1) disk writes that fully saturate drive write heads and benefit from OS page caching. - 02.Binary Offset Indexing: Instead of linear log scanning, the broker maintains companion
.indexbinary files storing 8-byte relative offset to physical file position mappings. Lookups utilize binary search (O(log N)) to locate the nearest floor offset before streaming directly from disk. - 03.Custom Binary Wire Framing Protocol: Communication runs over persistent TCP sockets using a length-prefixed binary framing layout:
The implementation comprises several tightly coupled systems:
- ―Socket Handler & Connection Pool: A multi-threaded worker pool accepting client connections, parsing raw byte streams, and enforcing framing boundaries without memory fragmentation.
- ―Thread-Safe Producer & Consumer APIs: Client libraries implementing partition hashing, payload serialization, consumer group offset commits (
__consumer_offsets.dat), and auto-rebalance mechanics. - ―Log Segment Rotation & Watermarks: Configurable segment file rolling (e.g. at 100MB thresholds) maintaining high-watermark pointers to track committed message offsets.
- ―Embedded Architecture Dashboard: A lightweight HTTP server running on a dedicated administrative port that renders real-time partition visualizers and queue telemetry with zero external UI dependencies.
Building this broker exposed the mechanics of operating system I/O:
- ―Sequential disk writes on modern solid-state and magnetic media perform orders of magnitude faster than random writes because kernel read-ahead buffers and write-back caches work cooperatively.
- ―Allocating direct byte buffers (
ByteBuffer.allocateDirect) prevents JVM garbage collector pauses during sustained high-throughput socket ingestion. - ―Offset index sparsity is a vital tradeoff: indexing every N-th message instead of every single message drastically reduces index file size with negligible scan overhead.
PlantIQ — Coffee Agronomy & Advisory Platform
Multimodal precision agronomy platform with hybrid RAG and vernacular pre-routing.