System design interview · backend

Design a distributed key-value store

"Build us Dynamo/Cassandra: a key-value store that holds 1 billion keys, survives servers dying weekly, and scales by just adding machines — no downtime, no full reshuffles."

⚡ The takeaway first

Put nodes and keys on the same hash ring, and let each key live with the first node clockwise from it. That's consistent hashing: adding or removing a server only moves the keys that belonged to that server's arc — about K/n of them — instead of rehashing everything. Add virtual nodes for even spreading, replicate each key to the next N nodes clockwise, and use quorums so reads and writes survive dead nodes.

📋 Requirements

Functional

Non-functional

🚫 Common misconception

"Adding a server means rehashing all the keys." Only with naive hash(key) % n — change n and nearly every key moves. Consistent hashing fixes exactly this: keys and nodes share one ring, and a new node only "steals" the arc between itself and its counterclockwise neighbor. With 100 nodes and 1B keys, adding node 101 moves roughly 1B/100 = 10M keys (1%), not a billion. The ring is the whole trick — try the widget below and watch the counter.

🧮 Back-of-the-envelope math

AssumptionValue
Keys1,000,000,000
Avg value size1 KB
Replication factor N3
Total stored1B × 1 KB × 3 ≈ 3 TB
Nodes100 → 30 GB per node — comfortable
Keys moved when adding 1 node≈ K/n = 10M keys ≈ 1% (×3 replicas = 30M key-copies)
Read quorum (N=3, R=2)2 of 3 replicas must agree — survives 1 dead node
Write quorum (N=3, W=2)R + W > N (2+2>3) → strong consistency on overlap

Per-node data is small; the interesting numbers are all about movement (K/n) and quorums (R+W>N). Interviewers will ask "what if two nodes die?" — answer: with N=3, W=2, a write needs 2 healthy replicas; lose 2 nodes holding the same key and that key's writes block until hinted handoff replays.

Go deeper: virtual nodes

A plain ring gives lumpy arcs — one node might own 3% of the ring, another 0.2%. Give each physical node ~100–200 virtual nodes (random points on the ring) and the law of large numbers evens the load out. Bonus: a beefy machine gets more vnodes than a small one — weighted capacity for heterogeneous hardware, free.

🏗️ Architecture

The centerpiece: no master, just a ring every node agrees on.

flowchart TD
    A[Client library
knows the ring] --> B{hash key} B --> C[Coordinator node
any node can coordinate] C --> D[Replica 1
primary: first clockwise] C --> E[Replica 2
next clockwise] C --> F[Replica 3
next clockwise] D --> G[(SSTables
per node)] E --> G F --> G H[Gossip protocol
membership + failure detection] -.-> C H -.-> D I[Hinted handoff
replay buffer] -.-> D

Think of it like a clock face: servers stand at various hour marks, and each key is a sticky note placed at its hash position — it belongs to the first server clockwise from it. A new server squeezes in between two marks and only collects the sticky notes in its little arc. Nobody else's notes move.

🔍 Component deep-dives

Write path — quorum + hinted handoff

sequenceDiagram
    participant C as Client
    participant N1 as Coordinator
    participant R1 as Replica A (primary)
    participant R2 as Replica B
    participant R3 as Replica C (down!)
    C->>N1: PUT(k, v)
    N1->>R1: write v1
    N1->>R2: write v1
    N1->>R3: write v1 (timeout)
    R1-->>N1: ack
    R2-->>N1: ack
    Note over N1: W=2 reached → success.
R3's write becomes a hint. N1-->>C: 200 OK Note over N1,R3: R3 recovers → hint replays.
No write was ever lost.

Read path — read repair

sequenceDiagram
    participant C as Client
    participant N1 as Coordinator
    participant R1 as Replica A
    participant R2 as Replica B
    C->>N1: GET(k)
    N1->>R1: read
    N1->>R2: read
    R1-->>N1: v5 (newer)
    R2-->>N1: v4 (stale)
    Note over N1: R=2 reached. Newest wins;
read repair pushes v5 to R2. N1-->>C: v5

Every read heals the system a little: stale replicas get fixed in the background, so entropy never accumulates. Vector clocks (or last-write-wins with timestamps, honestly) resolve the "which is newer" question.

🔌 API + data model

API

PUT /v1/keys/{key}   { "value": "..." }   → 200 (W acks)
GET /v1/keys/{key}                        → 200 { "value": "..." } (R acks)
-- admin --
POST /v1/nodes { "id": "n101" }            → node joins, streams K/n keys
DELETE /v1/nodes/{id}                     → data re-replicates, node leaves

Data model

⚖️ Trade-offs

DecisionOption AOption BPick
Consistency (N=3)R=1, W=3 (fast reads)R=2, W=2 (balanced)B default — R+W>N gives read-your-write; tune per use case
Conflict resolutionLast-write-winsVector clocks + app mergeLWW for simple values; vector clocks only if concurrent writes to the same key are semantically meaningful (shopping carts — the Dynamo paper's case)
MembershipCentral coordinator (ZooKeeper)GossipGossip — no single point of failure, matches the "no master" ethos; slower convergence is fine
Rebalance on joinEager full streamLazy + throttledThrottled — a new node pulling 30 GB at line rate would starve live traffic

🔥 Failure modes

🛠️ What I'd actually build

I wouldn't build this from scratch — I'd run Cassandra or ScyllaDB (same architecture: ring, vnodes, gossip, tunable quorums, hinted handoff — it's literally the Dynamo paper productized). Client: the DataStax driver with token-aware routing so coordinators are usually replicas themselves. N=3, R=2, W=2 across 3 racks. If the workload is pure cache-like KV with no range scans, DynamoDB (the AWS one) buys the same model with zero ops. Custom code: the operator that watches node health and throttles rebalances. The interview wants the concepts — the production answer is "use the thing that already implements the paper."

🎤 Interview tips

Go deeper: when NOT to use this

Need range scans ("all keys starting with user:")? The hash ring destroys ordering — you want Bigtable-style ordered partitioning instead. Need transactions across keys? You want Spanner/CockroachDB. The Dynamo model is for huge, simple, always-on KV — shopping carts, session state, user profiles — not for everything.

🎮 Interactive widget: the hash ring

120 keys (dots) live on the ring, owned by the first node clockwise. Add or remove a node and watch the counter: only ~K/n keys move. Toggle the replication factor to see replicas fan out.

Replication factor:
keys moved: 0 expected ≈ K/n: — naive hash%n rehash would move: — nodes: 4

primary copy replica node

💡 Replicas sit on the next N−1 nodes clockwise from the primary — the same rule real systems use for the preference list.

v2026.10.03-01