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.
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.
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”.
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.