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

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

  1. Hash each server identifier onto a circular space (for example 0 to 2⁶⁴−1).
  2. Hash each key onto the same space.
  3. 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):

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

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

References