Build a wide-column store (Cassandra / DynamoDB family)
13 scenes · ~91 min · build the primitive

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
  1. 01
  2. 02
  3. 03
  4. 04
  5. 05
  6. 06
  7. 07
  8. 08
  9. 09
  10. 10
  11. 11
  12. 11a
  13. 12

Why distribute?

One server has three ways to die — and bigger boxes never fix all three.

  1. 01
    One server, three ways to die — the single-node ceiling on capacity and throughput
    A single box hits a capacity wall, a throughput wall, and a death event — and your data is gone.
    ~7 min

Sharding

Hash mod N spreads load — until adding a server reshuffles 80% of keys.

  1. 02
    Split keys with hash mod N
    Hash the key, take it mod N, route to that server — keys spread evenly across the cluster.
    ~7 min
  2. 03
    Adding a server remaps everything — hash-mod-N and the cluster-wide rebalance
    When N changes from 4 to 5, almost every key's home changes — a cluster-wide migration.
    ~7 min

The ring

Consistent hashing + vnodes — adding a server moves only ~1/N of keys.

  1. 04
    The ring — keys and nodes on a circle — arc ownership and the clockwise neighbor
    Map both keys and servers onto a circle; each key belongs to the next server clockwise.
    ~7 min
  2. 05
    Vnodes flatten the lumpy ring
    Give 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.

  1. 06
    Copies on the next N servers — replication factor RF across clockwise replicas
    Store each key on the next RF distinct servers clockwise so a node death loses no data.
    ~7 min
  2. 07
    Replicas disagree, then converge
    Concurrent 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.

  1. 08
    W plus R greater than N — quorum overlap and read-your-writes
    Make every read overlap every write on at least one replica by sliding two knobs.
    ~7 min
  2. 09
    Partition forces a choice
    Split the cluster in two; one side keeps quorum, the other goes unavailable — or you accept divergence.
    ~7 min

Healing

Hinted handoff + read repair + anti-entropy + gossip keep the cluster honest.

  1. 10
    Hold the write until it wakes up — hinted handoff and the hint TTL
    When a replica is briefly unreachable, the coordinator stashes the write and replays it on return.
    ~7 min
  2. 11
    Heal on read, heal on a schedule — read repair and Merkle-tree anti-entropy
    When a read sees disagreement, push the winner to stale replicas; on a cron, compare every key.
    ~7 min
  3. 11a
    Gossip keeps everyone informed
    Each 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.

  1. 12
    Design canvas — pick every knob
    Given a real workload, pick partition key, RF, W/R, vnodes, and DC topology — and grade the result.
    ~7 min

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.

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