In a sharded cache serving millions of distinct keys, how do you detect hot keys in real time without an exact counter per key?
answer
- do not count everything
- hot keys survive sampling
- rows of counters, take the minimum
- collisions only inflate
- heap of k, reset each window
basics
~20 sSample 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 sExact 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 lineson_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
Recall that hot keys are found by counting requests per key, and that sampling makes the counting cheap enough to run in production.
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.
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.
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