F1 Consistent Hashing & Partitioning
Loading learning experience...
Lecture transcript
Read the narration for F1: Consistent Hashing & Partitioning
From estimating scale to choosing a shard function
Dr. Wei: Last time: estimates drive architecture; today: partitioning drives everything else. Once you have a rough sense of how many requests you serve, it forces you to decide how the data and traffic will be split across machines.
Sam: So the estimate is not just a number, it changes the design choices we even consider?
Dr. Wei: Exactly. At large scale, you cannot treat storage and routing as an afterthought. The key question becomes: how do we choose a shard function so each key lands somewhere predictable, the load stays balanced, and when machines join or leave, only a small fraction of keys have to move.
Sam: And that shard function has to be cheap enough to run on every request, right, otherwise routing becomes the bottleneck.
Why naive modular hashing melts caches
Dr. Wei: When we partition data across servers, we want two things at once: spread the load fairly, and keep most keys on the same server over time. The problem is that some very simple hashing choices make stability extremely fragile.
Sam: Is the fragile part mainly about what happens when N changes, like scaling up by one box or losing one box?
Dr. Wei: Yes. Quick prediction: if N equals 5 and we add one node so N becomes 6, what fraction of keys do you think stay on the same shard? A key stays put only when h of k mod 5 equals h of k mod 6, meaning both remainders match. Under a uniform hash, that happens for about one out of six keys, so roughly five out of six keys move.
Sam: So the system basically behaves like a full reshuffle, which means whatever you cached is suddenly on the wrong machine.
Dr. Wei: That is the core motivation for consistent hashing: avoid remapping almost everything when you add a server, remove a server, or replace a failed node. In the next moments, we will see why the naive modular approach causes a cache avalanche and pushes load to the database.
The core idea: a hash ring in a fixed space
Dr. Wei: Imagine we agree on one big, fixed address space for every key and every server, and we treat that space like a circle so the end wraps back to the start.
Sam: So the circle is just a way to avoid edge cases, like what happens after the largest hash value?
Dr. Wei: We run the same hash function on a server name and on a data key, and each one lands at some position on that ring.
Sam: And by keeping the space fixed, we are decoupling placement from how many servers exist at the moment.
Dr. Wei: To decide who owns a key, start at the key’s position and walk clockwise until you hit the first server; that server is responsible. The key idea is that the ring itself stays fixed, and only the server positions change when servers come and go.
A concrete ring example: who owns which keys?
Dr. Wei: Let’s make consistent hashing feel concrete by using a tiny ring where both servers and keys have specific numeric positions.
Dr. Wei: Here’s the ring: servers sit at positions 10, 30, and 70 on a mod 100 circle. And we have four keys at positions 5, 12, 68, and 99.
Dr. Wei: Sam, before I tell you the rule, try to assign an owner to each key. For each key, look clockwise and pick the first server you encounter.
Sam: Okay: key 5 should hit 10, key 12 should hit 30, key 68 should hit 70, and key 99 wraps around and hits 10.
Dr. Wei: That’s right. And you can see the intuition: each server owns the arc just before it, and only the local clockwise ordering matters for deciding who owns a key.
Why it is called "consistent": minimal key movement
Dr. Wei: Consistent hashing is designed so that when the set of nodes changes, most keys keep the same destination, and only a small, predictable slice of key to node assignments needs to change.
Sam: When people say it moves about K over N keys, is that an expectation over random placement on the ring, not a hard guarantee for one specific cluster snapshot?
Dr. Wei: Yes, it is an expectation under uniform hashing. To keep the notation straight, let N be the number of nodes before the change and K be the total number of keys. If you add one node, going from N to N plus 1, then about K over N plus 1 keys move because the new node takes responsibility for roughly one over N plus 1 of the ring, which is close to K over N when N is large.
Dr. Wei: If instead you remove one node, going from N to N minus 1, then roughly K over N keys move, specifically the keys that were assigned to the removed node, which get reassigned to the next clockwise node.
Dr. Wei: And if you use replication with a factor of R, you can think in terms of assignments moved: the number of key to replica placements that change is multiplied by about R, even though the fraction of the ring affected by the membership change stays small.
Virtual nodes: fix uneven load on small $N$
Dr. Wei: When a cluster is small, consistent hashing can look deceptively simple, but the load can be surprisingly uneven. A few nodes means each one covers a large stretch of the hash ring, and those stretches are not guaranteed to be equal.
Sam: So even with a good hash, randomness can still give you a bad draw where one node gets a much bigger arc than the others?
Dr. Wei: If one machine happens to own a much bigger arc than the others, it will receive a much larger fraction of keys and requests. That creates high variance in storage and traffic, even though every machine is healthy and the hash function is doing its job.
Sam: And that variance shows up as one machine hitting limits earlier, even though average utilization across the fleet looks fine.
Dr. Wei: Virtual nodes are the standard fix: instead of giving each physical machine one position on the ring, we give it many positions, each acting like an independent tiny slice. With, say, on the order of one hundred to two hundred positions per machine, the slices average out and the total load per machine becomes much more even.
Dr. Wei: The cost is bookkeeping: the routing information grows because there are more ring positions to track. In practice, the balance improvement is usually worth the larger routing table, especially when you have a small number of physical nodes.
Vnodes in practice: from lumpy to smooth
Dr. Wei: Let’s ground consistent hashing in something practical: how evenly traffic and data end up spread across machines when you add or remove capacity.
Sam: And the point is that with enough virtual nodes, the distribution tightens up so one server is less likely to become the obvious hotspot, right?
Dr. Wei: In the real world, perfect balance is rare. Small random differences can create noticeable hotspots, and those hotspots turn into paging, throttling, and emergency rebalancing work.
Dr. Wei: Virtual nodes are the simple trick that makes the distribution smoother: instead of each server owning one big contiguous chunk of the hash ring, each server owns many small slices, and the variance averages out.
Replication on the ring: $R$ successors for durability
Dr. Wei: To make the system durable, we do not store a key on just one node. Instead, we copy the same data to several nodes placed around the ring, so a single failure does not make the data disappear.
Sam: When you say the first R distinct nodes clockwise, is that also how you avoid putting two replicas on the same physical machine when you have virtual nodes?
Dr. Wei: The idea is simple: for a key k, start at the position given by the hash of k, then walk clockwise and choose the first R distinct nodes you encounter. Those R nodes are the replica owners for that key.
Dr. Wei: A quick note on symbols so the quorum language is not confusing: here R is the replication factor, meaning how many copies we keep, and W is the write quorum, meaning how many replica acknowledgements we wait for on a write. In some Dynamo style descriptions, the replication factor is called N, which can be easy to mix up with the number of nodes in the whole cluster.
Hot keys: bounded-load consistent hashing
Dr. Wei: Hot keys are the classic failure mode of naive partitioning: one popular key can concentrate traffic and storage on a single node, even if the rest of the cluster is mostly idle.
Sam: So even with a perfectly balanced ring, a single key can dominate because popularity is not uniform, and the ring only balances the count of keys, not the request rate per key.
Dr. Wei: Bounded-load consistent hashing is a strategy to keep any single node from drifting far above the average load. Think of it as enforcing a soft cap so load stays close to the mean instead of spiking unpredictably.
Sam: Is the main cost here that you need some shared view of load, so placement depends on runtime measurements instead of just a pure function of the key?
Dr. Wei: A common approach is multiple choices: for each key, you generate a small number of candidate nodes using independent hashes, then assign the key to the least-loaded candidate. In the ring and vnode setting, the ring still provides the candidates, but you need a feedback loop that tracks observed load to decide among them, and you can apply that decision per replica or just for the primary. Under reasonable capacity headroom, this keeps each node within about one plus epsilon times the average load, which is exactly what we want for hot keys.
Meta Memcache: why adding a server does not nuke the cache
Dr. Wei: When you run a cache at huge scale, the hard part is not storing values, it is keeping the system stable as machines come and go. If every change reshuffled where keys live, you would trigger a wave of cache misses and sudden load on the backing databases. Consistent hashing is the idea that prevents that kind of cache meltdown.
Sam: So the operational metric I should keep in my head is the moved fraction, like about one over N keys changing owners when we add capacity?
Dr. Wei: In this example, think about Meta's Memcache fleet: there are thousands of cache servers, and a lightweight router decides which server should hold a given key. The goal is simple: if you add one more server, only a small fraction of keys should move, so most of the cache stays warm. That is what people mean when they say adding capacity does not nuke the cache.
Sam: And when a node fails, you are trading a smaller remap for a temporary load spike on the neighbors that pick up that arc, so warmup and throttling still matter.
Dr. Wei: Consistent hashing also helps when something goes wrong. If a server fails, you want to reroute just the affected keys, and you often pair that with warmup tactics so the replacement path does not get overwhelmed. The big takeaway is that the partitioning scheme is designed to limit key movement during both scaling events and failures.
Alternatives: Jump, Rendezvous, and Maglev
Dr. Wei: Not every system needs a hash ring; the real question is how you trade off churn, key movement, and lookup cost.
Sam: Is it fair to say the ring is a good default, but if lookup cost or state management dominates, you might pick a different algorithm?
Dr. Wei: If you want very fast lookups and you mostly add capacity, Jump hashing gives constant time placement with good balance, but it is awkward when you need to remove arbitrary nodes.
Dr. Wei: If you expect both adds and removals and want a simple story about disruption, Rendezvous hashing scores each node and picks the winner; it avoids ring state, but each lookup considers all nodes, so cost grows with cluster size.
Dr. Wei: If your priority is extremely high queries per second, Maglev-style approaches precompute a table so the lookup is just an index, and the main trade off becomes how you rebuild that table when membership changes and how much movement you allow.
Interview checklist: what makes this $IC4$ vs $IC6$
Dr. Wei: To stand out in a consistent hashing interview, you want to explain not only what the ring is, but what it means for day to day operations when the system changes.
Sam: So the bar is not memorizing the ring rule, it is being able to reason about movement, balance, and failure behavior under realistic conditions.
Dr. Wei: A useful way to frame your answer is by seniority level: what would be expected from an IC4 answer versus an IC6 answer, and what extra depth shows stronger ownership.
Sam: So for IC6, I should proactively call out the ugly parts: membership churn, uneven load without vnodes, replica placement quirks, and hot-key failure paths.
Dr. Wei: So as you listen to yourself explain the design, check that you are covering behavior and tradeoffs: how many keys move when nodes change, how expensive lookups and rebalancing are, and what happens under failures, hot keys, and replication choices.
Exit ticket: quantify disruption and explain the failure path
Dr. Wei: To wrap up, we are going to test two things you should be able to say clearly: a quick disruption estimate, and a clear failure story. Both are about what happens to data and traffic when the cluster membership changes. Keep it simple and explain it like you are on a whiteboard.
Sam: So I should be able to quantify how much data moves when we add a node, and then describe what operational problem shows up first if the design is missing key safeguards.
Dr. Wei: Exactly. For the number, you are estimating the fraction of keys that change owners when the cluster grows by one node, assuming the hash space is evenly shared. For the story, you are tracing a realistic failure path: which component gets overloaded first, what symptom you see, and why missing virtual nodes or hot key controls makes it cascade.
Thank you for watching!
Thanks for watching. Subscribe and share if you found this useful—see you next time!