skip to content

In a sharded cache that places keys with consistent hashing, why does adding more nodes fail to relieve a node overloaded by one hot key?

level: middleimportance: must knowfreq 62%

answer

  1. what the ring actually assigns
  2. requests versus keys
  3. one owner per key
  4. key rate above node capacity
  5. split or serve locally

basics

~20 s

A sharded cache places whole keys, so all of one key's traffic goes to its single owner node. Adding nodes only moves whole keys around, so a key that needs more than one node's capacity stays overloaded wherever it lands.

solid answer

~40 s

Consistent hashing maps each **key** to one owner node. Adding nodes changes *which* keys each node owns. It never divides one key's requests. If a node can serve 100,000 requests per second and one key needs 150,000, that key overloads its owner in a 10-node cluster and in a 100-node cluster. More nodes do lower the background load from *other* keys, which helps a bit, but they can't bring a single key under one node's ceiling. The fixes change the key itself or where its reads are served: split it into several suffixed keys that hash to different nodes, or serve it from a short-lived in-process cache on each app server so most reads never reach the shared tier.

go deeper

for a junior

Remember that a sharded cache sends every request for a given key to one node, so one very popular key can overload that node.

for a middle

Explain that placement works per key, show the arithmetic of a key whose rate exceeds one node's capacity, and separate that case from ordinary skew that rebalancing fixes.

for a senior

Show how you would diagnose it: one hot node while the average stays low, no relief after scale-out, and sampled logs pointing at one key. Then pick splitting or a local tier for that key.

for a principal

Frame the trade-off: generic capacity levers scale key count, not per-key rate, so hot-key handling needs its own detection and mitigation path rather than a bigger cluster.

## The key is the unit of placement A **sharded cache** spreads data over many nodes. A placement function decides which node owns each key. Common choices are `hash(key) mod N` and **consistent hashing**, where nodes and keys are hashed onto a ring and a key belongs to the next node clockwise. Either way, the rule works on the **whole key**: - every read and write for `post:42` goes to the same owner node - the function spreads *different keys* evenly, not *requests* evenly - a node's load is the sum of the request rates of the keys it owns When request rates follow a heavy-tailed popularity curve, as they do for a viral post, a flash-sale item or a live-score page, one key can carry a large share of all traffic. That key is a **hot key**. ## Why adding nodes does not help Take an illustrative cluster where each node safely handles 100,000 requests per second: 1. A single key receives 150,000 requests per second. 2. With 10 nodes, the key's owner gets 150,000 plus its share of everything else, which is already over capacity. 3. Doubling to 20 nodes halves that background share. The key still brings 150,000 on its own, which is 1.5x one node's capacity. 4. No node count fixes this, because the ring gives the key to exactly one owner. Adding nodes helps only when the key's own rate *fits* under one node's capacity and the overload comes from unlucky neighbours sharing that node. Then moving neighbours away, or moving the hot key to a quiet node, is enough. That case is **key skew across nodes**, and rebalancing solves it. A key that exceeds a node on its own is a **single-key hotspot**, and rebalancing cannot solve it. ## What each lever changes | Lever | What it changes | Helps a key hotter than one node? | |---|---|---| | Add nodes | Fewer other keys per node | No | | Move the key to a quiet node | Removes neighbours' load | Only if the key alone fits | | More virtual nodes per server | Smoother spread of *keys* | No | | Read replicas of the owner | Read capacity for that shard | Partly, for read-hot keys only | | Split into suffixed keys | One key becomes N placed keys | Yes, for reads or for counters | | Per-process local tier | Most reads never reach the shared tier | Yes, for read-mostly keys | **Virtual nodes** are a common source of confusion. They give each physical server many points on the ring so that *keys* spread evenly. A single key still hashes to one point and has one owner. ## What actually works The working fixes either stop treating the hot data as one key or stop sending its reads to the shared tier: - **Key splitting.** Store copies under `post:42#0` … `post:42#N-1`. Each suffix hashes independently, so the copies land on different nodes, and readers pick one suffix at random. Read load per copy falls to about 1/N, but every write now has to update N copies. - **Local in-process tier.** Each application process keeps the hot value in memory for a short TTL, so the shared tier sees about one refresh per process per TTL instead of every client read. - **Sub-counters for write-hot keys.** A counter that takes a flood of increments is split into several counters that are summed on read. All three need to know *which* keys are hot. That is why hot-key detection comes before mitigation. ## Recognising the symptom A single-key hotspot has a recognisable signature: - one node runs near its CPU or network limit while the cluster average is low - that node's load does not fall after a rebalance or a scale-out - sampled request logs show one key taking a large fraction of that node's traffic - latency spikes affect *every* key on that node, including cold ones, because they share its queue That last point matters in an interview. A hot key does not only slow its own readers. Every unrelated key that happens to share the node gets slower too, so the blast radius is the whole node.

  • When does moving the hot key to a quieter node fully solve the problem?
    When the key's own request rate fits comfortably under one node's capacity and the overload comes from other keys sharing its node. Moving the key, or its neighbours, then removes the combined load. If the key alone exceeds a node, moving it just moves the overload to another node.
  • Why do cold keys on the same node suffer during a hotspot?
    They share the node's CPU, network and request queue. When the hot key saturates those resources, every request waits longer, including requests for rarely-read keys. The blast radius of a single hot key is the whole node, which is why tail latency climbs for unrelated features too.

Opening more checkout lanes does nothing if every shopper insists on queuing at one particular cashier. You either clone that cashier or hand out the item at the door.

saying these in an interview costs you the question

  • Adding nodes spreads one key's requests across more machines
  • Consistent hashing balances request load, not just key placement
  • More virtual nodes will split a hot key's traffic
  • A hot key only slows down requests for that key
  • Rebalancing the ring always fixes an overloaded node