Consistent Hashing
- Pronunciation
- kun-SIS-tunt HASH-ing
In short
Consistent hashing spreads keys across a changing set of servers so that adding or removing a server moves only a small share of the keys.
What is consistent hashing?
The simple way to pick a server for a key is hash(key) % number_of_servers. It spreads keys evenly, but when a server is added or removed, the result changes for nearly every key, so caches suddenly miss and data has to be moved everywhere. Consistent hashing, introduced in a 1997 paper for distributing web caches, avoids that.
Imagine the space of hash values as a ring. Each server is placed on the ring at the position of its own hash, and each key belongs to the first server found by moving clockwise from the key's position. When a server joins, it takes over only the keys between it and its neighbor; when one leaves, only its keys move to the next server. On average, about one key in n moves, where n is the number of servers.
With only a few servers the ring can be uneven, so each physical server is usually given many positions, called virtual nodes, which spreads load more evenly and lets stronger machines take more of it. Consistent hashing is used by DynamoDB and Cassandra to place data, by CDNs and caching layers to choose servers, and by load balancers that need the same client to keep reaching the same backend.
A common misconception is that consistent hashing balances load perfectly. It balances keys, not traffic: one very popular key, a hot key, can still overload the server that owns it. Systems add replication, splitting of hot keys or extra caching for those cases, and some use alternatives such as rendezvous hashing.
Key takeaways
- Consistent hashing maps keys to servers on a hash ring.
- Adding or removing a server moves only about 1/n of the keys.
- Plain modulo hashing moves almost every key when servers change.
- Virtual nodes give each server many positions for an even spread.
- DynamoDB, Cassandra, CDNs and caches use it; hot keys still need care.
Example
import bisect, hashlib
def h(value):
return int(hashlib.md5(value.encode()).hexdigest(), 16)
class HashRing:
def __init__(self, servers, vnodes=100):
self.ring = sorted((h(f"{s}#{i}"), s) for s in servers for i in range(vnodes))
self.keys = [position for position, _ in self.ring]
def server_for(self, key):
i = bisect.bisect(self.keys, h(key)) % len(self.ring) # first server clockwise
return self.ring[i][1]
ring = HashRing(["cache-a", "cache-b", "cache-c"])
print(ring.server_for("user:42"))
# Adding "cache-d" later moves only about a quarter of the keys.Readers ask
Why not just use hash(key) modulo the number of servers?
Because changing the number of servers changes the result for almost every key, so nearly all cached or stored data ends up on the wrong server at once. Consistent hashing limits the move to a small fraction.
What are virtual nodes?
Extra positions on the hash ring for each physical server. Using many of them per server evens out the distribution of keys and lets more powerful servers take a bigger share.
Where is consistent hashing used?
In distributed databases such as Cassandra and DynamoDB to decide which nodes store which data, in distributed caches, in CDNs and in load balancers that keep a client on the same server.
See also
- Hash TableData Structures, p. 18A hash table is a data structure that stores key-value pairs and uses a hash function to find the value for any key in constant time on average.
- ShardingDatabases, p. 39Sharding is a way of scaling a database by splitting its data across several servers, called shards, so each one stores and handles only part of the total.
- PartitioningDatabases, p. 34Partitioning splits a large table into smaller partitions by a rule such as date ranges, so queries can skip irrelevant data and old data is easy to remove.
- CacheBackend & APIs, p. 8A cache is a fast, temporary storage layer that keeps copies of frequently used data so later requests can be served quickly without repeating slow work.
- Load BalancerDevOps & Cloud, p. 34A load balancer is a server or service that spreads incoming traffic across several backend servers so no single one is overloaded and the app stays available.
- Distributed SystemSoftware Architecture, p. 14A distributed system is a set of computers that work together over a network and appear to their users as a single system.
Spotted a mistake or something missing on this page?Suggest an edit