Design a key-value store

Build a distributed hash map: consistent hashing for placement, replication for durability, and the quorum dial between consistency and availability.

9 min read

A distributed key-value store comes down to three decisions: where a key lives, how many copies of it exist, and how many of those copies have to agree before you answer. The third is where all the interesting behaviour is.

Every node you add makes an outage more likely

Spread a dataset across a hundred machines with one copy of each key and you have not built a resilient system. You have built one with a hundred single points of failure, because availability across independent parts multiplies.

Chance every key is reachable on a given day
90.5%
About 35 days a year with some part of the keyspace unreachable.
100
99.90%
1
A hundred nodes at 99.9% with one copy each: some data is down about one day in ten. Add copies and watch it vanish.

0.999100 is about 0.905, so roughly one day in ten some part of the keyspace is unreachable. Keep three copies of each key on different nodes and the same hardware gets you something like nine nines. Replication is not an optimisation; it is what makes a distributed store possible at all. But as soon as there is more than one copy, the copies can disagree, and you have to decide how many get a say.

W + R > N is the whole idea

Keep N copies. A write counts as done once W of them acknowledge it; a read asks R of them and takes the newest answer. If W + R is bigger than N, the write set and the read set must share at least one replica, so every read reaches a copy that saw the latest write. It is a pigeonhole argument, and it is the entire mechanism.

write set: first 3read set: last 3click a replica to fail it
Read returned v2 from replicas 3 to 5: fresh, because replica 3 is in both sets.
W + R = 6 > N = 5: every read overlaps the latest write.
5
3
3
N = 5, W = 3, R = 3: the sets share a replica, so reads are always fresh. Try N = 3, W = 1, R = 1.

Turning W and R down buys speed and costs guarantees. With N = 3, W = 1, R = 1, a write acknowledged by one node is lost if that node dies before replicating, and a read straight after your own write can land on two copies that never saw it. Both are legitimate choices for some workloads. They are just choices, and the dial makes them visible.

When the network splits, pick which promise to break

Everything so far assumed the replicas can reach each other. When a network partition cuts them apart, a quorum system has exactly two options, and neither is “carry on normally”.

Replica ABo
Replica BBo
Replica CAdarefusing
majority side, wrote Boalone, wrote Cy
A and B form a majority, so they accepted Bo. C cannot reach a quorum, so it refused the write of Cy and refuses reads it cannot make safe.
During a partition you choose who gets an error (CP) or who gets a conflict (AP). Heal it to see the bill.

Consistency first (CP) means the minority side refuses what it cannot make safe: nobody sees a stale or conflicting value, and some users see an error. Spanner, etcd and ZooKeeper work this way. Availability first (AP) means both sides accept writes, versions diverge, and they are reconciled after the partition heals with vector clocks plus last-write-wins or an application merge: nobody sees an error, somebody sees a conflict. Dynamo, Cassandra and Riak work this way.

That is what the CAP theorem actually says. Not “pick two of three”: partitions are not optional, they happen to you. During a partition you choose between consistency and availability; the rest of the time you can have both, which is why the quorum dial is the knob you really spend your time turning. Where each key lives in the first place is the consistent hashing half of the design.

The short version

  • One copy per key across many nodes makes outages more likely; replicate.
  • With N copies, W write acks and R read replies, W + R > N guarantees fresh reads.
  • Lower W and R trade durability and freshness for speed.
  • In a partition, CP returns errors and AP returns conflicts; choose on purpose.