Consistent Hashing
When data or requests are spread across many servers (cache nodes, database shards, storage partitions), you need a rule that maps each key to a server. The naive rule, hash(key) % N, works until N changes: add or remove one server and almost every key maps somewhere new, causing a cache stampede or a massive data migration. Consistent hashing fixes this. When a server joins or leaves, only about 1/N of the keys move.
Introduced by Karger et al. in 1997 for web caching, consistent hashing now underpins distributed caches, Dynamo-style databases (Cassandra, DynamoDB, Riak), CDNs, load balancers with session affinity, and message routing. It's a staple of system design interviews and a practical tool whenever you shard by key.
TL;DR
- Modulo hashing (
hash % N) remaps ~all keys when N changes, which is catastrophic for caches and shards. - Consistent hashing places servers and keys on a hash ring; each key belongs to the next server clockwise.
- Adding or removing a server moves only the keys in its arc, roughly K/N keys.
- Virtual nodes (many ring positions per server) smooth out uneven load and allow weighting by capacity.
- Replicas go on the next distinct servers clockwise.
- Alternatives: rendezvous (HRW) hashing, jump consistent hash, Maglev, and bounded-load variants.
Quick Example
A minimal hash ring with virtual nodes in Python:
Core Concepts
Why Modulo Hashing Fails
With server = hash(key) % N, going from 4 to 5 servers changes the result for about 80% of keys. For a cache, that means a sudden flood of misses hitting the database. For a sharded store, it means moving most of the data. Consistent hashing bounds the disruption to the keys that must move.
The Hash Ring
- Hash each server identifier onto a circular space (for example 0 to 2⁶⁴−1).
- Hash each key onto the same space.
- A key is owned by the first server clockwise from its position.
When a server joins, it takes over the arc between its predecessor and itself, and only keys in that arc move (from one neighbor). When a server leaves, its keys move to its successor. Every other key keeps its owner.
Virtual Nodes
With few servers, random ring positions produce very uneven arcs: one server might own 40% of the keyspace. Virtual nodes place each physical server at many positions (tens to hundreds):
- Load spreads more evenly, since variance drops as virtual nodes increase.
- When a server leaves, its load is spread across many others rather than dumped on one neighbor.
- Weighting: give bigger servers more virtual nodes.
Cassandra historically used 256 vnodes per node, and newer versions default to fewer with smarter token allocation.
Replication on the Ring
For replication factor R, store each key on its owner plus the next R−1 distinct physical servers clockwise, skipping virtual nodes of servers already chosen. Rack- and zone-aware placement goes further, choosing replicas in different failure domains. It's the scheme described in Amazon's Dynamo paper and used by Cassandra and Riak. See database replication.
Alternatives
Where It's Used
- Distributed caches: Memcached clients (ketama), cache fleets behind proxies.
- Databases and key-value stores: Cassandra, DynamoDB, Riak, ScyllaDB partitioning.
- Load balancing with affinity: routing the same user or session to the same backend (Envoy ring hash and Maglev, Nginx
hash … consistent). - CDNs and object storage: mapping objects to edge or storage nodes.
- Stream processing and messaging: assigning keys to partitions or workers.
Best Practices
Use Enough Virtual Nodes
With only a handful of physical servers, use more virtual nodes (100+) to keep load balanced, and measure key distribution. Fewer vnodes reduce metadata and rebalancing overhead in large clusters.
Use a Good, Stable Hash Function
Choose a fast, well-distributed, non-cryptographic hash (MurmurHash, xxHash, BLAKE2 for simplicity). All clients must use exactly the same function and server naming, or they'll disagree about ownership.
Handle Hot Keys Separately
Consistent hashing balances keys, not traffic. A single celebrity key still lands on one server. Replicate hot keys, add a local cache layer, split them with key suffixes, or use bounded-load hashing.
Plan Rebalancing
When nodes join, data must actually move (for stores) or caches warm up (for caches). Throttle data streaming, and add capacity before you're at the limit.
Common Mistakes
Using hash % N for a Growing Cache Fleet
Scaling a Memcached tier from 10 to 12 nodes with modulo hashing invalidates most keys at once, and the database takes the full load. Use a consistent-hashing client.
Too Few Virtual Nodes
One ring position per server with 3 servers commonly yields one server owning half the keys. Increase vnodes, or use rendezvous or jump hashing.
Clients Disagreeing on Membership
If some clients see a node as down and others don't, they map the same key to different servers, which means duplicated cache entries or split writes. Distribute membership consistently (a config service, gossip, a service registry) and converge quickly.
FAQ
What problem does consistent hashing solve?
It minimizes how many keys change servers when servers are added or removed. With N servers, only about 1/N of keys move, instead of nearly all keys with modulo hashing. That keeps caches warm and data migrations small as clusters scale.
Why are virtual nodes needed?
Placing each server at a single random point on the ring produces uneven arc sizes, so some servers get far more keys than others. Many virtual nodes per server average out the randomness, spread load more evenly, and make failures redistribute load across many servers instead of one neighbor.
What's the difference between consistent hashing and rendezvous hashing?
Both achieve minimal key movement. Consistent hashing uses a sorted ring and binary search (O(log N) lookups, with vnodes for balance). Rendezvous hashing scores every server per key and picks the best (O(N) lookups, no ring, naturally balanced). Rendezvous is simpler for small server sets.
Does Redis Cluster use consistent hashing?
Not the ring variant. Redis Cluster uses a fixed set of 16,384 hash slots (CRC16 of the key modulo 16,384) assigned to nodes. Rebalancing moves slots between nodes, which achieves the same goal (limited movement) with explicit slot ownership. See Redis Cluster.
Related Topics
- System Design — Designing large-scale systems
- Database Sharding — Partitioning data across servers
- Caching — Distributed cache fleets
- Load Balancing — Affinity-based routing
- CAP Theorem — Consistency trade-offs for replicated data
- Redis Cluster — Slot-based sharding