Design a Key-Value Store: Quorums, Versions and Repair
9 min readBytePatterns
Design a distributed key-value store for a system design interview: N, W and R quorums, why a failed write can still appear, vector clocks and read repair.
"Design a key-value store" sounds like the easiest prompt in the book: two operations, put(key, value) and get(key), and a hash map already does both. The interview is about what a hash map on one machine never faces: the data must survive a dead disk, the store must keep answering while a node is unreachable, and two copies of one key will, sooner or later, disagree.
The problem it solves
With illustrative assumptions: billions of small keys, values under a few kilobytes, many more reads than writes, and a latency budget measured in milliseconds. That forces three decisions.
- Partitioning. No single machine holds everything, so keys are spread over nodes. Consistent hashing is the usual tool: adding a node moves only the keys next to it on the ring.
- Replication. Each partition lives on
Nnodes, typically the nextNdistinct nodes clockwise on the ring, so one failure loses nothing. - Agreement. When is a write done, and which copy does a read believe? This is the part the interviewer is waiting for.
The intuition
Pick three numbers. N is the number of replicas, W is how many must acknowledge a write before the caller hears "ok", and R is how many a read consults. If R + W is greater than N, any set of R replicas shares at least one member with any set of W replicas, so every read touches at least one copy that saw the latest acknowledged write. Each value carries a version, and the reader keeps the highest one.
Three details separate a strong answer from a memorised one:
- A failed write is not rolled back. If only one replica acknowledged before the timeout, the caller hears "failed", yet that replica still holds the value, and a later read may return it. Quorums promise "at least as new as the last acknowledged write", not "exactly the acknowledged writes".
- Stale copies get repaired. A read that sees an older version on one replica writes the newer one back (read repair). A background job compares replicas, often by exchanging hash trees of key ranges, to catch keys nobody reads.
- Counters are not enough for concurrent writers. If two clients write through different nodes during a partition, neither write is newer. Vector clocks detect that case, and the store keeps both values as siblings for the application to merge.
Watch it run
The animation follows one key. It is hashed to a partition, and that partition lives on three machines, not one. A write arrives and the coordinator sends it to all three, then starts counting. Two acknowledge quickly, which is already enough for the write to be done, so the caller is told it succeeded while the third replica is still behind. Now a read: ask one replica, and the answer depends on which one you picked. So read two. With R + W greater than N, the two sets must share a replica, and one of the two replicas you touched is one the write reached, so the newest version is there. The disagreement is not a guess: the higher version wins and the stale copy is fixed. Then the dial moves. Raise W to three and durability improves, until any single replica is unreachable and the write fails. Drop W to one and writes never wait on a second replica, leaving two versions to reconcile later. Quorum is that dial; it does not remove the choice, it lets you set it.
Design a Key-Value Store
Step 1 of 11
The key is hashed to a partition, and that partition lives on three machines, not one.
The same interactive animation as the lesson — step through it with the controls.
The code
A toy model: one partition, three replicas as dictionaries, and a coordinator that numbers each key's writes. reached stands in for the network, listing the replicas a write got to before the caller stopped waiting:
import random
class Replica:
def __init__(self, name):
self.name, self.data = name, {} # key -> (version, value)
def read(self, key):
return self.data.get(key, (0, None))
def write(self, key, version, value):
if version > self.read(key)[0]: # never overwrite newer with older
self.data[key] = (version, value)
class QuorumStore:
"""Toy model: one partition, N replicas, W acks per put, R replies per get."""
def __init__(self, n=3, w=2, r=2, seed=0):
self.replicas = [Replica(f"r{i + 1}") for i in range(n)]
self.n, self.w, self.r = n, w, r
self.version = {} # the coordinator numbers each key's writes
self.rng = random.Random(seed)
def put(self, key, value, reached):
"""`reached` = replicas the write got to before the caller stopped waiting."""
v = self.version[key] = self.version.get(key, 0) + 1
for rep in reached: # a write that fails is NOT rolled back
rep.write(key, v, value)
return len(reached) >= self.w # success means W acknowledgements
def get(self, key, asked=None):
asked = asked or self.rng.sample(self.replicas, self.r)
replies = [(rep, rep.read(key)) for rep in asked]
version, value = max(reply for _, reply in replies)
for rep, (v, _) in replies: # read repair: fix the stale ones you saw
if v < version:
rep.write(key, version, value)
return value
store = QuorumStore(n=3, w=2, r=2)
r1, r2, r3 = store.replicas
print(store.put("cart", "v6", [r1, r2, r3]), store.put("cart", "v7", [r1, r2])) # True True
print(r3.read("cart")) # (1, 'v6') the third replica is behind
print(store.get("cart", [r3])) # v6 a single-replica read can be stale
print(store.get("cart", [r2, r3]), r3.read("cart"))
# v7 (2, 'v7') the quorum read found v7 and repaired r3
print(store.put("cart", "v8", [r1])) # False one ack is not W = 2 ...
print(store.get("cart", [r1, r3])) # v8 ... yet the "failed" write is visible
Checked on 600 seeded random configurations, N from 1 to 5 and every W and R, with writes that reach random subsets of replicas. A reference remembers the newest acknowledged version, and every read is compared with it:
def run(n, w, r, steps, rng):
"""Random puts that reach random subsets; every get is checked against the
newest acknowledged version. Returns how many gets returned something older."""
s = QuorumStore(n, w, r, seed=rng.random())
acked, stale = 0, 0
for _ in range(steps):
if rng.random() < 0.5:
reached = rng.sample(s.replicas, rng.randint(0, n))
if s.put("k", s.version.get("k", 0) + 1, reached):
acked = s.version["k"]
else:
got = s.get("k") or 0 # the value is the version number here
stale += got < acked
return stale
rng = random.Random(30)
overlap_ok, gaps_found = True, 0
for _ in range(600):
n = rng.randint(1, 5)
w, r = rng.randint(1, n), rng.randint(1, n)
stale = run(n, w, r, 60, rng)
if r + w > n:
overlap_ok &= stale == 0 # R + W > N: never older than the last ack
else:
gaps_found += stale > 0
print(overlap_ok, gaps_found > 0) # True True
Version counters assume one coordinator numbers every write. When two nodes accept writes independently, each keeps its own counter, and a vector clock holds all of them. One clock is "before" another only if none of its counters is larger:
def compare(a, b):
"""Vector clocks: 'before', 'after', 'equal' or 'concurrent'."""
nodes = set(a) | set(b)
le = all(a.get(x, 0) <= b.get(x, 0) for x in nodes)
ge = all(a.get(x, 0) >= b.get(x, 0) for x in nodes)
return "equal" if le and ge else "before" if le else "after" if ge else "concurrent"
def merge(*clocks):
return {x: max(c.get(x, 0) for c in clocks) for c in clocks for x in c}
# Both clients read the cart at A:1, then write through different nodes.
base = {"A": 1}
phone = {"A": 2} # coordinated by node A
laptop = {"A": 1, "B": 1} # coordinated by node B during a partition
print(compare(base, phone), compare(base, laptop), compare(phone, laptop))
# before before concurrent
siblings = {"phone": ({"milk", "eggs"}, phone), "laptop": ({"milk", "bread"}, laptop)}
cart = set().union(*(items for items, _ in siblings.values())) # the application merges
clock = merge(phone, laptop)
clock["A"] += 1 # the merged write is a new event on A
print(sorted(cart), sorted(clock.items()), compare(phone, clock), compare(laptop, clock))
# ['bread', 'eggs', 'milk'] [('A', 3), ('B', 1)] before before
Checked on 2,000 seeded random histories of local writes and messages between up to four nodes, against a brute force that computes "happened before" directly as reachability in the event graph:
def random_history(rng, nodes=3, events=12):
"""Local writes and messages. Returns each event's vector clock and the
happens-before edges (program order and send -> receive)."""
clocks, last, edges, inflight = [], {}, [], []
for e in range(events):
node = rng.randrange(nodes)
vc = dict(clocks[last[node]]) if node in last else {}
if node in last:
edges.append((last[node], e))
if inflight and rng.random() < 0.5: # receive a message: take the max
sender = inflight.pop(rng.randrange(len(inflight)))
vc = merge(vc, clocks[sender])
edges.append((sender, e))
vc[node] = vc.get(node, 0) + 1
clocks.append(vc)
last[node] = e
if rng.random() < 0.4:
inflight.append(e) # this event also sends a message
return clocks, edges
def reachable(edges, n):
reach = [{i} for i in range(n)]
for _ in range(n): # brute force: iterate to a fixed point
for a, b in edges:
reach[a] |= reach[b]
return reach
rng = random.Random(30)
ok = True
for _ in range(2_000):
clocks, edges = random_history(rng, nodes=rng.randint(1, 4), events=rng.randint(1, 14))
reach = reachable(edges, len(clocks))
for e in range(len(clocks)):
for f in range(len(clocks)):
if e == f:
continue
want = "before" if f in reach[e] else "after" if e in reach[f] else "concurrent"
ok &= compare(clocks[e], clocks[f]) == want
print(ok) # True
The complexity
- Write latency is the
W-th fastest replica's response, not the slowest; read latency is theR-th fastest. - Messages: a write goes to all
Nreplicas, a read to at leastR. - Storage:
Ncopies of every value, plus a vector clock whose size grows with the number of nodes that ever coordinated a write to that key, which is why real systems prune old entries.
Where it goes wrong
- Quoting
R + W > Nas "strongly consistent". It gives "no older than the last acknowledged write". A failed write can still surface, and two concurrent writers still conflict. - Last-write-wins by wall clock. Clocks drift, and the "loser" is silently dropped. Fine for a cache entry, wrong for a shopping cart.
W = N. Every write stops when any replica is down, which is the availability the replicas were meant to buy.- Forgetting repair. Without read repair and background comparison, a replica that missed writes stays stale indefinitely.
- Sloppy quorums. Some stores accept a write on a stand-in node while a home replica is down and hand it over later. That keeps writes available, but the overlap guarantee no longer holds until the handoff completes.
When it shows up in interviews
As "design a key-value store", "design a distributed cache", or a follow-up inside larger designs that need a durable map, such as a URL shortener. The partition-tolerance question underneath is the CAP theorem, the leader-based alternative is database replication, and the hot-key problem is covered in DynamoDB hot partitions.
How to say it in an interview
"I'd partition keys with consistent hashing and store each partition on N = 3 nodes. A put goes to all three and succeeds after W = 2 acknowledgements; a get reads R = 2 and returns the highest version, so R + W > N guarantees the read overlaps the last acknowledged write. A write that failed its quorum isn't rolled back, so the guarantee is 'at least that new'. Reads repair stale replicas, and a background job compares key ranges. For concurrent writers I'd use vector clocks and keep siblings for the application to merge, and I can tune W and R per use case: lower W for availability, higher for durability."