Consistent Hashing
Consistent hashing assigns cache keys to nodes on a hash ring so that most assignments survive a change in the node set: adding or removing one node in a pool of N remaps roughly 1/N of the keys instead of nearly all of them. CDNs use it to shard caches and keep them warm across resizes.
Also known as Ketama hashing, Ring hash.
Full Explanation
Consistent hashing is a key-distribution algorithm. It assigns cache keys to a set of nodes so that most assignments survive a change in the node set. It decides where a key lives: which edge server or cache shard owns it. It does not decide what that node evicts when it fills up. An LRU policy handles that. Consistent hashing also does not decide which node is closest to the user. Karger et al. named it in 1997: "a consistent hash function is one which changes minimally as the range of the function changes" (Karger et al., 1997).
The practical payoff is one number. Place keys with plain hash(key) mod n and a membership change moves (n - 1) / n of them. Place them with consistent hashing instead, and it moves "roughly 1 / n (where n is the new number of nodes)" (AWS). That is the difference between a resize that empties the cache and one the cache barely notices. The method is also called ketama hashing, after the memcached implementation that popularised it. It is also called ring hash, after the data structure.
How it works
- Every node is hashed to points on one shared space. Implementations use a circle. ketama calls it the continuum: "Conceptually, these numbers are placed on a circle called the continuum" (ketama). Picture a clock face running from 0 to 2^32. The original construction instead maps buckets onto the unit interval and sends each item to the bucket point "closest" to it. The ring, with a clockwise walk, is the form that shipped (Karger et al., 1997).
- Each physical node gets many points, not one. These points are the virtual nodes. This is not a later optimisation. The paper already states that "we actually need to have more than one point in the unit interval associated with each bucket" (Karger et al., 1997). ketama hashes each server string to "several (100-200) unsigned ints" (ketama). nginx's implementation is compatible with Cache::Memcached::Fast, "with the ketama_points parameter set to 160" (nginx).
- To find a key's owner, hash the key to a point. Take the next point clockwise from it. The node behind that point owns the key. If the key hashes past the last point, wrap to the first (ketama).
- Weighting falls out of the same structure. A bigger node simply gets more points. In ketama, "the weightings are realised by adding more or less points to the continuum" (ketama).
- When a node joins or leaves, only the keys next to that node's points move. That is the paper's monotonicity property: "When a new bucket is added, the only items that move are those that are now closest to one of the new bucket's associated points. No items move between old buckets" (Karger et al., 1997). A key never migrates between two nodes that both stayed.
- The disruption is about 1/N of the keyspace. Envoy states it for request routing: "the addition or removal of one host from a set of N hosts will affect only 1/N requests" (Envoy).
- Lookup is cheap. With a balanced search tree over the ring segments, "a single hash computation takes O(log (C)) time". The paper also has a trick: split the space into equal-length segments with a tree per segment. That trick makes "the expected running time for a single hash computation" O(1) (Karger et al., 1997).
There is no single formal definition to appeal to. The paper says "we will not commit to a precise definition". Instead it formalises four properties: balance (each node gets a roughly equal share), monotonicity, spread (an object lands on few nodes even across disagreeing views of membership), and load (no node is asked to hold too much across those views). Informally it also asks for smoothness: "When a machine is added to or removed from the set of caches, the expected fraction of objects that must be moved to a new cache is the minimum needed to maintain a balanced load across the caches" (Karger et al., 1997).
Why it matters for a CDN
Edge pools are not a fixed set. They grow for a traffic spike, shrink afterwards, lose machines to hardware faults, and roll through releases. Under plain mod-n placement, each of those events remaps (n - 1) / n of the keys. That turns almost the whole tier into misses at the worst possible moment. AWS describes the consequence bluntly: "A large number of cache misses results in hits to the database, which is already overloaded due to the spike in traffic" (AWS). Substitute "origin" for "database" and that is a CDN outage. Consistent hashing confines the damage to roughly 1/N of keys. That is why nginx says switching it on "helps to achieve a higher cache hit ratio for caching servers" (nginx). The cache hit ratio survives the scaling event.
The second use is deduplication rather than survival. Because the mapping is stable and shared, every cache in a tier agrees on which peer owns an object. So a region stores one copy instead of one per machine. This is how a tiered topology picks a parent: Cloudflare's regional content hashing "groups data centers by region and uses consistent hashing to route content to designated upper-tier caches" (Cloudflare).
What CDNs do
- Cloudflare: Generic Global Tiered Cache applies regional content hashing. As a result, "Same content always routes to the same upper-tier data center within a region", and it "Eliminates redundant copies across multiple upper-tier caches". Cloudflare's illustration is a popular image that used to sit in 3-4 upper tiers and now sits in one, for "3-4x fewer cache MISSes". That is their example, not a guaranteed figure. Tiered Cache has to be switched on. The Generic Global topology is listed as Enterprise-only (changelog) (tiered cache).
- nginx: widely used as a CDN cache tier. Adding consistent to an upstream's hash key directive means "the ketama consistent hashing method will be used instead". Then "only a few keys will be remapped to different servers when a server is added to or removed from the group". Without that keyword, the same directive "may result in remapping most of the keys to different servers". Note that hashing is never nginx's default: "By default, requests are distributed between the servers using a weighted round-robin balancing method" (nginx).
- Envoy: used as an edge and sidecar proxy. Its ring hash balancer "implements consistent hashing to upstream hosts", a technique "also commonly known as 'Ketama' hashing". Maglev is available as "a drop in replacement for the ring hash load balancer any place in which consistent hashing is desired". But it trades stability for speed. Envoy warns "it is not as stable as ring hash when upstream hosts change". In their simulations, "approximately double the keys" moved on host removal (Envoy).
- AWS ElastiCache for Memcached: a shared cache tier behind an origin. "We recommend that you configure your clients to use consistent hashing". That is because a node change then moves only "roughly 1 / n" of keys. It is a client-side setting. The default varies: off by default in spymemcached and in the Memcached PHP library, on by default in the Enyim .NET client (AWS).
- Maglev: Google's network load balancer, in front of services rather than caching them. It is also the source of the table-based algorithm Envoy borrows. It "is also equipped with consistent hashing and connection tracking features, to minimize the negative impact of unexpected faults and failures on connection-oriented protocols". It "has been serving Google's traffic since 2008" (NSDI 16).
Watch out for
- Even distribution is not automatic. One point per node gives arbitrary arc sizes. Envoy notes that placing one entry for a weight-1 host and two for a weight-2 host "doesn't actually provide the desired 2:1 partitioning of the circle, however, since the computed hashes could be coincidentally very close to one another". So it multiplies the entries: 100 and 200 in its own example (Envoy). This is the whole reason for virtual nodes.
- 1/N is an average. It is a large number when N is small. AWS: "Scaling from 1 to 2 nodes results in 1/2 (50 percent) of the keys being moved, the worst case" (AWS). Ring membership is cheap to change only once the ring is wide.
- It does nothing for a single hot object. In any one view a key has exactly one owner. So the method spreads different keys, never copies of the same one. The 1997 paper does not use it for hot pages either. It pairs consistent hashing with a per-page tree of caches. The tree exists "to ensure that no cache has many 'children' asking it for a particular page" (Karger et al., 1997). Replicate or shield hot objects with a separate mechanism.
- A node that leaves hands its share of keys over cold, all at once. If those objects are popular, the refill arrives at the origin as a single burst and can stampede it. That is the same misses-become-origin-hits failure AWS describes. Drain and pre-warm instead of pulling a node at peak (AWS).
- Disagreeing views of membership do not break the mapping, but they do duplicate work. Consistent hashing was designed for exactly this case. On the Internet, "clients may have incompatible 'views' of which machines are available to replicate data". Its spread property means "references for a given object are directed only to a small number of caching machines". Small, not one: each divergent view can cost another copy and another origin fill. So converge membership quickly (Karger et al., 1997).
- The mapping is only as stable as its input. Envoy notes that ring hash, "like all hash-based load balancers, is only effective when protocol routing is used that specifies a value to hash on" (Envoy). nginx lets the key "contain text, variables, and their combinations" (nginx). That also lets you pick a variable that changes per request. Do that, and the mapping reshuffles constantly. Nothing stays cached.
Best practice
- Turn it on explicitly wherever a pool can resize: hash key consistent in an nginx upstream, ring hash or Maglev in Envoy, a consistent or Ketama distribution in Memcached clients. In each of those the default is something else.
- Give every node many points and size the ring deliberately: 100-200 per server as ketama does, 160 in nginx's ketama-compatible mode. Envoy's guidance is to set minimum_ring_size and maximum_ring_size. Envoy also says to monitor the min_hashes_per_host and max_hashes_per_host gauges "to ensure good distribution" (Envoy).
- Hash a stable identifier: the components of the cache key that identify the object. Never hash a per-request or per-client attribute.
- Change membership one node at a time. Let the ring settle, and warm the successor before cutover. That way, the roughly 1/N of keys that move do not all miss cold together.
- Watch hit ratio and per-node load side by side. One node hot on an otherwise balanced ring is a hot object, not a hashing problem. Replicate that object rather than retuning the hash.
Examples
Consistent hashing in Python:
import hashlib
from bisect import bisect_right
class ConsistentHash:
def __init__(self, nodes, vnodes=150):
self.ring = []
self.node_map = {}
for node in nodes:
for i in range(vnodes):
key = hashlib.md5(f"{node}:{i}".encode()).hexdigest()
self.ring.append(key)
self.node_map[key] = node
self.ring.sort()
def get_node(self, cache_key):
h = hashlib.md5(cache_key.encode()).hexdigest()
idx = bisect_right(self.ring, h) % len(self.ring)
return self.node_map[self.ring[idx]]
# Adding a 4th node only remaps ~25% of keys
ch = ConsistentHash(["edge-1", "edge-2", "edge-3"])
print(ch.get_node("/video/trailer.mp4")) # edge-2
Nginx upstream with consistent hashing:
upstream cdn_cache {
hash $request_uri consistent;
server cache-1.internal:8080;
server cache-2.internal:8080;
server cache-3.internal:8080;
# Adding cache-4 only remaps ~25% of requests
}
Frequently Asked Questions
Consistent hashing assigns cache keys to nodes on a hash ring so that most assignments survive a change in the node set: adding or removing one node in a pool of N remaps roughly 1/N of the keys instead of nearly all of them. CDNs use it to shard caches and keep them warm across resizes.
Consistent hashing in Python:
import hashlib
from bisect import bisect_right
class ConsistentHash:
def __init__(self, nodes, vnodes=150):
self.ring = []
self.node_map = {}
for node in nodes:
for i in range(vnodes):
key = hashlib.md5(f"{node}:{i}".encode()).hexdigest()
self.ring.append(key)
self.node_map[key] = node
self.ring.sort()
def get_node(self, cache_key):
h = hashlib.md5(cache_key.encode()).hexdigest()
idx = bisect_right(self.ring, h) % len(self.ring)
return self.node_map[self.ring[idx]]
# Adding a 4th node only remaps ~25% of keys
ch = ConsistentHash(["edge-1", "edge-2", "edge-3"])
print(ch.get_node("/video/trailer.mp4")) # edge-2
Nginx upstream with consistent hashing:
upstream cdn_cache {
hash $request_uri consistent;
server cache-1.internal:8080;
server cache-2.internal:8080;
server cache-3.internal:8080;
# Adding cache-4 only remaps ~25% of requests
}
Yes. Consistent Hashing is also known as Ketama hashing, Ring hash. Consistent hashing assigns cache keys to nodes on a hash ring so that most assignments survive a change in the node set: adding or removing one node in a pool of N remaps roughly 1/N of the keys instead of nearly all of them. CDNs use it to shard caches and keep them warm across resizes.