Distributed Systems
Maglev Hashing: How Google Balances Load Consistently
Maglev Hashing is the consistent-hashing algorithm inside Google's software network load balancer, first described at NSDI 2016, that decides which backend server each network connection is sent to. It has to do two things that pull against each other at once: spread millions of connections almost perfectly evenly across a fleet of backends, and barely disturb the existing connection-to-backend mapping when a server is added or removed, so live TCP flows are not reset.
The clever part is that it achieves both with a single flat lookup table and an O(1) array index — no ring to search, no per-key tree walk. That combination of near-perfect balance, minimal disruption, and constant-time lookup is what makes it remarkable, and why variants now show up far beyond Google.
- OriginGoogle Maglev, NSDI 2016 (Eisenbud et al.); in production since 2008
- LookupO(1) — entry[ hash(5-tuple) mod M ]
- Table size MA prime (65537 or 655373); kept ≥ ~100× the backend count N
- Build costO(N·M) permutations + ~O(M log M) fill; ≈1.8 ms for M=65537, N=1000
- Load evenness≈1% max/min imbalance; each backend gets ~M/N slots (±1)
- Throughput~10 Gbps of small packets per Maglev machine at line rate
Interactive visualization
Press play, or step through manually. The visualization is yours to drive — try it before reading on.
Watch the 60-second explainer
A condensed visual walkthrough — narrated, captioned, under a minute.
Two goals that fight each other
A layer-4 load balancer sits behind a virtual IP (VIP) and forwards every incoming packet to one of N backend servers. It faces two demands that are in tension. First, evenness: no backend should get meaningfully more connections than any other, or it becomes a hotspot. Second, consistency (minimal disruption): when the backend set changes — a server is drained, crashes, or is added — the mapping for unaffected flows must stay put, because a TCP connection whose packets suddenly land on a different backend has no matching socket there and is reset.
The naive answer, backend = hash(flow) mod N, is perfectly even but disastrous on the second count: change N and nearly every flow remaps. Classic ring consistent hashing fixes disruption but is only even if you sprinkle each server across the ring as many virtual nodes, and it costs a logarithmic search per lookup.
Maglev's setting makes both goals sharper. Google's routers spray packets for a VIP across many Maglev machines using ECMP (equal-cost multipath) with the VIP announced by BGP anycast. When a Maglev machine is added or removed, ECMP reshuffles which machine sees a given flow — so every Maglev must independently compute the same flow-to-backend decision, or a mid-stream connection breaks. Consistent hashing is what lets stateless machines agree.
The lookup table and each backend's preference list
Maglev's central data structure is a flat array called the lookup table, of size M, where M is a prime much larger than N (Google used M = 65537 and M = 655373). Each slot holds the index of one backend. A packet is routed by hashing its 5-tuple (source/dest IP, ports, protocol) and taking that value mod M to pick a slot: backend = entry[ hash(5-tuple) mod M ]. That is the whole lookup — a single array read, O(1), with no comparisons or tree walks.
The interesting question is how to fill the table so it is both balanced and stable. Each backend i is given a preference list: a permutation of all M slot indices, ranked from most to least wanted. The permutation is generated cheaply from two independent hashes of the backend's name:
offset = h1(name[i]) mod M
skip = h2(name[i]) mod (M-1) + 1 // in [1, M-1]
permutation[i][j] = (offset + j * skip) mod MBecause M is prime and skip is in [1, M-1], the step skip is coprime to M, so j = 0,1,…,M-1 visits every slot exactly once — a genuine permutation. The offset gives each backend a different starting point and skip a different stride, so distinct backends prefer the slots in different, well-scrambled orders.
Populate: backends take turns claiming slots
The table is filled by a round-robin auction. Every backend keeps a cursor next[i] into its own preference list; the table entry[] starts empty (-1). Going around the backends in order, each one claims its most-preferred slot that is still empty, then it is the next backend's turn. This repeats until all M slots are taken:
for each i: next[i] = 0
for each c: entry[c] = -1
n = 0
while true:
for i = 0 .. N-1:
c = permutation[i][ next[i] ]
while entry[c] >= 0: // slot already taken
next[i] += 1
c = permutation[i][ next[i] ]
entry[c] = i // claim it
next[i] += 1
n += 1
if n == M: return // table fullThe invariant is simple and powerful: because backends claim in turn, one at a time, the counts can differ by at most one until the very end. After a full pass every backend has claimed exactly one more slot, so each backend ends up owning between floor(M/N) and ceil(M/N) slots — an almost perfectly even split with no virtual nodes and no tuning. With M = 65537 and N = 1000, that is roughly 65 or 66 slots each, an imbalance of about 1%.
The disruption property falls out of the same turn-taking. Remove a backend and rebuild: the survivors still visit their preference lists in the same order and mostly reclaim the same slots. The departed backend's ~M/N slots are handed out, plus a modest amount of secondary churn where the changed fill order nudges a few other assignments. Maglev deliberately favors evenness slightly over the theoretical disruption minimum — and the bigger you make M relative to N, the closer disruption gets to the ideal 1/N.
The complexity and the numbers, and why they hold
Lookup: O(1). One hash of the 5-tuple, one modulo, one array read. This is the operation done per packet, millions of times a second, so its constant cost is the point.
Build: two parts. Generating the permutations is O(N·M) — one linear pass per backend. Filling the table is the interesting term. Early claims are instant, but as the table fills, a backend increasingly finds its top choice already taken and must skip forward. This is a coupon-collector-style process: the expected total number of skips to fill a prime table with well-scrambled permutations is O(M log M), giving an expected fill cost of roughly O(M log M). The theoretical worst case, if permutations pathologically collide, is O(M^2), but a prime M with independent offset/skip hashes makes that astronomically unlikely. In Google's measurements, building the table for M = 65537, N = 1000 took about 1.8 ms; the larger M = 655373 took on the order of tens of milliseconds. Since the table is only rebuilt on a backend-set change, not per packet, that cost is easily amortized.
Evenness. The turn-taking invariant bounds per-backend counts to floor(M/N)–ceil(M/N), so the max/min ratio is about 1 + N/M. Keeping M at least ~100× larger than N holds imbalance near 1%. Disruption. The sizing rule matters here too: a table only a few times larger than N churns visibly on changes, while M >> N pushes disruption toward the minimal one-server-share. This is the single knob — the ratio M/N — that trades a bigger table (more memory, slower build) for tighter balance and stickier mappings.
Why it beats ring and rendezvous hashing
Versus ring consistent hashing (Karger et al., 1997). The ring places each server at a hashed point on a circle and maps a key to the next server clockwise. On its own the arcs are wildly uneven, so real systems give each server V virtual nodes; the load imbalance shrinks only like 1 + O(sqrt((log N)/V)), so pinning skew to a few percent needs V in the hundreds — hundreds of ring entries per server, and an O(log(N·V)) binary search on every lookup. Maglev gets ~1% skew from one flat table and an O(1) index, no virtual nodes to tune.
Versus rendezvous / HRW hashing (Thaler & Ravishankar, 1998). Rendezvous computes hash(key, server) for every server and picks the maximum. It is beautifully even and provably minimal-disruption — remove a server and exactly its 1/N of keys move — but the natural lookup is O(N) per key (skeleton variants get O(log N)). For a load balancer doing this per packet across thousands of backends, O(N) is a non-starter. Maglev trades rendezvous's exactly-minimal disruption for a slightly-above-minimal disruption in exchange for O(1) lookups and a small, cache-friendly table.
The design lens is: precompute a good assignment once into a table, then make the hot path a single array read. That is why Maglev-style hashing has been adopted well beyond Google — it appears in Cilium/eBPF load balancers, Katran (Meta's XDP L4 balancer), Envoy's maglev load-balancing policy, and various service meshes as the go-to consistent-hashing policy.
In production: connection tracking, encapsulation, and failure modes
Consistent hashing is only half of how Maglev keeps connections alive. Each Maglev machine also keeps a connection-tracking table — a local hash map from 5-tuple to the backend already chosen for that flow. The fast path checks this table first; consistent hashing is the fallback used for a flow the machine has not seen, or when tracking state is lost. This layering matters because the two mechanisms cover different failures: connection tracking keeps a flow pinned even when the backend set changes (which would otherwise remap it), while consistent hashing keeps flows pinned when the Maglev set changes and ECMP reshuffles packets onto a fresh machine that has no tracking entry — as long as the backend set is the same, both machines hash to the same backend and the connection survives.
Once a backend is chosen, Maglev forwards the packet using GRE encapsulation, and backends reply straight to the client via Direct Server Return, so return traffic — the bulk of the bytes — never passes back through the balancer. A single Maglev machine, using kernel-bypass packet processing, sustains roughly 10 Gbps of small packets at line rate.
Failure and edge cases. The hard case is simultaneous change on both axes — a Maglev machine added while a backend is also draining — where a flow can land on a new Maglev that both lacks tracking state and computes a slightly different table; a small fraction of such flows can reset. Sizing M >> N minimizes the table-churn side of this. Two other practical notes: M must be prime for the permutation trick to cover every slot, and the name hashes must be well-distributed — a weak hash that clusters offsets or skips degrades both evenness and build time. Maglev, running Google's public-facing and internal VIPs since 2008 and described publicly in 2016, is the canonical demonstration that you can have even load, sticky connections, and constant-time lookups all at once.
| Scheme | Lookup cost | Load evenness | Disruption when a backend is added/removed |
|---|---|---|---|
| Modulo hashing (hash mod N) | O(1) | Near-perfect | Catastrophic — almost every mapping changes |
| Ring consistent hashing (Karger 1997) | O(log(N·V)) binary search on the ring | Needs V≈100–200 virtual nodes per server for ~5% skew | Minimal — only the leaving server's arc moves |
| Rendezvous / HRW hashing | O(N) (or O(log N) with a skeleton) | Near-perfect | Minimal — provably only 1/N of keys move |
| Maglev hashing | O(1) flat array index | ≈1% imbalance from one table, no virtual nodes | Small — the failed server's ~M/N slots plus modest churn |
Frequently asked questions
Why does Maglev use a prime number for the table size M?
Each backend's preference order is generated as (offset + j*skip) mod M with skip in [1, M-1]. When M is prime, every possible skip is coprime to M, which guarantees the sequence visits all M slots exactly once — a true permutation. If M were composite, some skip values would share a factor with M and the sequence would only touch a subset of slots, breaking the algorithm.
How does Maglev keep the load so even without virtual nodes?
Backends fill the lookup table by taking strict turns, each claiming its most-preferred still-empty slot before yielding. Because they claim one at a time in rotation, no backend can get more than one slot ahead of another, so every backend ends with between floor(M/N) and ceil(M/N) slots — roughly a 1% imbalance. Ring consistent hashing, by contrast, needs hundreds of virtual nodes per server to approach that evenness.
What actually happens to live connections when a backend fails?
The table is rebuilt, handing out the failed backend's ~M/N slots and causing a modest amount of secondary churn as the fill order shifts. Flows that were already on surviving backends mostly keep their slots, so they are not reset. In addition, each Maglev machine's connection-tracking table pins flows it has already seen, so most established connections survive a backend change untouched.
How is Maglev hashing different from rendezvous (HRW) hashing?
Rendezvous hashing picks a backend by computing hash(key, backend) for every backend and taking the maximum, which is even and gives provably minimal 1/N disruption — but costs O(N) per lookup. Maglev precomputes one flat table so every lookup is a single O(1) array index, accepting slightly-above-minimal disruption in return for constant-time routing that scales to thousands of backends at line rate.
How large should M be relative to the number of backends N?
Google keeps M at least about two orders of magnitude larger than N (e.g., M=65537 for up to ~1000 backends). A larger M/N ratio tightens both load evenness (imbalance ≈ N/M) and consistency (disruption approaches the ideal 1/N), at the cost of more memory and a slightly slower table build. It is essentially the one tuning knob in the scheme.
Is Maglev hashing used outside of Google?
Yes. The algorithm from the 2016 NSDI paper has been widely reimplemented: Envoy offers a 'maglev' load-balancing policy, Cilium's eBPF datapath and Meta's Katran (XDP) L4 balancer use Maglev-style consistent hashing, and various service meshes adopt it as their default sticky-routing policy because it combines even load, minimal disruption, and O(1) lookup.