Distributed Stream Processor
A stream-processing system built from scratch in Go: SWIM failure detection, a replicated file system with consistent hashing, and exactly-once streaming on a 10-node Docker cluster.
Origin
Started in my graduate distributed systems course (CS 425). Afterward, I moved the cluster to Docker so the demos could run locally.
The problem
Stream-processing frameworks hide the hard parts: how nodes learn that a peer died, where replicas live, what exactly-once actually costs you. I had to build the whole stack from the ground up.
This was the cumulative project for my graduate distributed systems course. It ran on a 10-VM cluster there. Afterward I moved it to a 10-node Docker cluster so the demos could run locally with one command.
How it works
Three layers, each built on the one below. Membership uses SWIM-style failure detection with gossip dissemination, where suspicion and confirm windows give a six-second detection bound and no leader sits in the detection path.
Storage is a replicated file system on a consistent-hashing ring, three replicas per file, with pipelined replication on writes. On top of both sits a leader-worker streaming engine with exactly-once semantics: a failed task restarts on another node and produces identical output with zero duplicates.
What the measurements taught me
Kill a node and re-replication finishes in about a second, whether the lost files are small or large. That flatness is the interesting part: fixed coordination overhead dominates, not data transfer, and you only learn that by measuring.
On these workloads, it was slower than Spark but had more consistent runtimes. The full measurements live in the repo.