Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
69 changes: 62 additions & 7 deletions README.md
Original file line number Diff line number Diff line change
@@ -1,9 +1,17 @@
# Distributed Message Queue

A Kafka-style message queue built in Go, with exactly-once delivery via 2PC. Built to understand what Kafka trades away for throughput — and what it costs to get delivery guarantees back.
[![CI](https://github.com/ga11221/message-queue/actions/workflows/ci.yml/badge.svg)](https://github.com/ga11221/message-queue/actions/workflows/ci.yml)

A Kafka-style message queue built in Go. Phase 1 (single-node log engine + broker + TCP server) is complete. Phase 2 (exactly-once via 2PC) is in progress — design complete, implementation next.

<!-- REMINDER: Every section should explain WHY, not just WHAT. Show design decisions and tradeoffs. A tech HM skimming this should see depth of understanding, not a feature list. -->

![demo: produce with key routing, merged consume, on-disk validation](docs/demo.gif)

*The full lifecycle: keyed produce (watch FNV-1a partition
assignment), consume merged across partitions, and byte-level on-disk
validation. Run it yourself: `./scripts/demo.sh` — recorded 2026-08-26.*

## Documentation

| Document | What it covers |
Expand Down Expand Up @@ -34,19 +42,28 @@ Detailed ADRs with commit links: [docs/decisions/](docs/decisions/).

**Exactly-once via 2PC** — Kafka chose at-least-once + idempotent producers as a practical compromise (20-30% throughput cost). This project implements full 2PC to show what the real version costs: 10-100x throughput reduction, 2-5x latency increase, blocking on coordinator failure.

**Single writer per partition** — Avoids lock-free complexity (gap problems on crash, requires `pwrite()`). Parallelism comes from multiple partitions, not splitting one.
**Per-partition locking** — Each partition has its own mutex; parallelism comes from multiple partitions, not splitting one. Measured contention: ~12% under 4-goroutine parallel writes (PROGRESS.md) — negligible compared to the 86% syscall cost on the append path.

## Architecture

```
Producer --> Topic.Produce() --> partition.Append() --> active segment
|
v
Consumer <-------- partition.Read() <-------- OS page cache
Producer --> Server (goroutine per conn, synchronous)
|
v
Broker --> Topic --> Partition x N
|
v
append-only segment files
(.log + .index + .timeindex)
|
v
Consumer <-------- [read from offset] <----- OS page cache
```

Each layer has its own lock (see [Concurrency Model](docs/phase1/concurrency.md)).
Parallelism comes from multiple partitions, not from pipelining within one.
See [architecture overview](docs/phase1/architecture.md) for the full
diagram with layer responsibilities and invariants.

## What Kafka Gets Right (and What It Costs)

Expand All @@ -64,7 +81,7 @@ Parallelism comes from multiple partitions, not from pipelining within one.
| Phase | Scope | Status |
|-------|-------|--------|
| 1 | Single-node: log engine + broker + server/RPC | Complete |
| 2 | Exactly-once via 2PC (single-node) | Not started |
| 2 | Exactly-once via 2PC (single-node) | Design complete, implementation next |
| 3 | Replication (ISR, leader election, replica fetch) | Not started |
| 4 | Stream processing (conditional) | Not started |

Expand Down Expand Up @@ -96,6 +113,44 @@ docker compose up --build -d # server on :9092, data in mq-data volume
docker compose down # stop (volume persists)
```

### Demo: 60-second walkthrough

`scripts/demo.sh` runs the full lifecycle against the container — build the
CLI client, create a topic with 3 partitions, produce keyed messages (watch
the broker's FNV-1a partition assignment), consume everything merged across
partitions, and verify the bytes at rest with `mqvalidate`:

```bash
./scripts/demo.sh
```

Sample output:

```
=== produce keyed messages (note partition assignments) ===
produced to demo-events-1787715727 partition=1 offset=0
produced to demo-events-1787715727 partition=1 offset=1
...
-- 4 message(s) across 3 partition(s), merged by timestamp

4 topics, 12 segments, 6 messages, 246 bytes, 0 problems
on-disk data: OK
```

The same client works by hand:

```bash
go run ./cmd/mqclient produce -addr localhost:9092 -topic t -key user-1 -value hello
go run ./cmd/mqclient consume -addr localhost:9092 -topic t -all
```

Ordering guarantee (matching Kafka): messages are ordered per partition — and
per key when producing with keys, since a key always routes to the same
partition. There is **no global ordering across partitions**; `consume --all`
merges by timestamp on a best-effort basis. The raw asciicast of the demo is
available ([docs/demo.cast](docs/demo.cast)) for terminal players
(`asciinema play docs/demo.cast`).

### Integration test against the container

`TestExternalServer` runs a create/produce/consume round-trip against a live
Expand Down
83 changes: 83 additions & 0 deletions bench-results/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,83 @@
# Benchmark Results

Raw benchstat output and CSV exports for the message-queue log engine.
Hardware and environment are recorded per run (see `.env` files).

## Files

| File | Description |
|------|-------------|
| `baseline.env` | Environment snapshot for the baseline run |
| `baseline.txt` | Raw benchstat output (10 iterations, 2s benchtime) |
| `baseline.csv` | Same data as CSV — importable to spreadsheet tools |
| `sweep.env` | Environment snapshot for the batch-size sweep |
| `sweep.txt` | Raw benchstat output (5 iterations, sweep over batch sizes) |
| `sweep.csv` | Same data as CSV |

## Baseline Summary

Collected 2026-08-26 on bare metal (ext4 scratch disk), 0 Docker containers
running, Go 1.26.0, AMD Ryzen 7 7735HS x4.

| Benchmark | ns/op | MB/s | B/op | allocs |
|-----------|------:|-----:|-----:|-------:|
| Encode 16B | 76 | 657 | 96 | 2 |
| Encode 256B | 231 | 1253 | 608 | 2 |
| Encode 1024B | 694 | 1536 | 2304 | 2 |
| Decode 16B | 108 | 461 | 120 | 4 |
| Decode 256B | 277 | 1048 | 616 | 4 |
| Decode 1024B | 717 | 1511 | 2248 | 4 |
| Segment.Append | 3858 | 34.8 | 272 | 2 |
| Partition.Append | 3882 | 34.6 | 272 | 2 |
| PartitionRead 10K | 11308 | — | 433 | 5 |
| PartitionRead 1M | 16309 | — | 433 | 5 |
| IndexLookup 10 | 3351 | — | 0 | 0 |
| IndexLookup 10K | 9544 | — | 0 | 0 |
| IndexLookup 1M | 14556 | — | 0 | 0 |
| ConcurrentProduce | 4313 | 31.1 | 272 | 2 |

**Key observations:**
- Encode scales linearly with value size (76ns for 16B, 694ns for 1024B)
- Segment.Append and Partition.Append are identical — the mutex adds ~0ns under
single-writer load (contention matters only with parallel goroutines)
- Index lookup is disk-bound: 3.3us for 10 entries, 14.6us for 1M entries
(binary search probes ReadAt per step)
- ConcurrentProduce (4 goroutines) is ~11% slower than single-writer
Partition.Append — the measured contention from the profiler

## Batch-Size Sweep

The append path was measured 86% syscall-bound in CPU profiles (see
PROGRESS.md). This sweep finds the throughput knee for batching N records
into one write syscall.

| Batch size | ns/op total | ns/op per record | MB/s |
|------------|------------:|------------------:|-----:|
| 1 | 1294 | 1294 | 103 |
| 4 | 1696 | 424 | 316 |
| 16 | 3323 | 208 | 645 |
| 64 | 9238 | 144 | 928 |
| 256 | 32463 | 127 | 1057 |

**Knee at 16-64 records** (~2-8KB per write). Beyond 64, per-record
improvement flattens (144 -> 127 ns/op) while total latency climbs 3.5x.
Target for Phase 2.5 buffered writes: 32-128 records per flush.

## Reproducing

```bash
# Generate env header
scripts/bench-env.sh > bench-results/my-run.env

# Run baseline (10 iterations, 2s each)
go test ./pkg/log/ -bench=. -benchmem -count=10 -benchtime=2s > bench-results/my-run.txt

# Append env header
cat bench-results/my-run.env bench-results/my-run.txt > bench-results/my-run-combined.txt

# Run batch sweep (5 iterations)
go test ./pkg/log/ -bench=BenchmarkSegmentAppendBatch -benchmem -count=5 -benchtime=2s > bench-results/sweep.txt

# Generate CSV from benchstat output
# (use the parse script or manual extraction)
```
141 changes: 141 additions & 0 deletions bench-results/baseline.csv
Original file line number Diff line number Diff line change
@@ -0,0 +1,141 @@
benchmark,iterations,ns_per_op,throughput_mbps,b_per_op,allocs_per_op
BenchmarkEncodeMessage/value_16B,29983010,75.75,660.06,96,2
BenchmarkEncodeMessage/value_16B,29762554,75.45,662.72,96,2
BenchmarkEncodeMessage/value_16B,28896696,86.6,577.37,96,2
BenchmarkEncodeMessage/value_16B,31379400,76.81,650.99,96,2
BenchmarkEncodeMessage/value_16B,31298907,77.42,645.85,96,2
BenchmarkEncodeMessage/value_16B,31445527,77.59,644.44,96,2
BenchmarkEncodeMessage/value_16B,31776916,76.97,649.6,96,2
BenchmarkEncodeMessage/value_16B,32823192,73.08,684.19,96,2
BenchmarkEncodeMessage/value_16B,31958130,73.4,681.18,96,2
BenchmarkEncodeMessage/value_16B,32803478,74.32,672.76,96,2
BenchmarkEncodeMessage/value_256B,9934572,230.3,1259.12,608,2
BenchmarkEncodeMessage/value_256B,10919041,229.4,1263.99,608,2
BenchmarkEncodeMessage/value_256B,10864833,232.6,1246.76,608,2
BenchmarkEncodeMessage/value_256B,9986918,231.5,1252.68,608,2
BenchmarkEncodeMessage/value_256B,10444516,233.3,1243.02,608,2
BenchmarkEncodeMessage/value_256B,10550792,231.7,1251.49,608,2
BenchmarkEncodeMessage/value_256B,9700810,234.4,1237.01,608,2
BenchmarkEncodeMessage/value_256B,10922130,231.7,1251.83,608,2
BenchmarkEncodeMessage/value_256B,9909126,231.2,1254.34,608,2
BenchmarkEncodeMessage/value_256B,10543086,227.7,1273.38,608,2
BenchmarkEncodeMessage/value_1024B,3480708,702.3,1506.51,2304,2
BenchmarkEncodeMessage/value_1024B,3208869,709.4,1491.32,2304,2
BenchmarkEncodeMessage/value_1024B,3335572,716.8,1475.95,2304,2
BenchmarkEncodeMessage/value_1024B,3397845,735.8,1437.92,2304,2
BenchmarkEncodeMessage/value_1024B,3262086,730.7,1448.01,2304,2
BenchmarkEncodeMessage/value_1024B,3225657,650.0,1627.59,2304,2
BenchmarkEncodeMessage/value_1024B,3575229,676.1,1564.88,2304,2
BenchmarkEncodeMessage/value_1024B,3836512,673.8,1570.15,2304,2
BenchmarkEncodeMessage/value_1024B,3431493,665.3,1590.19,2304,2
BenchmarkEncodeMessage/value_1024B,3567741,678.3,1559.69,2304,2
BenchmarkDecodeMessage/value_16B,22345396,107.1,466.8,120,4
BenchmarkDecodeMessage/value_16B,22194754,107.7,464.26,120,4
BenchmarkDecodeMessage/value_16B,22761320,107.7,464.09,120,4
BenchmarkDecodeMessage/value_16B,21936830,108.5,460.72,120,4
BenchmarkDecodeMessage/value_16B,21380168,108.9,459.01,120,4
BenchmarkDecodeMessage/value_16B,22578217,108.5,460.81,120,4
BenchmarkDecodeMessage/value_16B,22133130,108.6,460.5,120,4
BenchmarkDecodeMessage/value_16B,21951352,108.3,461.74,120,4
BenchmarkDecodeMessage/value_16B,21495777,108.9,459.22,120,4
BenchmarkDecodeMessage/value_16B,21497160,109.3,457.36,120,4
BenchmarkDecodeMessage/value_256B,8837772,275.0,1054.65,616,4
BenchmarkDecodeMessage/value_256B,8706324,274.2,1057.51,616,4
BenchmarkDecodeMessage/value_256B,8850172,264.6,1096.0,616,4
BenchmarkDecodeMessage/value_256B,8737519,273.4,1060.57,616,4
BenchmarkDecodeMessage/value_256B,8662170,272.3,1064.96,616,4
BenchmarkDecodeMessage/value_256B,8125345,277.5,1045.0,616,4
BenchmarkDecodeMessage/value_256B,8597256,280.9,1032.56,616,4
BenchmarkDecodeMessage/value_256B,8679314,289.7,1000.96,616,4
BenchmarkDecodeMessage/value_256B,8416669,279.8,1036.52,616,4
BenchmarkDecodeMessage/value_256B,8389489,282.9,1025.17,616,4
BenchmarkDecodeMessage/value_1024B,3107478,751.2,1408.5,2248,4
BenchmarkDecodeMessage/value_1024B,3115168,751.1,1408.63,2248,4
BenchmarkDecodeMessage/value_1024B,3122564,748.9,1412.8,2248,4
BenchmarkDecodeMessage/value_1024B,3098383,770.7,1372.85,2248,4
BenchmarkDecodeMessage/value_1024B,3396823,665.6,1589.64,2248,4
BenchmarkDecodeMessage/value_1024B,3508438,668.0,1583.8,2248,4
BenchmarkDecodeMessage/value_1024B,3681781,660.7,1601.23,2248,4
BenchmarkDecodeMessage/value_1024B,3669466,683.9,1546.95,2248,4
BenchmarkDecodeMessage/value_1024B,3564024,677.5,1561.59,2248,4
BenchmarkDecodeMessage/value_1024B,3605251,688.0,1537.75,2248,4
BenchmarkSegmentAppend,522537,3875.0,34.58,272,2
BenchmarkSegmentAppend,540696,3882.0,34.52,272,2
BenchmarkSegmentAppend,534979,3875.0,34.58,272,2
BenchmarkSegmentAppend,531940,3861.0,34.71,272,2
BenchmarkSegmentAppend,537823,3850.0,34.8,272,2
BenchmarkSegmentAppend,534907,3849.0,34.82,272,2
BenchmarkSegmentAppend,541276,3843.0,34.87,272,2
BenchmarkSegmentAppend,522404,3854.0,34.77,272,2
BenchmarkSegmentAppend,533866,3849.0,34.82,272,2
BenchmarkSegmentAppend,538292,3857.0,34.74,272,2
BenchmarkPartitionAppend,537637,3866.0,34.67,272,2
BenchmarkPartitionAppend,526990,3856.0,34.75,272,2
BenchmarkPartitionAppend,534840,3868.0,34.65,272,2
BenchmarkPartitionAppend,526450,3879.0,34.55,272,2
BenchmarkPartitionAppend,530046,3883.0,34.51,272,2
BenchmarkPartitionAppend,543667,3872.0,34.61,272,2
BenchmarkPartitionAppend,534903,3879.0,34.55,272,2
BenchmarkPartitionAppend,530132,3947.0,33.95,272,2
BenchmarkPartitionAppend,516913,3885.0,34.49,272,2
BenchmarkPartitionAppend,527968,3886.0,34.48,272,2
BenchmarkConcurrentProduce,485064,4301.0,31.15,272,2
BenchmarkConcurrentProduce,477423,4300.0,31.16,272,2
BenchmarkConcurrentProduce,482035,4302.0,31.14,272,2
BenchmarkConcurrentProduce,480660,4333.0,30.93,272,2
BenchmarkConcurrentProduce,495822,4291.0,31.23,272,2
BenchmarkConcurrentProduce,489723,4302.0,31.15,272,2
BenchmarkConcurrentProduce,476989,4322.0,31.0,272,2
BenchmarkConcurrentProduce,485784,4339.0,30.88,272,2
BenchmarkConcurrentProduce,475093,4350.0,30.8,272,2
BenchmarkConcurrentProduce,481269,4289.0,31.24,272,2
BenchmarkPartitionRead/depth_10000,207650,11348.0,0,433,5
BenchmarkPartitionRead/depth_10000,209059,11293.0,0,433,5
BenchmarkPartitionRead/depth_10000,212390,11302.0,0,433,5
BenchmarkPartitionRead/depth_10000,210716,11309.0,0,433,5
BenchmarkPartitionRead/depth_10000,211742,11300.0,0,433,5
BenchmarkPartitionRead/depth_10000,210752,11287.0,0,433,5
BenchmarkPartitionRead/depth_10000,212746,11270.0,0,433,5
BenchmarkPartitionRead/depth_10000,211465,11249.0,0,433,5
BenchmarkPartitionRead/depth_10000,211916,11466.0,0,433,5
BenchmarkPartitionRead/depth_10000,212266,11257.0,0,433,5
BenchmarkPartitionRead/depth_1000000,145368,16332.0,0,433,5
BenchmarkPartitionRead/depth_1000000,147697,16246.0,0,433,5
BenchmarkPartitionRead/depth_1000000,147794,16340.0,0,433,5
BenchmarkPartitionRead/depth_1000000,147321,16306.0,0,433,5
BenchmarkPartitionRead/depth_1000000,146128,16267.0,0,433,5
BenchmarkPartitionRead/depth_1000000,147499,16323.0,0,433,5
BenchmarkPartitionRead/depth_1000000,146624,16348.0,0,433,5
BenchmarkPartitionRead/depth_1000000,147133,16311.0,0,433,5
BenchmarkPartitionRead/depth_1000000,146814,16254.0,0,433,5
BenchmarkPartitionRead/depth_1000000,147158,16363.0,0,433,5
BenchmarkIndexLookup/entries_10,670300,3352.0,0,0,0
BenchmarkIndexLookup/entries_10,679184,3358.0,0,0,0
BenchmarkIndexLookup/entries_10,684700,3382.0,0,0,0
BenchmarkIndexLookup/entries_10,680624,3345.0,0,0,0
BenchmarkIndexLookup/entries_10,676429,3343.0,0,0,0
BenchmarkIndexLookup/entries_10,675445,3340.0,0,0,0
BenchmarkIndexLookup/entries_10,681417,3351.0,0,0,0
BenchmarkIndexLookup/entries_10,676490,3332.0,0,0,0
BenchmarkIndexLookup/entries_10,681511,3356.0,0,0,0
BenchmarkIndexLookup/entries_10,671923,3350.0,0,0,0
BenchmarkIndexLookup/entries_10000,247064,9518.0,0,0,0
BenchmarkIndexLookup/entries_10000,246140,9515.0,0,0,0
BenchmarkIndexLookup/entries_10000,243652,9522.0,0,0,0
BenchmarkIndexLookup/entries_10000,246490,9562.0,0,0,0
BenchmarkIndexLookup/entries_10000,246704,9534.0,0,0,0
BenchmarkIndexLookup/entries_10000,248563,9558.0,0,0,0
BenchmarkIndexLookup/entries_10000,245401,9542.0,0,0,0
BenchmarkIndexLookup/entries_10000,249040,9589.0,0,0,0
BenchmarkIndexLookup/entries_10000,248149,9523.0,0,0,0
BenchmarkIndexLookup/entries_10000,249184,9547.0,0,0,0
BenchmarkIndexLookup/entries_1000000,153297,14529.0,0,0,0
BenchmarkIndexLookup/entries_1000000,155474,14598.0,0,0,0
BenchmarkIndexLookup/entries_1000000,153115,14553.0,0,0,0
BenchmarkIndexLookup/entries_1000000,152943,14652.0,0,0,0
BenchmarkIndexLookup/entries_1000000,151280,14597.0,0,0,0
BenchmarkIndexLookup/entries_1000000,151652,14638.0,0,0,0
BenchmarkIndexLookup/entries_1000000,154797,14496.0,0,0,0
BenchmarkIndexLookup/entries_1000000,156000,14541.0,0,0,0
BenchmarkIndexLookup/entries_1000000,154518,14508.0,0,0,0
BenchmarkIndexLookup/entries_1000000,165164,14642.0,0,0,0
8 changes: 8 additions & 0 deletions bench-results/baseline.txt
Original file line number Diff line number Diff line change
@@ -1,3 +1,11 @@
date: 2026-08-26T02:09:15Z
go: go version go1.26.0 linux/amd64
kernel: 6.6.114.1-microsoft-standard-WSL2
cpu: AMD Ryzen 7 7735HS with Radeon Graphics x 4 cores
disk: E:\ (WSL2 virtualized ext4)
loadavg: 0.86 0.67 0.42
docker: 0 container(s) running (quiesce before measuring!)

goos: linux
goarch: amd64
pkg: message-queue/pkg/log
Expand Down
26 changes: 26 additions & 0 deletions bench-results/sweep.csv
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
benchmark,iterations,ns_per_op,throughput_mbps,b_per_op,allocs_per_op
BenchmarkSegmentAppendBatch/batch_1,1861215,1286.0,104.22,0,0
BenchmarkSegmentAppendBatch/batch_1,1850292,1296.0,103.41,0,0
BenchmarkSegmentAppendBatch/batch_1,1847721,1297.0,103.33,0,0
BenchmarkSegmentAppendBatch/batch_1,1855597,1294.0,103.55,0,0
BenchmarkSegmentAppendBatch/batch_1,1858166,1297.0,103.34,0,0
BenchmarkSegmentAppendBatch/batch_4,1415066,1692.0,316.71,0,0
BenchmarkSegmentAppendBatch/batch_4,1415581,1695.0,316.18,0,0
BenchmarkSegmentAppendBatch/batch_4,1406553,1699.0,315.44,0,0
BenchmarkSegmentAppendBatch/batch_4,1411066,1696.0,316.03,0,0
BenchmarkSegmentAppendBatch/batch_4,1410406,1697.0,315.82,0,0
BenchmarkSegmentAppendBatch/batch_16,617504,3322.0,645.44,0,0
BenchmarkSegmentAppendBatch/batch_16,641960,3317.0,646.43,0,0
BenchmarkSegmentAppendBatch/batch_16,714495,3352.0,639.64,0,0
BenchmarkSegmentAppendBatch/batch_16,723044,3315.0,646.67,0,0
BenchmarkSegmentAppendBatch/batch_16,723578,3310.0,647.78,0,0
BenchmarkSegmentAppendBatch/batch_64,239803,9265.0,925.68,0,0
BenchmarkSegmentAppendBatch/batch_64,245210,9221.0,930.08,0,0
BenchmarkSegmentAppendBatch/batch_64,243982,9203.0,931.9,0,0
BenchmarkSegmentAppendBatch/batch_64,239110,9233.0,928.82,0,0
BenchmarkSegmentAppendBatch/batch_64,247848,9268.0,925.29,0,0
BenchmarkSegmentAppendBatch/batch_256,72699,32481.0,1056.13,0,0
BenchmarkSegmentAppendBatch/batch_256,73136,32783.0,1046.4,0,0
BenchmarkSegmentAppendBatch/batch_256,72991,32297.0,1062.14,0,0
BenchmarkSegmentAppendBatch/batch_256,72529,32306.0,1061.86,0,0
BenchmarkSegmentAppendBatch/batch_256,73081,32449.0,1057.18,0,0
8 changes: 8 additions & 0 deletions bench-results/sweep.txt
Original file line number Diff line number Diff line change
@@ -1,3 +1,11 @@
date: 2026-08-26T02:22:18Z
go: go version go1.26.0 linux/amd64
kernel: 6.6.114.1-microsoft-standard-WSL2
cpu: AMD Ryzen 7 7735HS with Radeon Graphics x 4 cores
disk: E:\ (WSL2 virtualized ext4)
loadavg: 1.63 1.55 1.15
docker: 0 container(s) running (quiesce before measuring!)

goos: linux
goarch: amd64
pkg: message-queue/pkg/log
Expand Down
Loading
Loading