How Consistent Hashing Minimizes Data Movement During Cluster Rebalancing
Consistent hashing minimizes data movement during cluster rebalancing by mapping both nodes and keys onto a circular hash ring, ensuring that only keys located between a newly added or removed node and its immediate predecessor must migrate, rather than triggering a full cluster reshuffle.
In distributed systems architecture, rebalancing data across a changing cluster topology presents significant challenges for performance and availability. The ByteByteGoHq/system-design-101 repository explains how consistent hashing solves this problem in data/guides/consistent-hashing.md, providing a foundation used by production systems like Amazon DynamoDB, Apache Cassandra, and Akamai CDN to maintain stability during scale-out operations.
The Problem with Modulo-Based Hashing
Traditional distributed systems often use a naïve modulo approach: hash(key) % N, where N represents the total number of servers. When a new node joins or leaves the cluster, N changes, causing nearly every key to remap to a different server.
According to the guide in data/guides/consistent-hashing.md, this creates a "storm of misses" where almost the entire dataset must migrate to new locations. This results in massive network traffic, cache invalidation cascades, and service degradation during routine maintenance or scaling events.
Mapping to the Consistent Hash Ring
Consistent hashing solves this by placing both servers and keys on the same logical circular space. As implemented in the ByteByteGoHq/system-design-101 documentation:
- Server placement: Each server's identifier (such as IP address or name) is hashed to a point on the ring.
- Key placement: Each data key is hashed to another point on the same ring.
- Assignment rule: A key maps to the first server encountered when moving clockwise from the key's position.
This architecture decouples key location from the total server count. The deterministic clockwise assignment ensures that removing or adding infrastructure only affects the local neighborhood on the ring, not the global key space.
Localized Rebalancing During Topology Changes
When a new server joins the cluster, only the key range between the new node and its predecessor on the ring requires redistribution. The guide provides a concrete example in data/guides/consistent-hashing.md: when inserting server s4 between existing nodes, only key k0 migrates from s0 to s4, while all other keys remain stationary.
This localized property holds for removals as well. If a node fails, only the keys that mapped to that specific server need reassignment to the next node clockwise, leaving the remainder of the distributed dataset undisturbed.
Virtual Nodes and Fine-Grained Distribution
Production implementations enhance basic consistent hashing through virtual nodes (replicas). Each physical server manages multiple points on the ring rather than a single location.
The fraction of keys that move during rebalancing equals roughly 1/(number of virtual nodes per physical server). By assigning dozens or hundreds of virtual nodes to each physical machine, systems achieve two goals simultaneously:
- Minimal movement: Only a small percentage of data migrates during rebalancing.
- Even distribution: Load spreads uniformly across the cluster, preventing hot spots.
Go Implementation Example
While the ByteByteGoHq/system-design-101 repository focuses on conceptual explanations rather than concrete implementations, the following Go code illustrates the core mechanics of a consistent hash ring with virtual nodes:
// ConsistentHash implements a simple ring with virtual nodes.
type ConsistentHash struct {
// hashFn maps a string to a uint32 value.
hashFn func(data []byte) uint32
// sorted list of ring points (virtual node hashes).
points []uint32
// mapping from point to real server identifier.
servers map[uint32]string
// number of virtual nodes per physical server.
replicas int
}
// New creates a new ConsistentHash with the given number of replicas.
func New(replicas int, fn func([]byte) uint32) *ConsistentHash {
ch := &ConsistentHash{
replicas: replicas,
hashFn: fn,
servers: make(map[uint32]string),
}
if ch.hashFn == nil {
// use a default hash (FNV‑1a)
ch.hashFn = func(data []byte) uint32 {
h := uint32(2166136261)
for _, b := range data {
h ^= uint32(b)
h *= 16777619
}
return h
}
}
return ch
}
// AddServer inserts a new physical server into the ring.
func (ch *ConsistentHash) AddServer(serverID string) {
for i := 0; i < ch.replicas; i++ {
// Create a distinct virtual node label.
virtualID := fmt.Sprintf("%s#%d", serverID, i)
point := ch.hashFn([]byte(virtualID))
ch.points = append(ch.points, point)
ch.servers[point] = serverID
}
sort.Slice(ch.points, func(i, j int) bool { return ch.points[i] < ch.points[j] })
}
// GetServer returns the server responsible for the given key.
func (ch *ConsistentHash) GetServer(key string) string {
if len(ch.points) == 0 {
return ""
}
point := ch.hashFn([]byte(key))
// Binary search for the first point >= key hash.
idx := sort.Search(len(ch.points), func(i int) bool { return ch.points[i] >= point })
// Wrap around the ring.
if idx == len(ch.points) {
idx = 0
}
return ch.servers[ch.points[idx]]
}
When AddServer() executes, the binary search in GetServer() ensures that only keys hashing between the new virtual node and its predecessor remap to the new server. This matches the behavior documented in the ByteByteGoHq/system-design-101 guide where adding s4 only captures k0 from s0.
Summary
- Modulo hashing causes near-total data migration when cluster size changes, creating network storms and cache misses.
- Consistent hashing maps servers and keys to a circular ring, localizing the impact of topology changes to adjacent nodes only.
- Virtual nodes reduce the fraction of moving data to approximately
1/replicasper physical server while ensuring even load distribution. - Systems like DynamoDB and Cassandra leverage this algorithm to achieve horizontal scaling without full cluster rebalancing.
Frequently Asked Questions
What is the difference between consistent hashing and modulo hashing?
Modulo hashing computes hash(key) % N, making every key's location dependent on the total node count N. When N changes, virtually every key relocates, causing a full cluster reshuffle. Consistent hashing places nodes and keys on a fixed circular space where a key maps to the next clockwise node, ensuring that only keys between the changed node and its neighbor move.
How many keys move when adding a node in consistent hashing?
Only keys located in the arc between the new node's position and its immediate predecessor on the ring must migrate. With virtual nodes (where each physical server occupies multiple ring positions), the fraction of moving keys equals approximately 1/(number of virtual nodes per server). For a configuration with 150 virtual nodes per physical machine, roughly 0.67% of keys would relocate.
Why do systems use virtual nodes in consistent hashing?
Virtual nodes solve the uneven distribution problem that occurs when physical servers have random positions on the ring. By assigning dozens or hundreds of virtual nodes to each physical server, load balances uniformly across the cluster. Additionally, virtual nodes reduce the volume of data movement during rebalancing by spreading a physical server's departure across many small ranges rather than one large contiguous block.
Which production systems use consistent hashing?
According to the ByteByteGoHq/system-design-101 documentation, Amazon DynamoDB, Apache Cassandra, and Akamai CDN all implement consistent hashing to distribute data across distributed clusters. These systems require the ability to add or remove nodes without triggering massive data migration, making consistent hashing essential for their scalability and fault tolerance architectures.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →