Distributed Systems

Jump Consistent Hash: A Consistent Hash in Five Lines

Jump Consistent Hashing is an algorithm that maps a key to one of N shards so consistently that growing or shrinking the cluster moves almost no keys — and it does this in about five lines of code, with zero memory and no data structure at all. Published by John Lamping and Eric Veach at Google in 2014, it needs no ring, no virtual nodes, and no lookup table: just a fast pseudo-random sequence seeded by the key and a loop that runs about ln N times.

The trick is to replay history. For a fixed key, the algorithm simulates every past moment when a new shard was added and asks, "would this key have jumped to the new one?" — jumping to shard j+1 with probability exactly 1/(j+1). Whatever shard the key last landed on before N is its answer. The result is perfectly balanced load and provably minimal disruption, distilled into arithmetic that fits in a tweet.

  • Published2014 · Lamping & Veach (Google)
  • Lookup timeO(ln n) expected
  • MemoryO(1) — zero stored state
  • Keys moved per added bucket≈ 1/n (optimal)
  • Code size~5 lines, no allocation
  • LCG multiplier2862933555777941757

Interactive visualization

Press play, or step through manually. The visualization is yours to drive — try it before reading on.

Open visualization fullscreen ↗

Watch the 60-second explainer

A condensed visual walkthrough — narrated, captioned, under a minute.

The problem: resharding without a stampede

Suppose you spread a billion cache entries or database rows across N servers by computing hash(key) % N. It is fast and perfectly balanced — until you add a server. Change N to N+1 and the modulus shifts underneath almost every key: roughly N/(N+1) of all keys — over 99% at any large scale — now map to a different server. In a cache that means a near-total miss storm; in a shard map it means copying nearly the whole dataset. The whole point of consistent hashing, introduced by David Karger and colleagues at MIT in 1997 — work that soon became the technical foundation of the Akamai CDN — is to bound that churn: when the bucket count changes by one, only about 1/N of keys should move, which is the mathematical minimum.

Karger's classic solution places servers and keys on a hash ring and walks clockwise to the next server. It works, but to get an even split it needs many virtual nodes per server — often 100–200 — stored in a sorted structure that costs memory proportional to n·V and a binary search per lookup. Jump consistent hash asks a sharper question: can we get the same guarantee — minimal movement and perfect balance — with no ring, no virtual nodes, and no memory at all?

The five lines

Here is the entire algorithm, essentially as printed in the paper:

int32_t JumpConsistentHash(uint64_t key, int32_t num_buckets) {
  int64_t b = -1, j = 0;
  while (j < num_buckets) {
    b = j;
    key = key * 2862933555777941757ULL + 1;          // one LCG step
    j = (b + 1) * (double)(1LL << 31) / (double)((key >> 33) + 1);
  }
  return b;
}

It takes a 64-bit key (you pre-hash your real key into this) and a bucket count, and returns a bucket in [0, num_buckets). Two variables carry all the state: b is the bucket the key currently sits in, and j is the next bucket count at which it would jump to a brand-new bucket. The line key = key * 2862933555777941757ULL + 1 is a 64-bit linear-congruential generator (LCG) — a deterministic pseudo-random step. Because it is seeded by the key and reset on every call, the same key always replays the exact same sequence of "random" numbers, which is what makes the mapping consistent across machines and across time. The loop advances b forward, jump by jump, until the next jump would land at or beyond num_buckets; the last bucket it settled on is the answer. No allocation, no table, no lock — the function is pure arithmetic over two registers.

Why it lands where it does

Start from a requirement, not from the code. If we have k buckets and every one must hold an equal 1/k share, then when we go from k−1 buckets to k the newest bucket must capture exactly its 1/k — and it can only take keys from the others, never give any back. So a given key jumps to the new bucket with probability 1/k at each expansion. That single rule delivers both properties at once: balance (each bucket ends with 1/n) and minimal disruption (only the fraction that must move to fill the new bucket ever moves).

Now flip it around. If a key currently rests in bucket b, the chance it has not jumped by the time there are k buckets is the product of surviving each step: (b+1)/(b+2) · (b+2)/(b+3) · … · (k−1)/k, which telescopes neatly to P(stay) = (b+1)/k. The algorithm doesn't test each step one at a time; it draws the destination of the next jump in closed form. Treat r = ((key>>33)+1)/2^31 as a uniform random number in (0,1] — the top 31 bits of the freshly stirred LCG value. Setting j = floor((b+1)/r) gives P(j ≥ k) = P(r ≤ (b+1)/k) = (b+1)/k — exactly the survival curve above. So each loop iteration samples, in one shot, the bucket the key jumps to next, and the code's (b+1) · 2^31 / ((key>>33)+1) is just (b+1)/r written with integer bit tricks.

Cost: logarithmic time, no space

The loop runs once per jump the key makes as the cluster grows from 1 bucket to n. The expected number of jumps is the sum of the per-step jump probabilities: 1/1 + 1/2 + 1/3 + … + 1/n, the harmonic number H_n ≈ ln n + 0.577. That is why the algorithm is O(ln n) expected time — for a thousand buckets it loops about seven times, for a million about fourteen. Each iteration is a multiply, an add, a shift, and a floating-point divide, so a lookup is on the order of tens of nanoseconds; the paper reports jump hash outrunning Karger-style ring hashing while using no memory whatsoever. Space is O(1): two stack variables and nothing else, so it needs no initialization, no warm-up, and no synchronization when threads share it.

Balance is not just good, it is essentially ideal. Because every key is assigned by the same fair 1/k coin at each step, bucket sizes follow the binomial distribution you would get from truly random assignment — the standard deviation is the theoretical minimum, which the authors confirm empirically. Ring hashing, by contrast, only approaches that evenness by paying for hundreds of virtual nodes per server; jump hash reaches it with zero.

The same coin as reservoir sampling

There is an elegant way to see why the design is correct: jump consistent hash is deterministic reservoir sampling with a reservoir of size one. In Vitter's Algorithm R, you scan a stream and keep the i-th item as your single sample with probability 1/i, replacing whatever you held. After the whole stream you are left holding each item with equal probability. Jump hash runs exactly that process over the stream of bucket indices 0, 1, 2, …, n−1: bucket j "replaces" the current choice with probability 1/(j+1), and the final holder is the assigned bucket. The difference is that the coin flips are not real randomness but a key-seeded LCG, so the stream can be replayed identically on any machine without communication — and the closed-form j = (b+1)/r lets it skip the steps where no replacement happens, turning an O(n) scan into an O(ln n) one. Recognizing the reservoir-sampling skeleton is also what makes the balance proof trivial: uniform final selection is the defining property of Algorithm R.

The catch: buckets must be numbered 0..n-1

Jump hash buys its simplicity with one hard constraint: buckets are anonymous integers 0 through n−1, and you may only add or remove at the top. Shrinking from n to n−1 is clean — the keys that were in bucket n−1 redistribute, and the paper's construction guarantees they land minimally. But there is no way to pull out bucket 17 while keeping 0–16 and 18–n. The mapping has no notion of which server is which; it only knows counts. Remove a middle bucket and you must renumber everything above it, which shifts a large fraction of keys — precisely the stampede consistent hashing exists to prevent.

This makes jump hash superb for the case it was built for — sharding a cluster you scale up and down at the tail, such as a growing set of database partitions or a horizontally-scaled cache tier — and unsuitable for arbitrary membership, where individual nodes with fixed identities join and leave in any order (a peer-to-peer overlay, a fleet of load-balancer backends that fail independently). Two more limits follow from the same anonymity: jump hash offers no native weighting (every bucket gets an equal share; heterogeneous server sizes need an external indirection layer), and it returns a bucket number, not a server, so you must maintain a stable index → server table on the side and never reshuffle it out from under the hash.

Where it runs, and its cousins

The most widely shipped implementation is in Google's Guava library as Hashing.consistentHash(long, int) — a direct port of this algorithm that Java services use to pin keys to shards. Ports exist in Go, C++, Rust, Python and more, and it shows up wherever a system shards a monotonically-numbered set of partitions and wants rebalancing to be cheap. It is best understood alongside its relatives. Rendezvous hashing (Highest Random Weight, Thaler & Ravishankar, 1996) hashes the pair (key, node) for every node and picks the maximum: it handles arbitrary add/remove and weighting that jump hash cannot, at O(n) per lookup. Google's Maglev load balancer (2016) takes yet another route — a fixed-size permutation lookup table that gives O(1) lookups and tolerates backends leaving by identity, trading a modest amount of table memory and a hair of imbalance for that flexibility.

So the field is a set of trade-offs rather than one winner: ring hashing for arbitrary membership with tunable weights, rendezvous for simplicity and per-node control, Maglev for O(1) stateless data-plane lookups, and jump consistent hash for the narrow but common case where you want the theoretically minimal movement and perfect balance, for free, in five lines. Open extensions — weighted variants, and adaptations that relax the sequential-numbering rule — remain an active corner of systems engineering, but none has matched the original's stark economy.

Jump hash vs. the two classic consistent-hashing schemes, for a cluster of n buckets
SchemeLookup timeMemoryRemove arbitrary bucket?Weighting
Ring / Karger (with V virtual nodes)O(log(nV))O(nV) sorted ringYesYes (vary vnode count)
Rendezvous / HRW (Highest Random Weight)O(n)O(n) node listYesYes (weighted HRW)
Maglev (lookup table)O(1)O(M) permutation tableYesYes
Jump consistent hashO(ln n)O(1) — noneNo (top only)No (equal only)

Frequently asked questions

Why can't I remove a bucket in the middle?

The algorithm never stores which server is which — it only maps a key to an index in 0..n-1 based on the count of buckets. Removing bucket 17 while keeping the rest would require renumbering every higher bucket, and that renumbering shifts a large fraction of keys, defeating the consistency guarantee. You can only shrink from the top (drop bucket n-1), which redistributes minimally.

How is it consistent if it uses a random number generator?

The generator is a linear-congruential sequence seeded by the key itself and reset at the start of every call. The same key therefore replays the identical sequence of pseudo-random values on any machine, at any time, so it always jumps to the same final bucket. It is deterministic pseudo-randomness, not true randomness.

What exactly does O(ln n) mean here in practice?

The loop runs once per 'jump' the key would make as the cluster grows from 1 to n buckets, and the expected number of jumps is the harmonic number H_n ≈ ln n. That is about 7 iterations for 1,000 buckets and about 14 for a million — each iteration a multiply, shift, and divide, so a lookup takes tens of nanoseconds.

How does it compare to ring-based consistent hashing?

Ring hashing (Karger 1997) needs a sorted ring with 100-200 virtual nodes per server to get even balance, costing O(nV) memory and a binary search per lookup. Jump hash achieves perfect balance with zero memory and O(ln n) lookups. The price is that jump hash can't remove arbitrary nodes or weight them, which ring hashing can.

Does it support servers of different sizes?

Not natively — every bucket receives an equal 1/n share. To give a bigger server more load you must add an indirection layer, for example mapping several bucket indices to the same physical server, or use a scheme built for weights like weighted rendezvous hashing. Plain jump hash assumes homogeneous, equal-weight buckets.

Where is jump consistent hash actually used?

Google's Guava library exposes it as Hashing.consistentHash(long, int), and it is ported to Go, Rust, C++, and Python. It fits systems that shard a sequentially-numbered, tail-scaled set of partitions — growing database shards or cache tiers — where cheap, balanced rebalancing matters more than arbitrary node membership.