Skip to content

Saving another 100TB of RAM with math (and Rust)

Summary

Cloudflare's Performance team found that Pingora Backend Router (PBR) — the internal load-balancing service that routes cacheable requests to servers by URL via consistent hashing — was using far more memory than expected (up to 6 GB in some cases) in structures owned by pingora-ketama, its open-source consistent-hashing library. Two changes reclaimed more than 100 TB of RAM globally: (1) a struct-packing fix that shrank each hash Point from 8 bytes to 6 bytes (a 25 % cut), and (2) a first-principles derivation of the load-balancing error formula that proved Cloudflare could cut the number of hashes per server by 90 % with no appreciable increase in load imbalance. The migration to the smaller ring was rolled out with a dual-ring, request-scoped, data-center-scoped scheme to avoid invalidating cached content network-wide. This is the second 100 TB RAM win in as many months — the DNS team shed another 100 TB the prior month.

Key takeaways

  1. Consistent hashing routes cacheable requests by URL so each data center keeps only one copy of a file. PBR maps servers (by IP) and tasks (by cache key / URL) onto a 32-bit hash number line; a task is served by the first server to its left, giving a stable file→server location that barely moves when servers are added or removed (Source: sources/2026-09-18-cloudflare-saving-another-100tb-of-ram-with-math-and-rust).

  2. A single hash per server gives terrible load balance — the coefficient of variation is ~99 % at N=100. Expected arc size is 1/N; standard deviation is (1/N)·√((N−1)/(N+1)), so CV = SD/Exp = √((N−1)/(N+1)). At N=100 that's ~99 %, meaning some servers handle nearly twice their fair share while others do almost nothing.

  3. Adding multiple hashes per server evens out load (law of large numbers). With 160 points per server (the value hardcoded in NGINX and used as the Pingora default), the CV at N=100 drops from ~99 % to ~8 %.

  4. Ketama weighting scales hashes by server capacity. To make server S1 take w× the traffic of S2, give it w× the hashes. Cloudflare weights by disk space (cacheable-content workload); elsewhere in the company weights are CPU/GPU count. A base-160 scale times a weighting factor (example w=625) can explode to k = 160×625 = 100,000 hashes per server.

  5. Feature combinations multiply rings, not just hashes. Because only a subset of servers can serve any given request (compliance, enabled caching features), each combination of features needs its own ring: 2^handful = dozens of separate rings. This combinatorial explosion — not any single ring — is what drove the 6 GB memory footprint.

  6. Struct packing cut ketama storage 25 % with zero semantic change. The Point { hash: u32, index: u32 } struct was 8 bytes; the index never needs more than 16 bits (PBR won't coordinate more than 2^16 ≈ 65k servers). But shrinking index to u16 alone does nothing — Rust alignment rules force the struct size to a multiple of its largest field (4-byte hash), so it stays 8 bytes. Storing the fields as a raw [u8; 6] byte array with getters (u32::from_ne_bytes / u16::from_ne_bytes) compiles to the same code as #[repr(packed)] without its well-known hazards, and yields the full 8→6 byte (25 %) reduction.

  7. The full CV-for-k-hashes formula shows diminishing returns — and 32-bit collisions make it worse past ~10k hashes. The derived exact formula is SD_k = √((k+1)/(N(kN+1)) − 1/N²), giving CV_k = √((N−1)/(N·k+1)). Each order-of-magnitude increase in k buys a smaller error reduction — the last 90,000 of 100,000 hashes bought only ~0.7 % improvement. Worse, with 32-bit hashes the birthday-paradox collision probability rises fast; simulations for 2048-server DCs showed error rate actually increasing between 10,000 and 100,000 hashes per server, because collisions randomly drop hash contributions. Conclusion: cut hashes per server by 90 % with no meaningful error cost.

  8. The migration used a dual-ring, request-scoped, DC-scoped rollout to protect origins. Changing the ring changes where cacheable requests land, so a single global flip would have invalidated almost all cached content and turned a memory win into an origin-traffic apocalypse. Instead PBR carried both the old ketama ring and the new smaller one in memory; each request deterministically chose a ring via the normal migration framework (stable per request hash, clean rollback path). Rollout controlled two independent dimensions — how much traffic used the new ring, and where it was allowed to move — starting at small validation locations and expanding by data-center groups. Watched: backend-selection traces, ring-version counters, PBR connection errors, process memory, startup time, cache behavior, origin traffic. Decommissioning the old rings dropped used memory by 100 TB.

  9. Shipped as a v2 ring in the open-source pingora-ketama crate. The v2 ring has the compacted 6-byte storage format, a faster sorting method, and a tunable base hash count per node, behind an (initially unadvertised) cargo feature. The v1 ring is byte-for- byte the historical behavior; the library can run both simultaneously and pick per request — the exact primitive the migration relied on.

Operational numbers

  • >100 TB RAM reclaimed globally by this change (on top of a separate ~100 TB DNS-team win the prior month).
  • Up to 6 GB per-instance ketama memory before the fix (Ivan's original ticket).
  • 8 → 6 bytes per hash Point (25 % reduction) from struct packing.
  • 90 % reduction in hashes generated per server, with no appreciable error increase.
  • 160 = NGINX/Pingora default hashes per server; example weighted k = 160 × 625 = 100,000.
  • Load-balancing CV at N=100: ~99 % (1 hash) → ~8 % (160 hashes).
  • 2^handful = dozens of consistent-hash rings (one per feature combination) held in memory.
  • 32-bit hash space; collisions materially degrade balance between 10,000 and 100,000 hashes per server for 2048-server DCs.

Caveats

  • The elegant continuous-ring math (expected value / standard deviation / coefficient of variation) is only exact for an idealized continuous hash ring. Real 32-bit integer hashes have collisions whose probability grows per the birthday paradox, so the theoretical "more hashes is always at-least-as-good" breaks down — this is empirical, shown by simulation, not by the closed form.
  • The 6 GB figure was "in some cases"; the memory footprint is dominated by the number of rings (feature-combination explosion) more than any single ring's size.
  • #[repr(packed)] would achieve the same layout but is controversial for good reasons (unaligned-reference UB); the raw-byte-array + getters approach was chosen as the safer equivalent.
  • The migration's correctness rested on the ring choice being stable per request hash — a plain global-percentage rollout would have spread cache churn everywhere at once; the data-center-scoped second dimension is what kept blast radius small.

Source

Last updated · 766 distilled / 2,225 read