CONCEPT Cited by 4 sources
Consistent hashing¶
Consistent hashing maps keys to buckets (shards,
servers, cohorts) in a way that minimises re-mapping when the
bucket set changes. Adding or removing a bucket moves only a
1/N fraction of keys — not the whole keyspace as with
naive hash(key) % N.
Originally an Akamai / Karger et al. construction (1997) for distributed web caching. The canonical implementation places buckets at hashed positions on a ring; each key maps to the next bucket clockwise. Modern variants use bounded-load, rendezvous (HRW), or jump consistent hashing.
Why it matters beyond sharding¶
The same "adding a bucket keeps most keys where they were" property is load-bearing in two very different contexts:
- Sharding / partitioning — consistent hashing means scaling a cache or database cluster doesn't evict the whole cache or rehash the whole dataset.
- Percentage rollouts — with consistent hashing on a context attribute, raising the rollout from 5% to 10% keeps every already-rolled-out user rolled out; it just admits a new 5% slice.
This second use is the one called out directly in Cloudflare Flagship:
"Rollouts use consistent hashing on the specified context attribute. The same attribute value (userId, for example) always hashes to the same bucket, so they won't flip between variations across requests. You can ramp from 5% to 10% to 50% to 100% of users, so those who were already in the rollout stay in it."
The property the post needs is weaker than the full
distributed-sharding ring — it just needs hash(userId) %
100 to be stable across calls. But "consistent hashing" is
the industry-standard name for the whole property cluster.
Core properties¶
- Stable mapping —
hash(k)is a pure function; the same key always produces the same position. - Minimal disruption on resize — adding or removing one bucket moves ~1/N of keys.
- Balanced load — with a good hash, keys distribute approximately uniformly across buckets (virtual nodes / multiple hash positions per physical bucket improve this).
Balancing math: virtual nodes, weighting, and the collision ceiling¶
Cloudflare's 2026 PBR post derives the load-balancing quality of a ring from first principles — a useful reference for how many virtual nodes ("hashes per server") to use:
- One hash per server is bad. For the fractional arc owned by
one of N servers:
Exp = 1/N,SD = (1/N)·√((N−1)/(N+1)). The coefficient of variationCV = SD/Exp = √((N−1)/(N+1))≈ 99 % at N=100 — some servers do nearly double their fair share. - k hashes per server evens it out (law of large numbers).
The exact formula:
SD_k = √((k+1)/(N(kN+1)) − 1/N²), soCV_k = √((N−1)/(N·k+1)). At N=100, going from 1 → 160 hashes/server (the NGINX/Pingora default) drops CV ~99 % → ~8 %. - Diminishing returns. Each order-of-magnitude increase in
kbuys progressively less; the last 90,000 of 100,000 hashes bought only ~0.7 % improvement. - Weighting = more hashes. To make server S1 take
w×the traffic of S2, give itw×the hashes (H1 = w·H2) — the ketama scheme. Weight by disk (cache), CPU, or GPU depending on the workload. The constant base-160 factor still sets a minimum error floor for the lowest-weight servers. - 32-bit collisions cap the benefit. The continuous-ring math assumes no collisions, but real 32-bit hashes collide per the birthday paradox; collisions randomly drop a hash's contribution. Cloudflare's simulations for 2048-server DCs showed error rate increasing between 10,000 and 100,000 hashes/server — so past a point, more hashes is worse, not just wasteful of RAM.
Common use cases on this wiki¶
- Distributed caches (memcached clients, Redis Cluster).
- Request routing to stateful shards ( Durable Objects route by hashed ID → single-writer instance).
- Partitioned streaming (Kafka partitioning by key).
- Percentage rollouts in feature flags — the 2026-04-17 Flagship launch made this a first-class on-wiki instance.
Seen in¶
- sources/2023-02-22-highscalability-consistent-hashing-algorithm
— canonical distributed-systems primer for this concept
page. High Scalability / systemdesign.one (2023-02-22)
walks the full partitioning-problem taxonomy (random /
single-global / key-range /
static-hash /
consistent) and derives consistent hashing plus its two
Google-paper variants (multi-probe 2015, bounded-load
2017). Names the canonical production lineage:
Memcached Ketama clients,
Amazon Dynamo /
DynamoDB / Apache Cassandra / Riak,
Vimeo LB, Netflix Open Connect CDN, Discord server-to-node
mapping, HAProxy bounded-load. Adds the
k/Naverage movement property, the BST-over-ring-positions implementation with O(log n) operations, the readers-writer lock for concurrent membership change, and the non-cryptographic hash recommendation (MurmurHash / xxHash / MetroHash / SipHash1-3 over MD5 / SHA-1 / SHA-256). - sources/2026-04-17-cloudflare-introducing-flagship-feature-flags-built-for-the-age-of-ai — canonical wiki instance of consistent hashing as the bucketing primitive under feature-flag percentage rollouts; the load-bearing property is "same attribute value always hashes to the same bucket" across ramps.
- sources/2026-09-18-cloudflare-saving-another-100tb-of-ram-with-math-and-rust
— canonical wiki instance of the balancing math (virtual
nodes, ketama weighting, coefficient of variation, 32-bit
collision ceiling) and of consistent hashing as cache-routing
by URL: Pingora Backend
Router keeps one copy of each file per data center by routing
cacheable requests through systems/pingora-ketama rings.
Cutting hashes/server by 90 % (justified by
CV_k+ collision analysis) plus a struct-packing memory fix reclaimed >100 TB of RAM globally. - — Zalando's PRAPI uses Consistent Hash Load Balancing (CHLB) at the Skipper ingress to partition 10M products across backend pods, so each pod owns a deterministic slice of the catalogue in its local Caffeine cache. Two upstream Skipper contributions come from this deployment: fixed-100-position placement (skipper#1712) reduces cache invalidation on pod scale events to 1/N; bounded-load (skipper#1769) caps per-pod traffic at 2× average so limited-edition product drops can't overload their ring owner. See bounded-load-consistent-hashing.
Related¶
- concepts/consistent-hashing — the ring topology the algorithm operates on.
- concepts/hot-path — the failure mode virtual nodes + bounded-load address.
- concepts/horizontal-sharding — the naive
hash mod Nscheme consistent hashing replaces. - percentage-rollout — uses consistent hashing to make ramps monotonic.
- concepts/feature-flag — the broader primitive consistent-hashed rollouts plug into.
- patterns/sharding — the original use-site.
- virtual-nodes-for-load-balancing — the canonical corrective for ring-arc variance.
- multi-probe-consistent-hashing — Google 2015 variant: linear space, no virtual nodes, slower lookup.
- bounded-load-consistent-hashing — the variant that caps hot-key overload at the ring owner.
- systems/libketama — 2007 Memcached-client consistent hashing library, one of the earliest open-source productisations.
- systems/pingora-ketama — Cloudflare's Rust ketama library; source of the CV math and the 8→6 byte struct-packing fix.
- systems/pingora-backend-router — routes cacheable requests by URL through ketama rings.
- systems/amazon-dynamo — canonical Dynamo paper deployment (consistent hashing + gossip BST).
- systems/cloudflare-flagship — rollout-ramp implementer.
- systems/skipper-proxy · systems/zalando-prapi — CHLB + bounded-load at Kubernetes-ingress altitude.
Merged aliases¶
consistent-caching-horizontal-scalehash-ring