YouTube Top K (Heavy Hitters / Trending)
Sharpened prompt. Design a service answering "what are the top 100 most-viewed videos" over the last 5 minutes, hour, day, and week — globally and per country and per category — from a stream of 6B view events/day across billions of distinct videos, with results under a second old, using bounded memory per node.
Exact counting is trivially correct and completely infeasible. This problem is a tour of the approximation techniques you use when the exact answer costs terabytes, and of the honesty required about what those approximations actually guarantee.
1. Problem framing
Functional requirements
- Top-K by view count over multiple time windows (5m, 1h, 24h, 7d).
- Sliced by country, category, and language.
- "Trending" — a ranking that weights velocity, not just volume.
- Sub-second query latency; results at most a few seconds stale.
Non-functional requirements
| Property | Target | Consequence |
|---|---|---|
| Memory | Bounded per node regardless of key cardinality | Sketches, not hash maps |
| Accuracy | Top-K correct with high probability; counts within a few percent | Approximation is acceptable and must be quantified |
| Freshness | < 5s | Streaming, not batch |
| Query fan-out | Hundreds of (window × slice) combinations | Precompute the combinations; do not compute per query |
Back-of-the-envelope
Events: 6B views/day = 70k/sec avg, ~250k/sec peak
Distinct: ~2B videos with at least one view in a week
Exact count: 2B keys × (16 B id + 8 B count + ~48 B hashmap overhead) = ~144 GB
per window per slice. With 4 windows × 200 slices = 115 TB. Absurd.
CMS: width 2^20, depth 5, 4-byte counters = 20 MB per sketch
ε = 2/2^20 ≈ 0.0002 (0.02% of stream size), δ = 1/2^5 ≈ 3%
4 windows × 200 slices × 20 MB = 16 GB total. Fits on one large node.
Space-Saving: exact-ish top-K with 10×K counters = 1,000 entries ≈ 50 KB per slice
The comparison — 115 TB exact versus 16 GB sketched — is the argument for the whole design. Lead with it.
2. High-level architecture
The verifier is the box most designs omit and the one that makes the answer defensible: approximate structures nominate candidates, and an exact recount over that small candidate set produces the published ranking.
3. Component inventory
| Component | Concrete choice | Why this one |
|---|---|---|
| Stream processor | Flink with RocksDB state backend | Event-time windows, exactly-once checkpoints, and state larger than memory |
| Frequency sketch | Count-Min Sketch (conservative update variant) | Bounded memory, mergeable, one-sided error you can reason about |
| Top-K structure | Space-Saving (Stream-Summary) | Deterministic error bound on the top-K set specifically |
| Alternative | HeavyKeeper | Better accuracy than CMS at the same memory for heavy hitters specifically |
| Decay | Exponentially decayed counters | Trending needs recency weighting, not raw volume |
| Serving | Redis with precomputed lists per (window, slice) | Queries are lookups; nothing is computed at read time |
4. The toughest parts
4.1 Counting billions of keys in megabytes
Why it's hard. A hash map keyed by video ID grows with the number of distinct videos — 2B keys is ~144 GB per window per slice, and the memory is dominated by keys and pointer overhead, not by the counts themselves.
Solution — Count-Min Sketch: a fixed-size 2D array of counters, indexed by several hash functions.
class CountMinSketch:
def __init__(self, width=1 << 20, depth=5):
self.w, self.d = width, depth
self.table = [[0] * width for _ in range(depth)]
self.seeds = [random.getrandbits(64) for _ in range(depth)]
def add(self, key, count=1):
# Conservative update: only raise the counters that are currently minimal.
# Reduces overestimation substantially for skewed streams at no extra memory.
idxs = [(i, mmh3.hash64(key, self.seeds[i])[0] % self.w) for i in range(self.d)]
cur = min(self.table[i][j] for i, j in idxs)
for i, j in idxs:
if self.table[i][j] < cur + count:
self.table[i][j] = cur + count
def estimate(self, key):
# Every cell may have collided upward, so the MINIMUM is the tightest bound.
return min(self.table[i][mmh3.hash64(key, self.seeds[i])[0] % self.w]
for i in range(self.d))
def merge(self, other):
# Linearity: sketches over disjoint streams add element-wise. This is why
# distributed aggregation works at all.
for i in range(self.d):
for j in range(self.w):
self.table[i][j] += other.table[i][j]
State the guarantee precisely, because that is the whole point of using it: with width w = ⌈e/ε⌉ and depth d = ⌈ln(1/δ)⌉, the estimate never undercounts, and it overcounts by more than ε·N (where N is the total stream size) with probability at most δ. At w = 2^20, d = 5, and a 6B-event stream, that is an error under ~1.2M with 97% confidence — negligible for a video with 50M views, and meaningless for one with 12, which is exactly the right shape of error for a top-K problem.
Mergeability is the property that makes it distributed. Because sketches add element-wise, each partition keeps its own and the aggregator sums them. Naming linearity as the reason is a good signal.
4.2 CMS finds counts; Space-Saving finds the list
Why it's hard. A Count-Min Sketch tells you the estimated count of a key you name. It cannot enumerate the top K, because it does not store keys at all. Iterating every possible video ID to query the sketch defeats the purpose.
Solution — pair the sketch with Space-Saving, which maintains a bounded set of candidate heavy hitters.
class SpaceSaving:
"""Keeps m counters. Guarantees every true top-K item is present when m >> K."""
def __init__(self, capacity=1000):
self.cap = capacity
self.counts = {} # key -> (count, overestimate_error)
def add(self, key, c=1):
if key in self.counts:
n, e = self.counts[key]
self.counts[key] = (n + c, e)
return
if len(self.counts) < self.cap:
self.counts[key] = (c, 0)
return
# Full: evict the minimum and give the newcomer its count as the error bound.
victim = min(self.counts, key=lambda k: self.counts[k][0])
vcount, _ = self.counts.pop(victim)
self.counts[key] = (vcount + c, vcount) # error <= vcount
def top(self, k):
# (count - error) is a guaranteed lower bound on the true frequency.
return sorted(((key, n, n - e) for key, (n, e) in self.counts.items()),
key=lambda t: -t[1])[:k]
The guarantee worth quoting: with m counters, any item whose true frequency exceeds N/m is guaranteed to be in the structure. Keeping m = 10K for a top-100 query means anything with more than 0.1% of the stream is certainly captured — comfortably covering trending videos, which are heavy hitters by definition.
The (count, error) pair matters: count - error is a guaranteed lower bound and count is an upper bound. When two candidates' intervals overlap, their relative order is genuinely unknown, and the honest engineering response is to resolve it with the exact verifier from 4.5 rather than to pretend.
HeavyKeeper is the modern alternative worth naming: it uses exponential-decay-based probabilistic eviction and empirically achieves substantially lower error than CMS at the same memory for the heavy-hitter task specifically.
4.3 "Trending" is a derivative, not a count
Why it's hard. Rank by raw 24-hour views and the list is dominated by the same catalogue hits every day — content that is popular but not trending. Rank by pure velocity and a video going from 2 to 200 views shows infinite growth and tops the chart. Neither is the product.
Solution — exponentially decayed counters plus a velocity score with a volume floor.
class DecayedCounter:
"""Lazily-decayed count: no background sweep over billions of keys."""
def __init__(self, half_life_s: float):
self.lam = math.log(2) / half_life_s
self.value, self.t = 0.0, time.time()
def add(self, c=1.0):
now = time.time()
self.value = self.value * math.exp(-self.lam * (now - self.t)) + c
self.t = now
def get(self):
return self.value * math.exp(-self.lam * (time.time() - self.t))
def trending_score(video) -> float:
fast = video.decayed_1h.get() # recent burst
slow = video.decayed_24h.get() # established baseline
# Acceleration, damped by a floor so tiny videos can't win on ratio alone.
velocity = fast / max(slow / 24.0, 1.0)
volume_floor = math.log1p(video.views_1h) # must actually be watched
novelty = math.exp(-video.age_hours / 48.0) # fresh content gets a boost
return velocity * volume_floor * novelty
Lazy decay is the implementation detail that makes this feasible: you cannot iterate billions of counters to decay them on a timer, so you store the last-update time and apply the decay factor when the counter is touched or read. This is the same trick as lazy token-bucket refill in the rate limiter, and pointing out the shared pattern is a nice cross-connection.
Note that decayed counters are still mergeable as long as all nodes decay to a common reference timestamp before summing — worth stating, because it is what keeps 4.4 valid.
4.4 Merging top-K across 100 partitions
Why it's hard. Each partition sees only its own slice of the stream and produces a local top-K. The union of local top-Ks is not the global top-K. A video ranked 11th on every one of 100 partitions has a large global total but appears in no local top-10. This error mode is subtle, real, and the thing interviewers probe.
Partition A top-3: X=100, Y=90, Z=85 ... W=50
Partition B top-3: X=100, Y=90, Z=85 ... W=50
(× 100 partitions)
Naive merge of local top-3s: X=10,000, Y=9,000, Z=8,500
Truth: W = 5,000 — a genuine #4 that never appeared locally.
Solution — over-fetch locally, merge sketches (not just heaps), and verify.
# 1. Each partition emits its local top (α · K) — over-fetch by 10×, not top-K.
LOCAL_K = 10 * GLOBAL_K
# 2. Emit the FULL CMS alongside the candidate list. The sketch is mergeable and
# lets the aggregator estimate any candidate's global count, including
# candidates it learned about from OTHER partitions.
def emit(partition):
return Payload(candidates=partition.ss.top(LOCAL_K), sketch=partition.cms)
# 3. Aggregator: union the candidate KEYS, sum the SKETCHES, then estimate each
# candidate against the merged sketch. This is what fixes the W case above.
def merge(payloads):
global_cms = reduce(CountMinSketch.merge, (p.sketch for p in payloads))
candidates = set(chain.from_iterable(k for k, _, _ in p.candidates for p in payloads))
return sorted(((k, global_cms.estimate(k)) for k in candidates),
key=lambda t: -t[1])[:GLOBAL_K]
The essential move is merging sketches rather than merging heaps. A heap merge can only rank keys some partition already nominated; a merged sketch can score any candidate globally, so a key that was 11th everywhere gets its true total as long as it was nominated somewhere. Over-fetching by 10× makes that nomination overwhelmingly likely.
Use a two-tier hierarchy (partition → regional → global) so the global merger receives dozens of payloads rather than thousands. Merging is associative, so the tree is correct at any depth.
4.5 When approximate is not good enough
Why it's hard. The published trending list is a product surface with real consequences — creators care, press covers it, and advertisers buy against it. "Number 3 and number 4 might be swapped" is a poor answer when someone asks why their video moved.
Solution — approximate to nominate, exact to rank.
async def publish_top_k(window, slice_, k=100):
# 1. Sketches nominate a generous candidate set. Cheap, streaming, always fresh.
candidates = merged_sketch_top(window, slice_, k * 20) # ~2,000 ids
# 2. Exact recount over ONLY those candidates. The events are already in
# ClickHouse; counting 2,000 known ids is a bounded, indexed query.
exact = await clickhouse.fetch("""
SELECT video_id, count() AS c
FROM view_events
WHERE ts >= {start} AND ts < {end}
AND video_id IN {candidates}
AND country = {slice}
GROUP BY video_id ORDER BY c DESC LIMIT {k}
""", ...)
# 3. Publish the exact ranking. Sketch error affects only which 2,000 were
# considered, and the probability that a true top-100 item misses a
# top-2,000 nomination is vanishingly small.
await redis.set(f"topk:{window}:{slice_}", serialize(exact))
The framing to give: "I use approximation for the part that is unbounded — scanning billions of keys — and exactness for the part that is bounded, which is counting two thousand known IDs. The published number is exact; only the candidate selection is probabilistic, and I can quantify that risk." That answer is both more accurate and more honest than either extreme, and it is what a production system actually does.
Publish the freshness and method alongside the list (as_of, window, method: verified) so downstream consumers and support can reason about discrepancies.
4.6 Late events, replays, and windows that must not lie
Why it's hard. Mobile clients buffer views offline and deliver them hours later. A Kafka partition lags and delivers a burst behind the others. A pipeline restart replays from the last checkpoint. Each can double-count, or can slot yesterday's views into today's window, corrupting a published ranking.
Solution — event time, watermarks, bounded lateness, and idempotent sinks.
views
.assignTimestampsAndWatermarks(
WatermarkStrategy.<View>forBoundedOutOfOrderness(Duration.ofMinutes(10))
.withIdleness(Duration.ofMinutes(1)) // idle partitions can't stall time
.withTimestampAssigner((v, ts) -> v.clientEventTimeMs))
.keyBy(v -> v.sliceKey())
.window(TumblingEventTimeWindows.of(Time.minutes(5)))
.allowedLateness(Time.hours(2)) // revise, don't discard
.sideOutputLateData(veryLateTag) // beyond 2h: reconcile in batch
.aggregate(new SketchAggregate())
.sinkTo(idempotentRedisSink); // upsert on (window, slice)
Four mechanisms, each fixing a specific lie: event time puts a view in the window it happened in, not the one it arrived in. Watermarks with idleness prevent one quiet partition from freezing progress for everyone. Allowed lateness revises a published window rather than silently dropping the data. And an idempotent sink keyed by (window, slice) makes a checkpoint replay overwrite rather than double-count.
Be explicit that a window can be revised after publication and that consumers must handle it. The alternative — freezing a window early and discarding late data — silently undercounts exactly the mobile-heavy regions where buffering is most common.
4.7 Serving hundreds of window-by-slice combinations
Why it's hard. 4 windows × 200 countries × 20 categories × several languages is thousands of distinct top-K lists. Computing any of them at query time is out of the question, and maintaining a separate sketch for every combination multiplies memory by the number of combinations.
Solution — exploit the additivity of both sketches and time windows.
Base granularity: (5-minute bucket) × (country) × (category)
-> ~200 × 20 = 4,000 sketches per 5-minute bucket, 20 MB each = 80 GB.
Too much. Reduce sketch width for fine slices: a per-country sketch sees
~1/200th of the stream, so w = 2^16 gives the same relative error at 1.25 MB.
4,000 × 1.25 MB = 5 GB per bucket. Keep 12 buckets (1 hour) hot = 60 GB. Fine.
Rollups by ADDITION, not recomputation:
1 hour = merge 12 five-minute sketches
24 hours = merge 24 hourly sketches
global = merge all country sketches
"Europe" = merge the European subset — an arbitrary slice union, for free
Two points worth making. Sketch width should scale with the stream each sketch sees, not be a global constant — a per-country sketch needs far less width for the same relative error, and this is what keeps total memory tractable. And arbitrary slice unions come for free from linearity: "top videos in the EU" requires no precomputation because you merge the member-country sketches on demand.
Materialise only the combinations that are actually queried (track query patterns and precompute the hot ones), and compute the long tail on demand from merges. Serving is then a Redis lookup for hot slices and a sub-second merge for cold ones.
5. What breaks first
| Event | First failure | Mitigation |
|---|---|---|
| A single video goes massively viral | Hot Kafka partition (keyed by video_id) | Salt the key for extreme heavy hitters; sketches merge regardless of partitioning |
| Cardinality explosion (bot-generated IDs) | Space-Saving churns, real hitters evicted | Sketches are memory-bounded by construction; filter obvious bots upstream |
| Pipeline restart | Duplicate counting on replay | Exactly-once checkpoints + idempotent sink keyed by (window, slice) |
| Mass late arrivals after an outage | Published windows revised hours later | Allowed lateness of 2h; side output plus batch reconciliation beyond it |
| Verifier query slow at peak | Publication lag | Fall back to publishing the sketch-only ranking, clearly labelled method: estimated |
| Query for an unmaterialised slice | Latency spike | On-demand merge with a small result cache; precompute slices that get queried twice |
6. Cheat sheet
- The number: exact counting is ~115 TB across windows and slices; sketches do it in ~16 GB. That ratio is the design.
- CMS gives mergeable, never-undercounting frequency estimates:
ε = e/w,δ = e^-d. Use conservative update. - Space-Saving gives the candidate list with a per-item
(count, error)interval; CMS alone cannot enumerate. - Trending = velocity × volume floor × novelty, with lazily-decayed counters (same trick as token-bucket refill).
- Merging: over-fetch 10×K locally and merge sketches, not heaps — otherwise a consistent #11 disappears globally.
- Publish exact: approximate to nominate ~2,000 candidates, then recount them exactly in the OLAP store.
- Time: event time, watermarks with idleness, 2h allowed lateness, idempotent sinks; windows can be revised.
- The one-liner: "Sketches to survive unbounded cardinality, an exact recount over a bounded candidate set to survive scrutiny — approximation where the input is infinite, exactness where the output is public."