Design YouTube Top K Videos
DifficultyHard | HelloInterview: problem breakdown
Problem Statement
Design a system that tracks video view counts in real-time and returns the top K most-viewed videos for a given time window (e.g., top 10 videos in the last 1 hour, 24 hours, 7 days).
📖Real-world: "Top K" / "heavy hitters" is one of the most elegant problems in systems design because you can't just count everything — at billions of events you don't have memory to keep an exact counter per video, so you use probabilistic data structures. The star is the Count-Min Sketch: a clever grid of counters and hash functions that estimates any item's frequency in fixed memory (it may over-count slightly but never under-counts), paired with a small min-heap holding the current top K. It's the same family as HyperLogLog (for counting distinct items) — the recurring theme being trade a little accuracy for enormous memory savings, which interviewers love because the naive "just sort the counts" answer doesn't scale. Security-adjacent relevance for this audience: the identical heavy-hitters machinery powers detection — finding the top talkers in NetFlow, the most-hit URLs in a DDoS, or the noisiest source IPs — so this pattern shows up directly in SOC tooling.
Requirements
Functional
- Track view events per video
- Query: "top K videos by view count in last [1h, 24h, 7d]"
- Near-real-time: results update within 1 minute
Non-Functional
- 1B view events/day → 11,500 views/s
- 1B+ videos total; K = 10-1000
- Query latency: < 100ms (pre-computed results)
- Approximate results acceptable for leaderboard
Core Design
This is a heavy hitters / top-K problem at stream scale.
Naïve Approach (doesn't scale)
Sort all videos by view count for a time window. At 1B videos and 11,500 events/s, a full sort is O(N log N) and can't be done in real-time.
Count-Min Sketch + Top-K Heap
Count-Min Sketchprobabilistic data structure for approximate frequency counting.
2D array: d rows × w columns
On event video_id:
for i in 0..d:
col = hash_i(video_id) % w
sketch[i][col] += 1
count(video_id) = min over all rows of sketch[i][hash_i(video_id) % w]Estimation error: with d=5 rows, w=10,000 columns → ~5 MB memory for 11,500 events/s. Error bound: ε = e/w probability of > ε*N error.
Top-K Heapmaintain a min-heap of K elements. When a video's estimated count exceeds the minimum in the heap, replace it.
sketch = CountMinSketch(d=5, w=10000)
top_k = MinHeap(k=100)
on_view_event(video_id):
sketch.increment(video_id)
count = sketch.estimate(video_id)
if count > top_k.min():
top_k.push(video_id, count)Memory: sketch (5 MB) + heap (100 video IDs) = trivially small.
Time-Windowed Aggregation
For multiple time windows (1h, 24h, 7d), use separate sketches per time bucket:
Approach: sliding window with 1-minute buckets
- Maintain 60 buckets for "last 1h" (one per minute, TTL 61 min)
- Maintain 1440 buckets for "last 24h" (one per minute, TTL 1441 min)
- 1h count for video = sum(sketch[minute-60 to minute-0].estimate(video_id))
OR: separate pre-computed aggregates
- 1-minute Flink job: compute top-K per minute
- 1-hour aggregator: merge last 60 minute results
- Store in Redis: top_k:1h, top_k:24h, top_k:7dArchitecture
View Events → Kafka (partitioned by video_id)
↓
Flink Stream Processor
(Count-Min Sketch + Top-K heap per partition)
↓
Merge step: combine per-partition top-K lists
↓
Redis: store pre-computed top-K results
↓
Query API: GET /top-k?window=1h&k=10 → Redis lookup (O(1))Exact vs Approximate
| Use case | Approach |
|---|---|
| Leaderboard display | Approximate (Count-Min + heap) |
| Billing (ads on viral videos) | Exact: batch aggregate from raw events in S3 |
| Fraud detection | Exact: need to verify real vs bot views |
Key Design Decisions
| Decision | Choice | Reason |
|---|---|---|
| Counting | Count-Min Sketch | O(1) update, constant memory |
| Top-K maintenance | Min-heap | O(log K) updates; K << total videos |
| Time windows | Pre-computed per window | O(1) query; update every minute |
| Exact counts | Batch job from S3 | For billing; separate from real-time path |
Interview Tips
explain the 2D array and min-over-rows estimation.
each Kafka partition processes its own sketch; merge the per-partition top-K lists with a secondary sort.
at 11,500 events/s, a hashmap is fine, but synchronizing it across a distributed cluster introduces coordination overhead that Count-Min avoids.
don't compute at query time; update every minute and serve from Redis.