Vibe Engines
YouTube
System Design

Design a Key-Value Store

Step 1 / 9

Learn system design by building a distributed key-value store like DynamoDB or Cassandra step by step.

The numbers to beat1/Nkeys moved on changevnodeseven spreadO(1)key → node

The whole design, in writing

Learn system design by building a distributed key-value store like DynamoDB or Cassandra step by step. An interactive guide covering consistent hashing, replication, quorum reads/writes, conflict resolution, gossip membership, and hinted handoff.

Every step of the build above, written out: the problem each piece solves, the option that was taken and the ones that were not, the numbers, and how it fails in production.

The big idea

What is a key-value store?

Sometimes you don’t need SQL — you need a giant, always-on dictionary: put(key, value) and get(key), spread across hundreds of machines, that never goes down even when machines do. A single database can’t offer that.

Clientget / put(key)
New in this step: Client.

Build a distributed, replicated store with no single master. Spread keys across many nodes, keep several copies of each, and let any node serve a request. The hard part isn’t the dictionary — it’s staying available and consistent while machines fail.

What the new pieces do

Clientclient
Wants a simple contract: store a value under a key, get it back later, at massive scale and without ever going down. No joins, no SQL — just key in, value out.

Step 1 · The skeleton

Any node, one value

A client wants to put and get a key without knowing or caring which of hundreds of machines actually holds it. Routing every request through one master would just recreate the single point of failure we’re trying to avoid. So how does a request find its key?

CoordinatorReplica A
New in this step: Coordinator, Replica A. · swipe to pan the diagram

A client must put/get a key across hundreds of machines, and it must never go down. Routing?

  1. The master is exactly the single point of failure you’re trying to escape — when it dies the whole store stops. Masterful designs trade away the availability you can’t sacrifice here.

  2. If that node dies those clients are stranded, and load can’t rebalance. Ownership must be able to move, and any node must be able to step in.

  3. Every node runs the same code and can locate replicas and serve the request. No special master means no single failure stops the world — symmetry is the whole philosophy.

Let the client hit any node, which acts as the coordinator for that request: it finds the right home for the key and reads/writes there. Every node can coordinate, so there’s no special, fragile master.

Why this piece earns its place

Symmetry buys availability, and it quietly costs you a hop and an exactly-once story. The hop first: a client that picks a node at random usually lands on one that doesn’t hold the key, so every request pays coordinator → replica on top of the replica round trip. The fix isn’t a load balancer in front, it’s a ring-aware client driver that learns the ring from the cluster and sends the request straight at a node that owns the key, so the extra hop only reappears during topology changes. The retry story matters more. When a coordinator dies mid-write the client retries somewhere else, and that second attempt can reach replicas the first one already wrote. That is survivable only because the contract is put and get: a whole-value overwrite applied twice is the same as applied once. The moment someone asks for a server-side increment or append, the very retry that made this design available starts double-counting, and the operation needs the version metadata of a later step before it is safe to repeat.

What the new pieces do

Coordinatorbackend
Whichever node the client hits. It owns the request: locate the replicas, talk to them, and assemble the answer. There is no special master — every node can coordinate.
Replica Astore
The first node responsible for a key (the one the ring lands on). Holds the data on local disk, typically as an LSM-tree for fast writes.

Step 2 · Where does a key live?

Consistent hashing

With hash(key) % N, changing the number of nodes remaps almost every key — a catastrophic reshuffle every time you add or lose a machine. At scale, nodes change constantly.

CoordinatorHash RingReplica A
New in this step: Hash Ring. · swipe to pan the diagram

Nodes are added and lost constantly. How do you map keys to nodes?

  1. Changing N remaps almost every key — a catastrophic full reshuffle every time a machine joins or dies. Modulo hashing and a churning cluster are incompatible.

  2. Adding or removing a node moves only the keys between it and its neighbor — about 1/N of the data, not all of it. Virtual nodes spread that load across many survivors.

  3. A key→node table for billions of keys is huge, hot, and another single point of failure. The ring computes ownership locally from the hash — no lookup table needed.

Place nodes and keys on a hash ring. A key belongs to the first node clockwise from its hash. Add or remove a node and only the keys between it and its neighbor move — a tiny fraction, not the whole dataset.

Why this piece earns its place

Virtual nodes are the knob this step hands you, and it has a real setting at both ends. Give each machine too few points on the ring and the split is lumpy: ranges land unevenly, so one node owns more data and more traffic than its neighbours, and when it dies its whole range falls on the single node behind it — the pile-up the ring was supposed to prevent. Push the count high and different things break. Every node now carries ring state proportional to the total number of tokens, and that state is exactly what gossip has to spread and every node has to hold; streaming and repair get chopped into many tiny ranges, each with its own per-range overhead; and one machine’s recovery fans out to essentially every other machine, so a local event becomes cluster-wide work. Set it from the cluster size you expect to grow into, not the one you have today: enough tokens that a departing node’s load scatters across many survivors, few enough that adding a machine stays cheap.

  • 1/Nkeys moved on change
  • vnodeseven spread
  • O(1)key → node

What the new pieces do

Hash Ringservice
Maps both keys and nodes onto a circle. A key is owned by the next node clockwise, so adding or removing a node moves only a small slice of keys.

Back of the envelope

add/remove node ⇒ ~1/N keys move
vs ~all keys with hash % N
each node ⇒ many vnodes
load spreads evenly; a death scatters across many survivors
key → node = O(1) hash
computed locally, no directory lookup

Step 3 · Survive a death

Replication

If a key lives on exactly one node, that node’s failure means the key is gone and unreachable. Disks and machines fail constantly at scale — single copies are not an option.

Hash Ringconsistent hashingReplica Bnode N+1Replica Cnode N+2
New in this step: Replica B, Replica C.

Disks and machines fail constantly. Where do you keep each key?

  1. A single copy means the key is unreachable the instant that node fails, and backups are minutes-to-hours stale. At this scale a node is always down somewhere.

  2. Full replication is ruinously expensive on storage and write cost, and pointless — you only need enough copies to survive a few failures, not all of them.

  3. Three consecutive ring nodes (N=3) each hold a copy, so the key survives up to two failures and reads can hit the healthiest replica. The preference list skips vnode duplicates so copies land on distinct machines.

Store each key on the N nodes following its position on the ring (its “preference list”). With N=3, three consecutive nodes each hold a copy, so the key survives failures and can be read from whichever replica is closest or healthiest.

Why this piece earns its place

Tolerating two failures is arithmetic that assumes the three failures are independent, and the ring knows nothing about physical reality. Consecutive positions on the ring are consecutive hashes, not distinct racks, power feeds or availability zones, so nothing stops all three copies of a key sitting behind one switch. The day that switch goes you did not lose two of three, you lost the key — and every key in that range — while the console still reports N=3. So the preference list has to walk further: skip a candidate whose failure domain the list already used and take the next distinct one. That costs something honest. The replica set is no longer the tidy ring neighbourhood, some writes now cross a zone boundary and pay for the crossing, and a domain holding fewer machines takes a disproportionate share of the copies. Raising N is the tempting alternative and it is the weaker one — more copies multiply every write, every repair and every stored byte while leaving the correlation exactly where it was. Independence is worth more here than count.

  • N=3copies / key
  • ring-adjacentreplicas
  • 2failures tolerated

What the new pieces do

Replica Bstore
The next node clockwise on the ring. Holds a copy so the key survives one node’s death and can serve reads in parallel.
Replica Cstore
The third copy. With N=3, the key lives on three consecutive ring nodes, so it tolerates two failures and spreads read load three ways.

Back of the envelope

N=3 copies on adjacent ring nodes
survives 2 simultaneous failures
preference list skips vnode dups
copies land on distinct machines / racks
read any healthy replica
spreads read load N ways

Step 4 · Consistent or available?

Quorum reads & writes

With three copies, who’s the source of truth? Waiting for all replicas means one slow node stalls every request; trusting one risks reading stale data right after a write. CAP says you can’t have perfect consistency and availability under partitions.

CoordinatorQuorum R/WReplica B
New in this step: Quorum R/W. · swipe to pan the diagram

With 3 copies, how many do you wait for on a read and a write?

  1. One slow or dead replica now stalls every request — you’ve traded away availability entirely. Strict all-replica waits defeat the point of replicating.

  2. Fast, but a read right after a write can hit a replica that hasn’t seen it — stale data. With no overlap guarantee you can’t promise the latest value.

  3. When the read set and write set must overlap, every read sees the last write. R=W=2 on N=3 gives strong consistency with one-node fault tolerance — and you can re-tune per workload.

Use a quorum: a write succeeds after W replicas ack; a read consults R. Choose R + W > N and every read set overlaps the last write set — so you always see the latest value. Tune R and W to slide between fast and strongly-consistent.

Why this piece earns its place

The overlap rule tells you what a quorum guarantees; it says nothing about when. So pick the numbers with latency in mind too. At N=3 with W=2 you wait for the second-fastest replica, which means the slowest of the three is hedged out of every request — and on most days that, not the consistency property, is the biggest thing quorums do for your tail. Push W up to N and you throw the hedge away: you are back to waiting on the worst replica, and one slow disk stalls every write to that key. Pull W down to 1 and writes get fast and thin — acknowledged from a single machine, so the write survives only as long as that machine does, until the others catch up. Which end you lean on should follow the workload, with one caution: a read-heavy service can afford R=1 only if it is willing to fail writes whenever a replica is down, and for a store whose promise is to stay writable that is the wrong slack to spend. Note also that the rule is about sets of replicas, not about time — two clients can each satisfy W=2 concurrently and both succeed, which is precisely the divergence the next step exists to clean up.

  • R + W > Noverlap rule
  • W=2,R=2common N=3
  • tunableC vs latency

What the new pieces do

Quorum R/Wservice
The tunable knob: a write waits for W replicas, a read for R. Make R + W > N and every read overlaps the latest write — consistency you can dial.

Back of the envelope

R + W > N ⇒ sets overlap
every read intersects the last acknowledged write
N=3: W=2, R=2
strong consistency + 1-node fault tolerance
W=1 fast writes / R=1 fast reads
slide the knob per workload

Step 5 · Two truths

Resolving conflicts

During a network partition, two replicas can each accept a write to the same key. When the partition heals, they disagree — and with W < N this is expected, not a bug. Which value wins?

Quorum R/WR + W > NConflict ResolverversionsReplica Bnode N+1Replica Cnode N+2
New in this step: Conflict Resolver.

A partition lets two replicas each accept a write to the same key. On heal, who wins?

  1. Rejecting writes during a partition sacrifices the availability that’s this system’s whole reason to exist. In an AP store concurrent writes are expected, not preventable.

  2. LWW is simple but silently drops a real concurrent update, and clock skew makes "last" unreliable. Fine for some data, data-loss for a shopping cart.

  3. Vector clocks capture causality, so the system distinguishes a newer write from genuine divergence and can surface both siblings for the app (or a merge function) to reconcile — carts merge, counters add.

Attach version metadata to every write. Last-write-wins (by timestamp) is simple but can drop data. Vector clocks capture causality, detecting true conflicts and surfacing both versions (siblings) for the application — or a merge function — to reconcile.

Why this piece earns its place

Detecting conflicts instead of dropping them means the conflicts now have to go somewhere, and where they go is into the value. Every divergence nobody reconciles adds a sibling, and siblings come back together on the next read, so a key the application keeps writing and never merges grows until reading it is slow and eventually until it cannot be read at all — the classic way a cart becomes a payload nobody can fetch. Reconciliation therefore belongs on the read path and has to be written back: merge the siblings, store the merged value under a version that descends from all of them, and the branch collapses. It also explains why the store cares about the shape of your data — values that merge deterministically, a set union or an add-only counter, can be resolved by the system, while an opaque blob always needs the application in the loop. And the version metadata is not free: a clock grows an entry per node that has coordinated a write to that key, so a long-lived hot key gets its clock truncated at some cap, and a truncated clock has lost the link that proved one write descended from another. That shows up as a conflict between versions that were really ordered — extra siblings, never lost data.

What the new pieces do

Conflict Resolverservice
Two replicas can hold different values for the same key after a partition. Vector clocks (or last-write-wins) decide which version survives, or surface both.

Step 6 · Who’s alive?

Gossip membership

Coordinators need to know which nodes are up and who owns which ring range — but a central membership registry would be another single point of failure and a bottleneck.

CoordinatorQuorum R/WConflict ResolverReplica BReplica CGossip Membership
New in this step: Gossip Membership. · swipe to pan the diagram

Nodes gossip: each periodically exchanges health and ring state with a few random peers. Within seconds the whole cluster converges on a shared view of membership — fully decentralized, no coordinator required.

Why this piece earns its place

Gossip makes membership cheap and makes it late, and the lateness is the part you design around. While the news spreads, two coordinators hold different pictures of the cluster, so one can route at a node the other has already written off. That is normal, and it is the reason ownership of data must not be something a failure detector is allowed to change. Keep the two apart: gossip about liveness marks a node temporarily unavailable, which is all the other mechanisms need in order to route around it, while moving ring ranges stays a deliberate operation that someone or something decides to run. Then the knob is the suspicion threshold. Make it twitchy and a node that is merely busy — a long garbage-collection pause, a saturated network card — gets declared dead, the cluster reacts, it returns seconds later, and you have bought churn and repair work for nothing. Make it patient and every request aimed at a genuinely dead node burns a full timeout before anything else is tried, which surfaces as a latency cliff long before it surfaces as an outage. The version worth building judges a node against its own recent gossip intervals rather than against a fixed number of seconds.

What the new pieces do

Gossip Membershipbus
Nodes periodically swap health and ring info peer-to-peer. No central registry — the cluster collectively knows who is up and who owns what.

Step 7 · Heal while degraded

Hinted handoff & read repair

If a target replica is temporarily down during a write, do you reject the write (hurting availability) or accept it and risk that replica permanently missing the update?

ClientCoordinatorHash RingQuorum R/WConflict ResolverReplica AReplica BReplica CGossip Membership
The system as it stands at this step. · swipe to pan the diagram

A target replica is down during a write. Reject the write, or accept it?

  1. That makes writes fail whenever any replica is down — and something is always down at scale. You’d be choosing consistency over the availability you promised.

  2. Hinted handoff keeps writes always-succeeding: a stand-in node stores the write and forwards it when the rightful replica returns. Read repair and Merkle-tree anti-entropy converge the rest.

  3. Then that replica is permanently stale and quorums silently weaken. You must track the missed write (a hint) and replay it — accepting without healing just hides divergence.

Accept it: a healthy node stores a hint and replays the write to the rightful replica once it returns (hinted handoff). Meanwhile, reads that notice a stale replica push the fresh value back to it (read repair), and background anti-entropy (Merkle trees) reconciles the rest.

Why this piece earns its place

Be precise about what a hint buys, because it is easy to over-claim: hinted handoff is an availability mechanism, not a durability one. The hint lives on the stand-in, which is a machine like any other — if it dies before the replay, the write exists only on the replicas that did take it and nothing records that anyone owed a copy. That is acceptable because the write was already acknowledged under the quorum you promised and the hint rode on top; it stops being acceptable the moment you count the stand-in toward W, because that is counting a node which does not own the key. Hints also have to stop. A node down for an hour is a different animal from one down for a week: survivors accumulate everything addressed to it, and when it finally returns they all replay at once, at full speed, at the exact moment it is least able to absorb load. So cap how long hints are kept, throttle the replay, and past that window let the returning node get its data the slow honest way, through the Merkle comparison, which does not care how long it was gone. Read repair has its own blind spot — it only touches keys somebody reads, so the cold ones drift until the background job reaches them. That job is not optional.

You did it

You just designed a key-value store.

ClientCoordinatorHash RingQuorum R/WConflict ResolverReplica AReplica BReplica CGossip Membership
The finished design, end to end. · swipe to pan the diagram

Everything you assembled, in order

  • Any node coordinates — no master, no single point of failure.
  • Consistent hashing maps keys to nodes; only 1/N keys move on change.
  • Each key replicated to the next N nodes on the ring for durability.
  • Quorum reads/writes with R + W > N give tunable consistency.
  • Vector clocks detect and resolve conflicting concurrent writes.
  • Gossip spreads membership and ring state with no central registry.
  • Hinted handoff and read repair keep replicas converging through failures.

Where an interviewer pokes next

Getting the boxes right is the easy half. These are the questions that separate a candidate who drew the diagram from one who has run the thing. Answer each one out loud before you open it.

  1. Why choose AP (available) over CP here?

    A key-value store’s whole promise is "never down at scale." Under a partition CAP forces a choice; Dynamo-style stores pick availability and reconcile later via quorums, vector clocks and read repair. If you need strict consistency you raise R+W>N or pick a CP store — it’s a per-workload dial, not a fixed law.

  2. How are hot keys / hot partitions handled?

    A single wildly popular key can overwhelm its N replicas. Mitigations: more vnodes for finer spread, a cache in front of hot keys, or splitting a hot key’s value (e.g. sharded counters). The ring spreads the average; hotspots still need bespoke handling.

  3. What does a node do when it rejoins the cluster?

    Gossip propagates that it’s back; it streams the key ranges it owns from peers and accepts any hinted-handoff writes buffered for it. Anti-entropy (Merkle-tree diffing) reconciles whatever drifted while it was gone, so it converges without a full copy.

  4. Can a get() still be stale even with R+W>N?

    In edge cases yes — failed nodes and "sloppy quorums" that write to fallback nodes can break the clean overlap. R+W>N guarantees overlap in the normal case; read repair fixes stragglers it touches; apps that need certainty use vector clocks to detect and merge. "Strong-ish, eventually exact" is the bargain.

  5. How is data stored on each node?

    Typically an LSM-tree: writes append to an in-memory memtable + commit log, flush to immutable SSTables, and compact in the background. Writes are fast and sequential — ideal for a write-heavy distributed store — with Bloom filters keeping reads cheap.

Check yourself — the answers, and why

Eight steps in, these are the calls you should be able to make cold. Pick one, then read why.

  1. There is no master node because…

    Symmetry: every node runs the same code; a dead node is routine, not catastrophic.

  2. Consistent hashing matters because on a node change…

    A key moves only between a node and its ring neighbor — not the whole dataset like hash % N.

  3. R + W > N guarantees that…

    Overlap means every read sees the latest acknowledged write.

  4. Vector clocks are used to…

    They capture causality so the system surfaces siblings instead of silently dropping data.

  5. Hinted handoff lets a write…

    A stand-in holds the write and forwards it when the rightful replica returns — availability now, consistency soon.

How you’d open this design in an interview

Before any boxes: agree what it must do, pin the qualities that shape everything, then build — naming each trade-off as you make it. The walkthrough above is that exact order.

What it must do

Agree on these before drawing a single box.

  • Get / Put: put(key, value) and get(key) — a giant dictionary, no SQL, no joins.
  • Any node serves: a client hits any node, which coordinates the request — no special master.
  • Survive failure: a key stays available and durable through node deaths.
  • Tunable consistency: callers dial between fast and strongly-consistent per workload.
  • Reconcile conflicts: concurrent writes to a key during a partition are detected, not silently lost.

The qualities that shape everything

Each one names the mechanism that buys it.

Never go down at scale (availability first)
A masterless, symmetric design — every node runs the same code and can coordinate, so a dead node is routine, not an outage.
Minimal reshuffle when nodes come and go
A consistent-hash ring: adding or removing a node moves only ~1/N of keys, and virtual nodes spread that load across many survivors.
Durability through node loss
Replicate each key to the next N nodes on the ring (N=3), so it survives up to two simultaneous failures and reads spread across replicas.
Read-your-writes when you need it
Quorum with R + W > N makes every read set overlap the last write set — turn the knob toward consistency or latency per workload.
Always writable, even mid-failure
Hinted handoff lets a stand-in accept a write and replay it later; read repair and Merkle anti-entropy converge replicas afterward.
Decentralized membership
Nodes gossip health and ring state peer-to-peer — no central registry to become a bottleneck or single point of failure.

The trade-offs you say out loud

Senior signal isn’t the boxes — it’s naming what you gave up and why it was the right price.

Consistent-hash ring over hash(key) % N

Modulo hashing remaps almost every key when the node count changes — a catastrophic reshuffle in a churning cluster. The ring moves only the keys between a node and its neighbor.

Quorum R + W > N over wait-for-all or trust-one

Waiting for all replicas lets one slow node stall every request; trusting one risks stale reads. R + W > N guarantees the read set overlaps the last write — consistency you can tune.

Vector clocks over last-write-wins

LWW by wall-clock is simple but silently drops a real concurrent update, and clock skew makes “last” unreliable. Vector clocks capture causality, surfacing true conflicts for the app to merge — carts merge, counters add.

AP — available over CP — strictly consistent

The store’s whole promise is “never down.” Under a partition CAP forces a choice; a Dynamo-style store favors availability and reconciles later. Need strict consistency? Raise R + W > N — it’s a dial, not a law.

Hinted handoff over rejecting the write

Rejecting writes whenever a replica is down means writes fail constantly at scale. A stand-in holds a hint and replays it when the rightful replica returns — availability now, consistency soon.

The answer, out loud

What a strong answer to “Design a Key-Value Store (like DynamoDB)” sounds like, first question to last trade-off. It is about 7 minutes of talking; the whiteboard and the interviewer fill the rest of the 45. Read it aloud once, then close the page and give it yourself.

  1. 0–4 min

    What we’re actually building

    Let me pin down what we’re building first, because “key-value store” covers two very different systems. I’m reading this as the always-on kind: put a value under a key, get it back, no joins, no SQL, spread over hundreds of machines, and the thing it must never do is go down. That last word does most of the design work. If availability is the headline requirement then I already know that under a partition I’m going to keep taking writes and reconcile afterwards, and I’d confirm that’s the trade you want, because it changes almost every box on the board.

  2. 4–9 min

    Take the master out

    So the first move is to remove the master. A client sends its request to any node, and whichever node it lands on coordinates: work out where the key lives, talk to the machines that hold it, assemble the answer. Every node runs the same code, so there’s nothing special to fail over and nothing special to fail. The test I like here is simple — if I kill the node you’re currently talking to, does anything happen beyond one retry? It doesn’t, and that isn’t a nice side effect, it’s the whole point.

    Built in step 1: Any node, one value
  3. 9–15 min

    Where a key lives

    Next, where a key actually lives. The obvious answer is hash the key modulo the node count, and it’s wrong for the same reason the master was wrong: it assumes the cluster holds still. Change the number of nodes and nearly every key moves. So I put both nodes and keys on a ring and give a key to the first node clockwise from it. Now adding or losing a machine only moves the keys between it and its neighbour — roughly one N-th — and I’d give each machine many positions on that ring so the slice scatters instead of landing on one unlucky survivor.

    Built in step 2: Consistent hashing
  4. 15–21 min

    Copies, and how many

    A single copy is one disk away from gone, so each key goes to the next three distinct nodes clockwise — its preference list. Two of them can be dead and I’m still serving. I’d say out loud that I’m choosing three rather than deriving it: it’s the setting where you survive the ordinary double failure without paying to write everything everywhere, and it’s per-table, not a constant. Reads then go to whichever copy is healthiest, so this buys read capacity as well as durability.

    Built in step 3: Replication
  5. 21–28 min

    How many do I wait for

    Now the interesting question: with three copies, who’s right? I don’t want to wait for all three, because then the slowest machine sets my latency and a dead one stops me outright. I don’t want to trust one, because a read can land on a copy that hasn’t heard the last write. So: quorums. A write waits for W copies, a read asks R, and I set R plus W greater than N — two and two on three. Any read then overlaps whatever the last acknowledged write touched, so I see it. And it stays a dial: a workload that wants cheap writes turns W down and accepts what that means.

    Built in step 4: Quorum reads & writes
  6. 28–34 min

    When two writes both happened

    Quorums don’t stop two people writing at once during a partition, and in a system that stays available they will. The easy answer is last-write-wins on a timestamp, and it quietly deletes somebody’s update — and “last” leans on clocks agreeing, which across a fleet they don’t. So I version each write with a vector clock that records what the writer had already seen. The system can then separate “this came after that” from “these genuinely happened together”, and only the second case needs a decision — which I’d hand to the application, because a cart merges and a counter adds.

    Built in step 5: Resolving conflicts
  7. 34–39 min

    Who’s alive

    All of that needs the nodes to know about each other — who’s up, who owns which range. I’d deliberately not build a registry for it, because a registry is a master wearing a different hat. Instead nodes gossip: each one periodically picks a few peers and swaps what it knows, and the cluster converges on a shared picture with nobody in charge. That’s why I like it, and in the same breath I’d name its weakness — the picture is always a little out of date, so the rest of the design has to tolerate two nodes briefly disagreeing about reality.

    Built in step 6: Gossip membership
  8. 39–43 min

    Writing while degraded

    Last piece: what happens to a write when one of a key’s three homes is down. Rejecting it means writes fail routinely, because at this size something is always down. So a healthy node takes the write on that replica’s behalf and holds a hint, and replays it when the replica returns. Alongside that, a read that notices one copy is behind pushes the fresh value back to it, and a background job compares replicas with Merkle trees and fixes whatever nobody happened to read. Available now, converged shortly after.

    Built in step 7: Hinted handoff & read repair
  9. 43–45 min

    The trade, and what’s left

    The trade in one line: I’ve chosen to stay up and reconcile afterwards, and everything from the quorum numbers to the vector clocks is machinery for making “afterwards” short and honest. It won’t give you a transaction across two keys, and I’d say that before somebody builds a ledger on it. With more time I’d go at hot keys, because the ring spreads the average and does nothing for one wildly popular key, and I’d want to talk about what sits on each node’s disk — an LSM-tree is what makes a write-heavy store like this cheap.

What this teaches

Learn system design by building a distributed key-value store like DynamoDB or Cassandra step by step. An interactive guide covering consistent hashing, replication, quorum reads/writes, conflict resolution, gossip membership, and hinted handoff.

Key takeaways

  • Any node coordinates — no master, no single point of failure.
  • Consistent hashing maps keys to nodes; only 1/N keys move on change.
  • Each key replicated to the next N nodes on the ring for durability.
  • Quorum reads/writes with R + W > N give tunable consistency.
  • Vector clocks detect and resolve conflicting concurrent writes.
  • Gossip spreads membership and ring state with no central registry.
  • Hinted handoff and read repair keep replicas converging through failures.

Concepts covered

  • What is a key-value store?
  • Any node, one value
  • Consistent hashing
  • Replication
  • Quorum reads & writes
  • Resolving conflicts
  • Gossip membership
  • Hinted handoff & read repair
RUN IT YOURSELF

Consistent hashing, in Python & TypeScript

A distributed KV store decides which node owns each key with consistent hashing. Here is the ring in both languages, running live. Switch tabs, read the comments, edit the nodes, and hit Run.

HOW TO READ THE CODE — 4 IDEAS
  1. Hash both nodes and keys onto a ring (here, 0–359 degrees).
  2. A key belongs to the first node clockwise from its position (steps 2–3).
  3. Adding or removing a node only moves the keys between it and its neighbour — not everything.
  4. Contrast plain hash(key) % N, where changing N reshuffles every key.
CPython · WebAssembly
built to be reasoned about, not memorized — make the calls, drop a replica, run the gauntlet.
Finished this one? 0 / 65 System Designs done

Explore the topic

See this alongside everything else on the same subject — handbooks, system designs, challenges and tools, in one place.

More System Designs