skip to content

In a sharded cache serving millions of distinct keys, how do you detect hot keys in real time without an exact counter per key?

level: seniorimportance: should knowfreq 48%

answer

  1. do not count everything
  2. hot keys survive sampling
  3. rows of counters, take the minimum
  4. collisions only inflate
  5. heap of k, reset each window

basics

~20 s

Sample a small fraction of requests at the cache clients, feed them into a count-min sketch plus a small top-k heap, and each window report keys whose scaled rate crosses a threshold, then reset or decay the counts.

solid answer

~50 s

Exact counting means one counter per distinct key, which is too much memory and work on the hot path. Instead, **sample** perhaps 1 in 100 requests where the key is known (the client library or a proxy). Feed each sampled key into a **count-min sketch**: `d` rows of `w` counters, with one hash per row. You increment one counter in each row and read the estimate as the minimum of those counters. The estimate can over-count because of collisions, but it never under-counts. Alongside the sketch, keep a **top-k min-heap** of the heaviest keys seen. At the end of each window, multiply estimates by the sampling rate, report keys over a hot threshold, and reset or decay the counts so old heat fades. Merge reports from many clients, because each client sees only part of the traffic.

code

pseudocode · 16 lines
pseudocode
on_request(key):
  if random() >= 1 / SAMPLE_FACTOR: return
  for i in 0..d-1:
    sketch[i][hash_i(key) mod w] += 1
  est = min over i of sketch[i][hash_i(key) mod w]
  if heap.contains(key): heap.update(key, est)
  else if heap.size < k: heap.push(key, est)
  else if est > heap.min().count:
    heap.pop_min()
    heap.push(key, est)

every WINDOW_SECONDS:
  for (key, est) in heap:
    rate = est * SAMPLE_FACTOR / WINDOW_SECONDS
    if rate >= HOT_RATE: report(key, rate)
  sketch.reset(); heap.clear()

go deeper

for a junior

Recall that hot keys are found by counting requests per key, and that sampling makes the counting cheap enough to run in production.

for a middle

Explain how a count-min sketch adds and estimates, why its error only goes upward, and why it needs a top-k heap beside it to name keys.

for a senior

Show the production choices: sample at the client, scale estimates by the sampling factor, reset or decay per window, merge reports across clients, and set thresholds from node capacity.

for a principal

Weigh detection latency against noise and cost, and decide where counting lives so that mitigation such as a local tier never hides the evidence the detector depends on.

## The problem with exact counting A cache that serves millions of distinct keys at hundreds of thousands of requests per second can't afford a precise hash map of `key -> count` on every client: - memory grows with the number of distinct keys, including millions that are read once - every request pays for a map update on the latency-critical path - most of those counters describe **cold keys** that nobody cares about The question is narrower: *which few keys are carrying a disproportionate share of traffic right now?* That is the **heavy hitters** problem, and approximate algorithms answer it with a small, fixed amount of memory. ## Step 1: sample **Sampling** records only a random fraction of requests, for example 1 in 100: - a key at 50,000 requests per second yields about 500 samples per second, which is plenty to spot - a key at 5 requests per second yields about 0.05 samples per second, so it barely registers, which is fine - the estimated true rate is `samples x sampling_factor / window_seconds` Sampling works *because* hot keys are hot. The keys that matter are the ones that survive heavy sampling. Where you sample matters. Counting at the **client or a proxy** sees the key name before any local caching hides it. Counting only on the cache node misses reads that were served from an in-process tier. ## Step 2: count-min sketch A **count-min sketch** is a `d x w` grid of counters with `d` independent hash functions: 1. To add a key, increment `sketch[i][hash_i(key) mod w]` for every row `i`. 2. To estimate a key, take the **minimum** of those `d` counters. 3. Collisions only ever *add* to a counter, so the estimate is never below the true count. The minimum picks the row with the fewest collisions. The standard sizing is `w = ceil(e / epsilon)` and `d = ceil(ln(1 / delta))`. With probability `1 - delta`, the over-count is at most `epsilon x total_count`. For example, epsilon = 0.001 and delta = 0.01 give `w = 2,719` and `d = 5`. That is 13,595 counters, about 54 KB with 4-byte counters, whether the stream has a thousand distinct keys or a billion. ## Step 3: top-k heap A sketch estimates counts but can't list keys, because it stores no key names. So you pair it with a **min-heap of size k** that holds the current top candidates: - after each sampled increment, look up the key's estimate - if the key is already in the heap, update its count - otherwise, if the heap is full and the estimate beats the heap's minimum, evict the minimum and insert the key ## Step 4: windows and decay Heat is temporary. A post that went viral an hour ago may be cold now. Detectors therefore either **reset** the sketch and heap each window (say every 1-10 seconds) or **decay** them by periodically halving every counter. Without this, counts grow forever and yesterday's hot keys crowd out today's. ## Step 5: aggregate and act Each client sees only its own slice of traffic, so the hot-key picture comes from **merging** reports: - clients ship their top-k lists (key, estimated rate) to a collector every window - the collector sums the rates per key and applies a threshold tied to node capacity, such as "more than 20% of one node's budget" - the resulting **hot set** drives mitigation: promoting the key to a local tier or splitting it | Approach | Memory | Accuracy | Latency cost | |---|---|---|---| | Exact per-key map | Grows with distinct keys | Exact | Map update on every request | | Sampling only, exact map of samples | Smaller, still unbounded | Good for heavy keys | Low | | Sampling + count-min + top-k | Fixed, tens of KB | Over-count bounded by epsilon | Low | | Per-node metrics only | Tiny | Finds the hot *node*, not the key | None | Per-node metrics (CPU, network, request rate per node) are still worth watching. They tell you *where* to look, and the sketch tells you *which key* is responsible. ```pseudocode on_request(key): if random() >= 1 / SAMPLE_FACTOR: return for i in 0..d-1: sketch[i][hash_i(key) mod w] += 1 est = min over i of sketch[i][hash_i(key) mod w] if heap.contains(key): heap.update(key, est) else if heap.size < k: heap.push(key, est) else if est > heap.min().count: heap.pop_min(); heap.push(key, est) ```

  • Why count at the client rather than only on the cache node?
    The client sees every request before any in-process caching hides it, and it knows the key name without extra work on the node. Once a hot key is served from a local tier, a node-side counter sees its rate collapse and would wrongly declare it cold. Per-node metrics still help find which node is suffering.
  • How long should the detection window be?
    Short enough to react before the node saturates, usually a few seconds, and long enough to collect meaningful samples at your sampling rate. A shorter window reacts faster but is noisier, so pair it with hysteresis: promote after one hot window, demote only after several cool ones.

saying these in an interview costs you the question

  • A count-min sketch can under-count a key's frequency
  • Keep an exact counter for every key on every client
  • A count-min sketch can list the hottest keys by itself
  • Counts never need resetting because popularity is stable
  • Per-node CPU graphs are enough to identify which key is hot