Build your own wide-column store (Cassandra / DynamoDB family)
One server is not enough — disk fills, throughput maxes, the box dies. Build a multi-node store from first principles: hash sharding, the consistent-hash ring, vnodes, replication, eventual consistency, tunable W+R quorum, hinted handoff, read repair. Every modern Dynamo-style store is a point in this design space.
- Scenes
- 13 interactive scenes
- Time
- about 91 minutes
- Topic
- Storage Engines & Databases
What you are building, and why
You have run a single-node KV or SQL database. Maybe you have watched the disk fill. Maybe you have watched throughput peg the ceiling. Maybe you have watched the box die and taken your data with it. The single server has three ways to die — and bigger disks, faster CPUs, and more RAM make those three problems quantitatively softer, but never qualitatively softer. To survive any of them, you need more than one machine. That is the brief.
The wide-column store is the family of distributed databases that took that brief seriously: Cassandra, DynamoDB, ScyllaDB, Riak. They share a load-bearing design — Dynamo-style hash partitioning, tunable consistency, anti-entropy — that is so widely deployed it is the default mental model for "database that survives a node death." Once you have built one from first principles, every later distributed database collapses to a point in this design space.
Resist the urge to memorize Cassandra's command-line flags. Build the system from the failure that motivates each step. The first server fails three ways; the cheapest fix shards by key; sharding by hash mod N reshuffles 80% of keys when you add a server; the ring fixes it but creates uneven arcs; vnodes fix the unevenness; copies survive a node death; copies that disagree need a winner; counting replicas (W+R>N) gives you reads that see writes; partitions force a choice; hinted handoff softens the everyday case; read repair and anti-entropy heal the rest; gossip tells the cluster who is alive. Twelve scenes; twelve choices; one design canvas at the end.
What you will be able to explain afterwards
- single-node ceiling — capacity, throughput, durability
- hash mod N sharding and its (N-1)/N rebalance catastrophe
- consistent hashing — ring + arcs
- virtual nodes (vnodes) — flatten load, spread death
- replication factor and successor replicas
- eventual consistency + last-write-wins (LWW)
- tunable consistency — W+R>N, ONE/QUORUM/ALL
- CAP under partition — per-request availability vs consistency
- hinted handoff — optimistic recovery + hint TTL
- read repair + anti-entropy via Merkle trees
- gossip — O(log N) cluster-state convergence
- design canvas — pick every knob against a workload
Why distribute?
One server has three ways to die — and bigger boxes never fix all three.
Sharding
Hash mod N spreads load — until adding a server reshuffles 80% of keys.
The ring
Consistent hashing + vnodes — adding a server moves only ~1/N of keys.
- 04The ring — keys and nodes on a circle — arc ownership and the clockwise neighborMap both keys and servers onto a circle; each key belongs to the next server clockwise.~7 min
- 05Vnodes flatten the lumpy ringGive each physical server many small ring positions; arcs become uniform and deaths spread their load.~7 min
Replication
Copies on the next RF servers; eventual consistency under concurrent writes.
- 06Copies on the next N servers — replication factor RF across clockwise replicasStore each key on the next RF distinct servers clockwise so a node death loses no data.~7 min
- 07Replicas disagree, then convergeConcurrent writes hit replicas at different times; for a moment they disagree, but they converge under LWW.~7 min
Tunable & CAP
W+R>N for strong reads; per-request choice between A and C under partition.
Healing
Hinted handoff + read repair + anti-entropy + gossip keep the cluster honest.
- 10Hold the write until it wakes up — hinted handoff and the hint TTLWhen a replica is briefly unreachable, the coordinator stashes the write and replays it on return.~7 min
- 11Heal on read, heal on a schedule — read repair and Merkle-tree anti-entropyWhen a read sees disagreement, push the winner to stale replicas; on a cron, compare every key.~7 min
- 11aGossip keeps everyone informedEach node, every second, swaps state with a few peers; the cluster picture converges without a master.~7 min
Design canvas
Pick every knob against a workload and grade the result.
More in Storage Engines & Databases
Open the box every design diagram labels "DB": pages, logs, LSM trees, wide-column, documents, graphs and columnar scans, built from scratch.
- Build Build a Bitcask-style KV storeThe simplest possible KV store that still works: an append-only log on disk + an in-memory hash index. Build it from first principles and feel which trade-offs every later store inherits.
- Build Build an LSM-tree storage engine (LevelDB / RocksDB style)The simplest possible storage engine that gives you BOTH ordered reads AND more keys than fit in RAM, by accepting a deal: write to RAM at memory speed, log to disk for safety, then merge sorted files in the background forever.
- Build Build a B-tree storage engine (SQLite-style)What actually happens when you run INSERT INTO users(...). One file of fixed-size pages, organized as B-trees, with a write-ahead log that turns commits into appends. Build it from a SQL writer's perspective and feel why every knob exists.
- Build Build a graph database (Neo4j / Dgraph-style)When the workload is 'friends of friends', a relational join melts. Build a store where edges are first-class — index-free adjacency, traversals that follow pointers instead of joining tables, and a query language (Cypher / GraphQL+) that thinks in patterns. Feel why graph storage shines for traversal-heavy work and stumbles on full-graph aggregates.
- Build Build a columnar OLAP store (ClickHouse / Druid style)OLTP picks one row by key; OLAP scans a billion rows of one column and asks for a percentile. Build the analytical engine that makes that fast: columnar layout, dictionary/RLE/delta compression, vectorized execution, late materialization, MPP shuffle. Internalize why Postgres is 1000× slower than ClickHouse on the same query and why the inverse is also true.
Prefer to design it yourself?
The same subject as a staged workspace: draw the architecture, and a simulator traces requests through the boxes you drew.
Open the Build a wide-column store (Cassandra / DynamoDB family) workspace