Course Content
System Design Interview
31 sections · 71 lessons
Key-Value Store: failure handling and the read and write paths
Nodes fail constantly at this scale. The system's job is to make failure ordinary: detect it without a central coordinator, keep accepting writes through it, and repair the divergence afterwards.
This lesson finishes the key-value store. The first half covers failure: how nodes learn who is alive, how writes keep succeeding while a replica is away, and how replicas find and fix the keys on which they disagree. The second half answers the question that closes this problem more often than any other — "walk me through a write" — by tracing both paths through every component built in this section.
Membership detection: gossip
There is no coordinator that knows which nodes are alive — a coordinator would be a single point of failure and a bottleneck. Instead, nodes tell each other, in a gossip protocol:
- Every node holds a membership list: each node's identifier, a heartbeat counter, and when that counter was last seen to change.
- Every second, each node increments its own counter and sends its list to a few randomly chosen peers.
- A receiving node merges the lists, keeping the highest counter it has seen for each member.
- If a member's counter has not advanced for longer than a threshold, it is marked suspect, then dead, and that judgement propagates the same way.
Information spreads exponentially — each round roughly doubles the number of informed nodes — so a 1,000-node cluster converges in around log₂(1000) ≈ 10 rounds, about ten seconds at a one-second interval. Bandwidth per node stays constant regardless of cluster size, which is the property that makes gossip scale.
The honest caveat: gossip is eventually consistent about membership. For a period, different nodes hold different views of who is alive, and requests can be routed to a node others consider dead. Systems that need agreement on membership — leader election, for instance — use a consensus protocol such as Raft instead, and pay for it in coordination.
Temporary failure: sloppy quorums and hinted handoff
A node in a key's preference list is unreachable. Strictly, a W = 2 write to a 3-replica key with two nodes down must fail.
A sloppy quorum relaxes this: the write goes to the first W reachable nodes on the ring, even if they are not in the key's preference list. The write succeeds; availability is preserved.
The substitute node stores the data with a hint recording who it was really for. When the intended node returns, the substitute hands the data over and deletes its copy — hinted handoff.
The trade is explicit and worth stating: a sloppy quorum breaks the R + W > N overlap guarantee, because the write may not be on any node the read consults. You have exchanged consistency for write availability during the failure window.
Permanent failure: Merkle trees for anti-entropy
Hinted handoff covers brief outages. A node that is down for hours, or one that dropped writes, needs its data compared against a replica — and comparing a terabyte of keys directly is prohibitive.
A Merkle tree is a tree of hashes: leaves hash individual key ranges, and each internal node hashes its children. Two replicas compare trees top-down:
- Root hashes match → the entire ranges are identical. One comparison, done.
- Root hashes differ → compare the two children, and recurse only into subtrees that differ.
Divergence in a few keys is found in O(log n) comparisons instead of scanning everything. For a range of a million keys with three differing, that is roughly 20 levels of comparison rather than a million.
The cost: the tree must be kept current as data changes, and the ranges must be aligned between replicas — when the ring changes, trees must be rebuilt, which is one more reason membership changes are throttled.
Alongside this, read repair does opportunistic maintenance: when a read collects R responses and finds one stale, it writes the newest value back to the lagging replica. Frequently read keys therefore repair themselves; anti-entropy exists for the ones nobody reads.
Multi-datacentre replication
Place at least one replica in another datacentre (Distributing the data) and replicate across regions asynchronously, because a synchronous cross-region write costs 100–150 ms against 0.5 ms locally. The consequence to state: a regional loss can lose the last few hundred milliseconds of acknowledged writes. If that is unacceptable, you must pay the cross-region round trip on every write, and there is no third option.
The write path, end to end
With every component in place, here are both paths through all of them.
- Client sends
put(key, value)to any node in the cluster. There is no leader, so any node can serve as coordinator for this request. Clients that hold a copy of the ring can contact a node in the preference list directly and skip a hop. - Coordinator locates the key. Hash the key onto the ring, walk clockwise, collect the first N = 3 distinct physical nodes — the preference list.
- Coordinator assigns a version. It increments its counter in the value's vector clock, producing a version that can later be compared for causality.
- Coordinator sends the write to all N replicas in parallel, and waits.
- On each replica: append the record to the write-ahead log (sequential disk write), then insert into the memtable. That is the whole durable write — no random I/O.
- Coordinator waits for W = 2 acknowledgements, then returns success to the client. It does not wait for the third, so the write completes at the speed of the second-fastest replica — 5 ms in the example from Quorums and tunable consistency, not 40 ms.
- If a replica is unreachable, the coordinator sends the write to the next reachable node on the ring with a hint (sloppy quorum), so the write still succeeds.
- Later, in the background: the memtable fills, is flushed to an immutable SSTable, and compaction merges SSTables. Hinted data is handed back when the intended node returns. Merkle tree comparison repairs anything still divergent.
The read path, end to end
- Client sends
get(key)to any node, which becomes coordinator. - Coordinator finds the preference list exactly as for a write.
- Coordinator requests the value from R = 2 replicas in parallel.
- On each replica: check the memtable first. Then, for each SSTable newest to oldest, consult its Bloom filter — a definite "not here" skips the file entirely, avoiding roughly 99% of unnecessary reads — and for a possible hit, use the sparse index to read one block.
- Coordinator compares the returned versions using their vector clocks. If one descends from the other, return the descendant. If they are concurrent, return both and let the client merge them (Conflict resolution).
- Read repair: if a replica returned a stale version, write the newest one back to it asynchronously. The read has already returned; the repair costs the user nothing.
The one-sentence version
If the interviewer wants it compressed:
"A write is coordinated by any node, placed on the ring, sent to three replicas where it is an append to a log plus a memtable insert, and acknowledged once two replicas confirm. A read is coordinated the same way, asks two replicas, checks memtable then Bloom-filtered SSTables on each, compares vector clocks, returns the descendant or both siblings, and repairs any stale replica in the background."