Virtual Nodes Won't Save You From One Viral Key
Your distributed cache ring uses 256 virtual nodes per physical server. You've verified the key-count distribution is close to ideal: every server owns roughly 1/N of the total keys, within a percent or two. Overnight, one specific key — a celebrity's profile — starts receiving 30% of all read traffic on the system.
A teammate suggests: "let's bump virtual nodes from 256 to 1,024 per server — that should smooth out the load spike." Will increasing the virtual node count fix this specific problem?
No — virtual nodes balance key count, not per-key access frequency.
Virtual nodes solve a specific, different problem: without them, a physical server occupies one point on the ring and owns one contiguous, unevenly-sized arc of the key space, so different servers can end up owning very different numbers of keys just from where their single hash position happens to land. Spreading each physical server across hundreds of virtual positions averages that out, so every server ends up owning close to 1/N of the keys — assuming each key is accessed with roughly equal frequency.
That assumption is exactly what breaks here. Consistent hashing (with or without virtual nodes) is a deterministic function of the key: a given key always maps to exactly one owner, full stop, regardless of how many virtual nodes exist. Adding more virtual nodes changes which server owns which slice of key-space on average, and makes that ownership more evenly distributed by key count — it does nothing to change the fact that this one celebrity key still maps to exactly one server, and that server alone absorbs 30% of read traffic no matter how finely the ring is subdivided elsewhere.
The real fix targets the hot key directly, not the ring: an
application-level (L1, in-process) cache on each app server for the
known-hot key, replication of that specific key across multiple
backend nodes with reads randomly routed among the copies, or
key-splitting (storing sharded sub-copies like celeb:taylor:0
through celeb:taylor:9 and routing reads round-robin across the
shards). This is the same "hot key" failure mode covered for the
product-page cache in Caching Strategies — consistent hashing's job
is spreading key count evenly; it was never designed to spread the
traffic of a single key across multiple owners.
Share this question