Ch.12: System Design: Key-Value Store Across 100 Machines

Outline

Transcript

0:00 Welcome to Learning Podcasts. System Design for Backend Engineers: Building a Key-Value Store Across a Hundred Machines. You call put, you call get, the value comes back. The whole API is three operations. And the work is in the cluster underneath. Today we walk every layer that makes those three calls behave at scale. Start with the finished diagram. A ring of nodes wraps the canvas, with virtual node positions scattered around it. Each key lands at a position on the ring and lives on the next three nodes clockwise.

0:32 A coordinator forwards the request, gossip arrows connect every pair of nodes for membership, a Merkle tree sketch sits to the side for repair, and a hinted-handoff buffer catches writes when a replica is offline. That is the whole map. The chapter visits each region in order. Right. The map will return at the end, and the chapter walks each region between now and then. Single-machine SQL hits a wall. The working set stops fitting in one box, writes need to land in many places at once, and the operations team starts writing custom shard scripts.

1:06 I have shipped that cluster. The one where the sharding script becomes its own oncall service. Yep. At that point the choice is to queue the work and reduce the write rate, or to distribute the storage itself. This chapter is the distribute path. And the distribute path turns into a key-value store almost by gravity. The system that scales horizontally without joins is also the system you reach for first. Right. So the question is what the cluster has to do, not what to call it. The API is small on purpose.

1:40 Three calls. Put writes a value at a key. Get reads it back. Delete removes it. No joins. No range queries. No transactions across keys. And no secondary indexes by default, although some systems bolt them on. The point is the surface stays narrow so the cluster underneath has freedom. That is the trade. You give up the shape of a relational query, and you get linear scale-out in return. Let us pin down the functional requirements before we pick any technology. The system has to store and retrieve a value by key.

2:14 Replicate every key for durability. Tolerate node loss without losing data. And allow online resize. Capacity grows by adding nodes; the cluster has to absorb them without downtime. Optional secondary indexes if the team wants them, optional time-to-live so keys can expire, but those are not the core. The core is durable, replicated, resizable storage of a value at a key. Right. Everything else is a feature on top. Non-functional requirements set the harder bar. Latency in single-digit milliseconds at high percentile, not at the median.

2:52 Availability across a region, ideally a continent. Tunable consistency from fast and loose to strong reads. Linear scale-out as nodes are added. And cost discipline because storage and IOPS are real money at this scale. And the consistency knob is the one that surprises new operators. A key-value store does not have one consistency story; it has whatever the caller asked for. Right. The cluster lets the application decide per call. Now the design questions. Where does a key live? The first answer every engineer reaches for is hash modulo N.

3:29 Hash the key, mod the result by the number of nodes, that is the owner. Which works fine until you add a node. Then nearly every key changes owner because the modulus changed. At a hundred nodes, adding one means moving roughly 99 percent of the keys. How does the cluster even stay online during a resize like that? It mostly does not. You either run a maintenance window or you slice the move and rate-limit it. The operations team is blocked either way. At larger scale it is a deterministic disaster.

4:03 Right. Resize is the move that exposes it. We need a placement scheme that survives resize. Consistent hashing fixes the resize cost. Hash every node onto the same ring of positions. Hash the key onto the ring. The next node clockwise from that position owns it. Add a node and it claims the slice between its position and the previous owner. Only that slice moves; everything else stays put. Which means resize moves order one over N keys, not all of them. Removing a node is the same move in reverse.

4:36 Oh, interesting. I have shipped clusters with virtual nodes for years without really understanding why the spread was uneven. The virtual count was the answer. Each physical node owns many virtual positions, so the slices average out and a heavier node can simply own more virtuals. This is the move every distributed key-value store builds on. Now durability. Every key needs more than one copy. The replication factor, usually written as N, says how many. And the placement is mechanical. Walk clockwise from the owner and put the next two replicas on the next two distinct physical nodes.

5:15 Three is the common default. Three machines, two of which can be lost before a key is gone. The replication and consistency tradeoff was its own episode earlier in the series; this is where we apply it. The consistency knob is the next move. A write waits for W acknowledgments before it returns. A read waits for R responses. And the rule of thumb is W plus R is greater than N. If the write set and the read set overlap by at least one node, every read sees the latest write. With N equal to three, you can pick W equals 1 and R equals 1.

5:52 That is fast and loose; reads might miss a recent write. So which combination do most teams pick by default? Or is that a per-workload call from day one? Most teams default to W equals two and R equals two with N equal to three. That is the safe knob. The fast-and-loose options come out for analytics or pure read-heavy paths where stale reads are tolerable. Wow. So the same cluster carries different guarantees just by adjusting two integers per call. Exactly. The application picks per call. Same cluster, different guarantees.

6:24 Any node can be the coordinator for a request. The client SDK picks one, often by, Hashing the key locally, and that node forwards reads and writes to the right replicas. Which means there is no leader. The fan-out shape is the same one the caching episode used at the read layer, just spread across the ring. Sloppy quorum extends that. If a replica is down, an alive node accepts the write on its behalf. Hinted handoff stores the write with a note saying which replica it really belongs to. And when the original replica returns, the alive node replays the buffered writes.

6:59 Right. The cluster keeps writing through partial failure. Availability does not flinch. Two writes for the same key from two coordinators, in the same window, is the normal failure mode. The cluster has to pick a winner. Last-writer-wins is the cheap answer. Compare wall-clock timestamps; the later one wins. And the test that was supposed to catch this race... ...passed. Wall clocks drift. The later write by the clock might be the earlier write in real time. Vector clocks are the correct answer.

7:28 Each write carries a per-node counter, and the cluster can tell whether two versions are causally related or genuinely concurrent. Concurrent versions are surfaced to the application; the caller decides how to merge. So when is last-writer-wins still defensible in production? When the data is naturally idempotent. A monotonic counter, a status field whose updates only move forward, a cache where stale is recoverable. For shopping carts or anything financial, vector clocks earn their cost. Cheap and lossy, or correct and expensive.

7:57 Right. Pick at design time, not in production. Membership is the next problem. Who is in the cluster right now, and who is alive? Gossip handles both. Each node periodically picks a small random set of peers and exchanges state digests. Membership lists, version vectors, recent failures. After a few rounds, every node has a roughly current view of the cluster. The chatter scales as log N because the rumor spreads exponentially. And there is no leader, no central registry to crash. The cluster knows itself by talking to itself.

8:34 After a partition heals, replicas have drifted. Two copies of the same key range now have different contents. The brute-force fix is to scan every key on both nodes and ship the differences. At a billion keys, that is a non-starter. Merkle trees fix it. Each replica summarizes its keys into a tree of hashes. Comparing two trees finds the diverging ranges in log N comparisons instead of scanning every key. Oh wow. And only the diverging ranges are exchanged. The repair pass becomes cheap enough to run continuously in the background.

9:12 One more failure mode. The cluster scales linearly only when traffic spreads evenly across the ring. Which it does not, the moment one row gets a hundred times the read volume of every other row. The hot-key problem. And virtual nodes do not save you. Virtuals balance ownership, not request rate. A single celebrity row lands all of its reads on whichever node owns it. The schema design is what fixes it. Compose the partition key from something that fans out, or pull the hot key into a dedicated cache layer.

9:50 Right. The ring can do a lot. It cannot do schema for you. The real systems agree on more than they disagree. DynamoDB hides the ring behind a managed API but exposes the consistency knob. Cassandra exposes everything: ring topology, replication factor, consistency per call. Riak ships with vector clocks turned on by default; many shops still flip on last-writer-wins for simplicity, knowing the cost. What does it actually look like when two concurrent writes race in production? A small percentage of records carry conflicting versions until anti-entropy sweeps through.

10:25 Last-writer-wins hides them; vector clocks surface siblings. Right. So the choice is whether the application sees the conflict at all. Exactly. The interesting move is what they share. Consistent hashing on a ring, replication factor N, tunable consistency, gossip for membership, anti-entropy for repair. The differences are the last 10 percent. Same diagram as the start. The ring still wraps the canvas. The virtual node positions still scatter around it. Every key still lives on three nodes clockwise.

10:56 But now the coordinator earns its arrow. The gossip arrows earn theirs. The Merkle tree on the side is the repair pass that runs while everything else runs. It really does. And the consistency knob is what the application picks per call, on top of the same physical layout. This is going to feel slow on a strict-consistency read across regions. That is the price. Right. The read path is calm because the layout, the replication, the gossip, and the repair all paid their cost up front. That is the capstone.

11:27 Consistent hashing for placement. Replication for durability. Quorums for tunable consistency. Gossip for membership. Anti-entropy for repair. 5 moves. Every distributed key-value store you have used picks them in some combination, and the differences are at the edges. Next up, the distributed job scheduler. How a cluster runs millions of jobs at the right time, exactly once, without losing any of them. Thanks for listening to Learning Podcasts.