Course Content
System Design Interview
31 sections · 71 lessons
Key-Value Store: distribution, quorums and conflict resolution
One node cannot hold 10 TB or serve 100,000 operations per second reliably. Data must be split across machines and copied for durability, and those are two separate decisions.
This lesson is the core of the key-value store deep dive. It places the data on a ring and replicates it. It then shows how N copies let each operation choose how many copies to involve — the decision that makes consistency a dial rather than a fixed property. It ends with what happens when two clients write the same key at the same moment, both succeed, and the replicas disagree.
Placement by consistent hashing
Partition with a hash ring (Section 7). Every node is placed at many positions on a circular hash space; a key hashes onto the ring and belongs to the first node clockwise.
Why this and not hash(key) % N: adding a node to a 20-node cluster with the modulo scheme moves roughly 95% of the data — for a 10 TB cluster, 9.5 TB across the network while serving traffic. With the ring it moves about 1/21 of the data, or under 500 GB, and it moves from a small number of neighbours rather than from everywhere.
Use 100–500 virtual nodes per physical node, which keeps the busiest node within roughly 10–15% of average and spreads a failed node's load across many peers rather than dumping it on one neighbour.
Replication factor N, and the preference list
Each key is stored on N nodes, conventionally 3. Walk clockwise from the key's ring position and collect the first N distinct physical nodes; that ordered list is the key's preference list.
The word "distinct" is doing real work, as Replication on the ring explained: with virtual nodes, consecutive ring positions often belong to the same physical machine, and taking the next three positions can give you two copies on one box. You would believe you had three-way replication and actually have two machines of durability.
Extend the same skip to failure domains. Three replicas in one rack have the annual failure probability of that rack; three replicas in three racks have approximately its cube. For a 1% annual rack failure rate that is 1-in-100 against 1-in-a-million, from nothing but placement.
What replication factor 3 costs
- Storage: 3 × 10 TB = 30 TB. Say this out loud when quoting a storage number; candidates routinely quote the logical size and forget the multiplier.
- Write bandwidth: every write is sent to three nodes, so internal network traffic is roughly three times the client write traffic.
- Read capacity: it buys you 3× read capacity, since any replica can serve a read.
Multi-datacentre placement
Extend the preference list rule to span datacentres: of the three replicas, place at least one in a second datacentre. Now a datacentre loss does not lose data.
The cost is precise: a write that must be acknowledged by a replica in another region pays that region's round trip — 100–150 ms across continents against 0.5 ms locally, a factor of 200. The standard resolution is to acknowledge on the local quorum and replicate across regions asynchronously, accepting that a regional failure can lose the last few hundred milliseconds of writes. That is a real trade with a real cost, and stating it is the point.
Quorums and tunable consistency
With N copies of every key, the system can decide how many copies to involve in each operation.
Why R + W > N gives strong reads
Every write lands on some set of W nodes. Every read consults some set of R nodes. If R + W > N, those two sets cannot be disjoint — by the pigeonhole principle they must share at least one node. That shared node has the newest write, so the read sees it.
The read then needs a rule for choosing among the versions it receives, since R − 1 of the responses may be stale. That rule is a version number or timestamp, which is where the problems of conflict resolution (the end of this lesson) begin.
If R + W ≤ N, the sets can miss each other entirely, and a read can legitimately return a value older than an acknowledged write. That is eventual consistency.
Worked configurations, with N = 3
| W | R | R + W > N? | Write latency | Read latency | Tolerates | Use for |
|---|---|---|---|---|---|---|
| 1 | 1 | No (2 ≤ 3) | Fastest | Fastest | 2 nodes down | Metrics, view counts, caches |
| 2 | 2 | Yes (4 > 3) | Moderate | Moderate | 1 node down, both paths | The general-purpose default |
| 3 | 1 | Yes (4 > 3) | Slowest | Fastest | 0 nodes down for writes | Read-dominated, writes rare |
| 1 | 3 | Yes (4 > 3) | Fastest | Slowest | 0 nodes down for reads | Write-dominated, reads rare |
The latency cost, computed
Quorum latency is the W-th fastest response, not the average. Suppose the three replicas respond in 2 ms, 5 ms, and 40 ms — the last being a node doing compaction, or one across a slower link.
- W = 1 → 2 ms
- W = 2 → 5 ms
- W = 3 → 40 ms
Requiring all three replicas means every write is held hostage by the slowest one, and in a cluster of hundreds of nodes there is always a slow one. That is why W = N is rare in practice, and why W = 2, R = 2 with N = 3 is the near-universal default: strong reads while tolerating one slow or dead replica on both paths.
The availability cost, computed
Take a per-node availability of 99.9% (each node is down about 43 minutes a month).
- W = 3: all three must be up. 0.999³ ≈ 99.70% — about 2.2 hours of write unavailability a month. Worse than a single node.
- W = 2: at least two of three up. That fails only if two or more are down simultaneously, which is roughly 3 × 0.001² ≈ 3 × 10⁻⁶ — about 99.9997%, or 8 seconds a month.
The gap between 2.2 hours and 8 seconds is produced by changing one integer. This is the arithmetic that makes the quorum conversation concrete instead of theoretical, and it is worth doing on the whiteboard.
Tuning per operation
The dial can be set per call, which is the practically important point:
"A write to a user's session uses W = 1 because losing one is acceptable and latency matters. A write to their account settings uses W = 2, R = 2 so the next read is guaranteed to see it. Same cluster, same data model, different call parameters."
What R + W > N does not give you
It gives strong reads under normal operation, not linearizability in the full sense. Two concurrent writes to the same key can both satisfy W = 2 with different values, and the quorum rule provides no way to decide which is "later" — that is conflict resolution, below. And during a partition, a sloppy quorum (Failure handling) may accept writes on nodes outside the preference list, which breaks the overlap guarantee in exchange for availability. Naming both limitations unprompted is a strong senior signal.
Conflict resolution: why conflicts are unavoidable
Two clients write the same key at the same moment, each satisfying its quorum. Both writes succeeded. The replicas now disagree, and something must decide what the value is.
In a leaderless system there is no single point that serialises writes, so concurrent writes to one key are ordinary rather than exceptional. Under a partition it is guaranteed: two halves of the cluster each accept a write, and when the partition heals both versions exist.
The system has three options: pick one, keep both, or ask the application.
Last-write-wins, and its data-loss problem
Attach a timestamp to every write and keep the one with the larger timestamp. It is trivial to implement, needs no extra storage, and it is what many production systems do by default.
The problem is that it silently discards a successful write. Two users edit a shared document field at the same second; one write is acknowledged, then thrown away during repair. The user who made it saw success and their change is gone, with no error anywhere.
The problem is worse than it looks, because clocks are wrong. Timestamps come from machine clocks, which drift and are corrected by the Network Time Protocol (NTP). Skew of tens to hundreds of milliseconds between machines in a datacentre is ordinary. So:
- Write A happens at real time T on a node whose clock is 200 ms fast.
- Write B happens at real time T + 50 ms on a node with an accurate clock.
- B is genuinely later, but A's recorded timestamp is higher, so A wins.
The later write loses. And if a clock jumps backwards during an NTP correction, a node can write timestamps that are permanently below its neighbours', so its writes are systematically discarded until the clock catches up.
When last-write-wins is acceptable: the value is immutable, or the whole value is overwritten each time by a single writer, or losing a concurrent update is genuinely harmless (a "last seen at" timestamp, a cached rendering). Say the condition rather than the label.
Vector clocks: detecting concurrency instead of guessing at time
A vector clock is a list of (node_id, counter) pairs carried with each version. When a node coordinates a write, it increments its own counter.
Given two versions X and Y:
- If every counter in X is greater than or equal to the matching counter in Y, X descends from Y — X is genuinely newer and Y can be discarded.
- If neither dominates — X is ahead on one node's counter and Y is ahead on another's — the writes were concurrent, and neither can be declared later.
Worked. Client 1 writes and gets [(A,1)]. Client 2 reads that and writes back, producing [(A,2)] — a clean descendant, no conflict. Now suppose instead that Client 1 writes via node A producing [(A,1)] and, at the same moment, Client 2 writes via node B producing [(B,1)]. Neither dominates. The system knows, with certainty and without consulting any clock, that these were concurrent.
What it buys: the system never silently discards a concurrent write, because it can tell concurrency from succession.
What it costs: the clock grows a pair for every node that coordinates a write to the key, so long-lived keys accumulate entries. Implementations truncate the oldest entries past a threshold, which reintroduces a small chance of a false conflict. And the version must be carried by the client and returned on the next write, which pushes complexity into every client library.
Pushing resolution to the client
Having detected a conflict, the store returns all concurrent versions (siblings) to the client on read. The application merges them, because only the application knows what merging means for this data.
The canonical example is a shopping cart: two concurrent versions, one with items {A, B} and one with {A, C}. The store cannot know the right answer. The application does — take the union, {A, B, C}. The user may see a removed item reappear, which is a far better failure than losing an added one.
Also worth naming: conflict-free replicated data types (CRDTs) are data structures — counters, sets, ordered lists — whose merge function is defined mathematically so that concurrent versions always converge without application involvement. They are the principled version of this idea and they carry real overhead in metadata, which is why they appear in collaborative editing (Section 17) rather than in every store.