Skip to main content
Create your own
Lesson illustration

Consistent Hashing for Fault-Tolerant Caching

Welcome! In our last lesson, we focused on what happens inside a single cache server, exploring eviction policies like LRU and LFU to manage its limited memory. Now, we're zooming out to address a fundamental challenge in large-scale systems: how to manage data across an entire fleet of cache servers. This is where your goal of designing scalable, distributed systems really comes into focus.

This lesson tackles a critical problem: how to add or remove servers from a distributed cache cluster without causing a system-wide meltdown. We will explore the design of a caching layer using consistent hashing, a clever and widely-used technique that ensures your system can scale gracefully. This is a cornerstone concept for building resilient, high-performance distributed caches and a frequent topic in system design interviews.

1. The Challenge of Distributing Data

Imagine you have a distributed cache with a handful of servers. The first question you need to answer is: for any given piece of data (identified by a key), which server should store it?

A simple and intuitive approach is to use a hash function and the modulo operator. You hash the key to get a number, and then you calculate hash(key) % N, where N is the number of servers. The result gives you the index of the server to use. This method works perfectly well as long as the number of servers, N, remains constant.

But what happens when your system needs to scale?

Master Consistent Hashing for System Design Interviews

The video "Master Consistent Hashing for System Design Interviews" does an excellent job of illustrating the problem with this simple approach.

Watch the first segment from the introduction. The key takeaway is that when you add or remove a server, the value of N changes. This change in the divisor completely alters the result of the modulo operation for almost every key, causing a massive reshuffling of data.

This massive reshuffle is a disaster for a caching layer. It leads to:

  • Mass Cache Invalidation: The vast majority of existing cached items are now effectively lost, as requests for them will be routed to the wrong servers.
  • Database Overload: With the cache suddenly "cold," a flood of requests hits the underlying database, potentially overwhelming it. This is often called a "thundering herd" problem.
  • High Latency: System performance degrades significantly until the cache is repopulated.

We need a distribution scheme where adding or removing a server only affects a small, predictable fraction of the keys. This is precisely the problem consistent hashing solves.

2. Consistent Hashing: The Hash Ring

Instead of mapping keys to a discrete list of servers, consistent hashing maps both servers and keys onto a conceptual circle, or hash ring.

Here’s how it works:

  1. The Hash Space: We imagine a large range of hash values (e.g., from 0 to ) arranged in a circle. This is the hash ring.
  2. Placing Servers: Each cache server is assigned a position on the ring by hashing its unique identifier (like its IP address or name).
  3. Placing Keys: To determine where a piece of data lives, you hash its key. This also gives it a position on the ring.
  4. Assigning Ownership: To find which server owns a key, you start at the key's position on the ring and move clockwise until you encounter the first server. That server is the owner of the key.
This infographic illustrates the core concepts of consistent hashing. Panel 2, "The Hash Ring," shows how both servers (A, B, C) and keys (K1, K2, K3, K4) are placed on the ring. Following the clockwise rule, K1 belongs to Server A, K2 and K3 to Server B, and K4 to Server C.

The Magic of Minimal Disruption

Now, let's see why this design is so powerful when the cluster size changes.

  • Adding a Node: Imagine we add a new server, "Node D," between "Node C" and "Node A" (as shown in Panel 3 of the infographic). Before, keys in that arc belonged to Node A. Now, they belong to Node D. Critically, keys belonging to Node B and Node C are completely unaffected. Only a small, localized slice of keys needs to be moved from Node A to the new Node D.

  • Removing a Node: If "Node B" is removed (Panel 4), the keys it owned (like "Key K") now simply continue their clockwise walk and are assigned to the next server, "Node C." Again, the impact is isolated. Keys on Node A and Node D are not affected.

Master Consistent Hashing for System Design Interviews

Let's return to the "Master Consistent Hashing" video for a clear explanation of this process.

Watch the section explaining the hash ring model and how it works. Then, see the demonstration of how this design minimizes disruption when adding or removing a server. The key insight is that the change is localized to a single segment of the ring.

3. Fixing Imbalance with Virtual Nodes

The basic hash ring model has one potential flaw: what if the servers land on the ring in a way that creates very unevenly sized segments? A server might get a huge arc, making it responsible for a disproportionate number of keys and creating a "hotspot," while another server gets a tiny arc and sits mostly idle.

The solution is to use virtual nodes (or vnodes).

Instead of mapping each physical server to a single point on the ring, we map it to hundreds or even thousands of points. We do this by hashing derived names like "Server-A-1", "Server-A-2", etc.

This diagram shows a ring with three physical servers (A, B, C), but many virtual nodes (A1, A2, B1, B2, B3, B4, etc.). Server B has more capacity and is assigned twice as many virtual nodes as Server A, thus handling roughly twice the load. This smooths out the key distribution and allows for weighting.

By using a large number of virtual nodes, we achieve two things:

  1. Better Load Balancing: The physical server's load is now the sum of many small, scattered segments on the ring. By the law of large numbers, this tends to average out, leading to a much more even distribution of keys across the physical servers.
  2. Weighted Distribution: We can assign more virtual nodes to more powerful servers, so they naturally take on a larger share of the workload, as shown with Server B in the image above.

When a physical server is added or removed, its many virtual nodes are added to or removed from the ring, ensuring that the rebalancing effort is spread thinly across many other servers instead of burdening a single neighbor.

4. Implementation in Practice

With your background in Go, it's helpful to see how this translates into code. You don't walk around a "circle" literally. Instead, the ring is typically implemented as a sorted array or slice containing the hash values of all the virtual nodes.

Building a Distributed Key-Value Store in Go

This article provides a clean, minimal implementation of consistent hashing in Go. It's a great practical example of the concepts we've discussed.

Read the section titled Partitioning with Consistent Hashing. Pay close attention to the Go code snippets for AddNode and Get. The AddNode function shows how a physical node is added to the ring as multiple virtual nodes. The slice of hashes is then sorted. The Get function demonstrates the lookup logic. It hashes the key and then uses sort.Search (a binary search) to efficiently find the first virtual node whose hash is greater than or equal to the key's hash. This is the O(log N) lookup we want.

Extending for Fault Tolerance: Replication

A distributed cache also needs to be fault-tolerant. What if the server holding a key goes down? Consistent hashing makes replication feel very natural.

To store N copies of a key (a replication factor of N), you first find the primary owner by walking clockwise. Then, you simply continue walking clockwise, skipping any virtual nodes that belong to the same physical server you've already seen, until you've found N unique physical servers.

Building a Distributed Key-Value Store in Go

The same article shows how to implement this replication strategy.

Now, read the next section, Ensuring Data Durability via Replication. The GetNodesForKey function is the key piece of code here. It finds the initial position and then scans clockwise, collecting unique physical nodes until the desired replication factor is met.

This strategy ensures that replicas are stored on different physical machines, protecting your data against single-server failures.

Conclusion

In this lesson, we designed a scalable and resilient caching layer by moving beyond simple data distribution schemes. You've seen how consistent hashing provides a robust solution to one of the fundamental problems in distributed systems: how to partition data in a way that gracefully handles dynamic changes in the number of servers.

Key Takeaways:

  • Simple Modulo Hashing is Fragile: It works for a fixed number of servers but causes catastrophic key reshuffling when servers are added or removed.
  • The Hash Ring is the Solution: By mapping both servers and keys to a circular hash space and using a "clockwise-first" rule, consistent hashing localizes the impact of scaling events.
  • Virtual Nodes Ensure Balance: Using multiple virtual nodes per physical server smooths out data distribution, prevents hotspots, and allows for server weighting.
  • Implementation is Efficient: The ring can be implemented with a sorted array and binary search, making key lookups fast (O(log M), where M is the number of virtual nodes).
  • Replication is a Natural Extension: Finding replicas is as simple as continuing the clockwise walk around the ring to find subsequent unique physical servers.

In our next module, we will shift our focus from data storage (databases and caches) to data processing. We'll begin by exploring asynchronous processing and event-driven architectures, another critical set of patterns for building scalable and resilient systems.

Can't find a good explanation? Sign up and we'll make it for you

Sign up