SYSTEM DESIGN #18 · INTERVIEW GUIDE
Design a Top-K (Heavy Hitters) System
What are the top 10 trending hashtags on Twitter right now? Which products are most-clicked on Amazon in the last 5 minutes? Which IPs are sending the most requests to your API gateway? These are all Top-K Heavy Hitters problems — and they need to be solved in real time, at millions of events per second, with sub-second freshness. This system uses Count-Min Sketch for probabilistic frequency estimation with bounded memory (O(ε⁻¹ log δ⁻¹)), local Min-Heaps per Flink worker to track local top-K, and a central merge coordinator that combines all sketches every 30 seconds into a global leaderboard served via Redis Sorted Sets. You will learn the exact data structures used by Twitter Trending, Amazon Best Sellers, and Cloudflare’s DDoS detection pipeline.
Count-Min SketchApache FlinkRedis ZADDMin-HeapKafka
💡
The Gist — What Problem Are We Solving?
What’s trending right now? Computing the top K things from billions of events
A Top-K system answers: what are the K most popular items right now? Trending hashtags, most-searched terms, most-played songs. The hard part is computing this continuously over billions of events per day with fixed memory — you can’t store a counter for every possible hashtag. The Count-Min Sketch is the probabilistic data structure that makes this possible: accurate to within 1%, using kilobytes of memory regardless of how many distinct items exist.
💬Think of it as a ‘most popular’ leaderboard that updates every 30 seconds, never runs out of memory no matter how many distinct items appear, and never misses a genuine trending item.
These are the capabilities the system must deliver — what users and operators can actually do with it.
📊
Heavy Hitter Detection
📊Return top K items by frequency in a sliding time window
⏰
Time Windows
⏰Tumbling (hourly) and sliding (rolling 60-min) window support
🔢
Approximate Counting
🔢Count-Min Sketch for fixed-memory frequency estimation
🌊
Distributed Merge
🌊Merge sketches across partitions for global top-K
📡
Real-Time API
📡REST endpoint for current top-K; WebSocket for live leaderboard
⚡
Non-Functional Requirements
These define how well the system must perform — the quality attributes that separate a toy from a production system.
💾 Memory
💾Fixed memory regardless of distinct item cardinality
⏱️ Latency
⏱️Top-K results available within 30 seconds of events
📈 Throughput
📈Process billions of events per day; 500K events/sec
🎯 Accuracy
🎯<1% error rate for items appearing in top-K
⏰ Windows
⏰1-min, 5-min, 1-hr, 24-hr windows all supported
📊
Key Metrics — The Numbers That Define This System
The headline numbers to know cold — and be ready to explain how each one is achieved.
Fixed RAM
vs any cardinality
🏗️
System Architecture Diagram
Full data flow from source to serving. Each layer scales independently.
Ingestion
Top-K System Flow
Kafka
partitioned by item_id
→
Flink Workers
local Count-Min Sketch per partition
→
→
↓
Processing
Min-Heap top-K extraction
→
→
🗺️
End-to-End User Journey
Trace a single request end-to-end — the story interviewers want you to tell fluently.
1
Events stream in
— Kafka receives 500K events/sec; partitioned by item_id for locality
2
Local counting
— Each Flink worker maintains a Count-Min Sketch for its partition; increments 5 cells per item per event
3
Partial top-K
— Each worker maintains a local min-heap of size K; updated on each sketch query
4
Global merge
— Every 30s: all workers broadcast their local CMS to a merge coordinator; CMS cells summed globally (CMS is additive)
5
Top-K extracted
— Min-heap of size K built from global CMS estimates; evict minimum when new item has higher estimate
🔭
High-Level Design — Component Breakdown
Core components — each with a single, well-defined responsibility. The key architectural insight: each layer scales independently, and failure in one component is isolated from the rest.
1 — Kafka
Distributed event bus with RF=3 for durability. Partitioned by user_id_hash for per-user ordering. LZ4 compression reduces storage cost by 60%. Exactly-once semantics via idempotent producers and transactional consumers.
2 — Flink
Stateful stream processor with 30-second checkpointing to S3. KeyedProcessFunction provides per-key state (Union-Find, EWMA baseline, experiment assignment). Exactly-once processing via two-phase commit sink.
3 — Merge
Handles responsibilities for the Merge layer. Designed for independent horizontal scaling — additional instances added without architectural changes. Communicates asynchronously with adjacent components to maximise throughput and fault isolation.
4 — Min-Heap
Handles responsibilities for the Min-Heap layer. Designed for independent horizontal scaling — additional instances added without architectural changes. Communicates asynchronously with adjacent components to maximise throughput and fault isolation.
5 — Redis
In-memory data structure server handling hot-path lookups in <1ms. Used for: canonical ID cache, rate limiting (token buckets), session state, feature store, and leaderboards. Cluster mode with 6 shards for horizontal scale.
🔬
Low-Level Design — Deep Dives
Deep dives worth explaining in detail in any senior engineering interview. For each: know the data structure, the algorithm, the why, and the trade-off you made.
1 — Kafka Ingestion
500K events/sec · Keyed Partitions
Events partitioned by item_id % N_PARTITIONS. Each Flink parallel task consumes one partition — all events for the same item processed by the same worker (no cross-worker coordination for per-item counts). Kafka topic retention: 7 days. Each event: {item_id, event_type, ts, user_id}. Item_id is the primary grouping key for frequency tracking. Flink checkpoints to S3 every 30 seconds for exactly-once processing.
# Kafka producer: partition by item_id
producer.produce(
topic=’events’,
key=event[‘item_id’].encode(), # ensures per-item ordering
value=json.dumps(event).encode()
)
2 — Count-Min Sketch Worker
O(1) Update · 2KB State
Each Flink worker maintains a local Count-Min Sketch (w=2000 columns, d=5 rows) and a local min-heap of top-K=1000 items. On each event: update CMS cells for all d hash functions. If item’s CMS-estimated count > heap minimum: replace minimum with this item. CMS state per worker: 2000 × 5 × 4 bytes = 40KB. Heap: 1000 items × (8 + 8) bytes = 16KB. Total per-worker state: ~56KB — extremely memory-efficient.
class CMSWorker:
def __init__(self):
self.cms = CountMinSketch(w=2000, d=5)
self.heap = MinHeap(k=1000)
def process(self, item_id):
self.cms.update(item_id)
freq = self.cms.estimate(item_id)
self.heap.update(item_id, freq)
3 — Sketch Merge Coordinator
Cell-wise Min · Every 30s
Every 30 seconds, all Flink workers send their CMS matrices and local top-K heaps to the merge coordinator. Merge algorithm: cell-wise minimum of all worker CMS matrices (CMS merge property). Merged CMS used to estimate global frequency of all items in the union of all local heaps. Items re-ranked by global CMS estimate → global top-K. Merge result written to Redis ZADD global_topk (ZREMRANGEBYRANK to keep only top-K).
def merge_sketches(worker_sketches):
# Cell-wise minimum: CMS merge property
global_cms = np.minimum.reduce(
[s.matrix for s in worker_sketches]
)
# Merge all local heaps
all_items = set().union(*[s.heap_items for s in worker_sketches])
return [(item, global_cms.estimate(item)) for item in all_items]
4 — Redis Leaderboard
ZADD · Sub-ms Reads
Global top-K stored in Redis Sorted Set: ZADD global_topk frequency item_id. API: ZREVRANGE global_topk 0 99 WITHSCORES returns top-100 with counts. Leaderboard updated by merge coordinator every 30 seconds via MULTI/EXEC transaction: DEL + ZADD in single atomic operation to prevent stale reads. Read replicas serve leaderboard queries — 50ms replica lag imperceptible at 30-second refresh rate. TTL: entries older than 1 hour pruned to prevent stale trending items.
def update_leaderboard(top_k_items):
pipe = redis.pipeline()
pipe.delete(‘global_topk’)
for item_id, freq in top_k_items:
pipe.zadd(‘global_topk’, {item_id: freq})
pipe.execute()
# Serve API
def get_trending(limit=100):
return redis.zrevrange(‘global_topk’, 0, limit-1, withscores=True)
⚖️
Trade-offs & Decision Log
Every senior interview comes down to these decisions. Know the exact trade-off, the reasoning, and the specific numbers that justify each choice.
⚖️ Count-Min Sketch vs Exact Counter for Top-K
✓
Count-Min Sketch (CMS) ✅ Chosen
- O(1) update and O(1) query — constant regardless of cardinality
- Memory: O(1/ε × log(1/δ)) — 1% error needs only 2KB per sketch
- Handles 1M distinct items with same memory as 100
- Never overcounts — only overestimates (bounded error)
→
Exact HashMap counter
- Zero error — perfect frequency counts
- Memory grows linearly: 1M distinct items = 8MB minimum
- At 500K events/sec with 10M distinct keys: 80MB/worker
- Merging across workers requires full map transfer
💡Decision: CMS with w=2000 columns, d=5 rows (ε=0.001, δ=0.007) per Flink worker; merge by cell-wise min into global sketch every 30s
⚖️ Local Min-Heap vs Shared Sorted Set for Top-K
✓
Local Min-Heap → Merge ✅ Chosen
- No cross-worker coordination during processing
- Each worker maintains local top-K (K=1000) in O(log K)
- Merge coordinator combines N heaps every 30s: O(N×K×log K)
- Small window delay (30s) before global top-K is fresh
→
Global Redis ZADD (all workers write)
- Always consistent — single source of truth
- ZADD contention: 500K writes/sec overwhelms Redis
- Network cost: every event requires a Redis round-trip
- Redis becomes bottleneck and single point of failure
💡Decision: Local min-heaps during processing; merge to global Redis ZADD every 30s for serving; Redis sorted set served by read replicas
🎯Interview Questions — Answered
The exact questions interviewers ask — with production-grade answers
Q1
Why is Count-Min Sketch used instead of a HashMap for frequency counting?
At 500K events/sec with 10 million distinct items (hashtags, URLs, product IDs), a HashMap would require ~500MB per Flink worker just for frequency counts. With 100 workers, that’s 50GB of state that must be checkpointed to S3 every 30 seconds. Count-Min Sketch with ε=0.001 error uses ~2KB of memory regardless of cardinality — a 250,000× reduction. The trade-off is a probabilistic overcount (never undercount) of at most ε×N where N is total events. At 500K events in a 30-second window, the maximum overcount is 500 — acceptable for trending detection where we care about relative ranking, not exact counts.
Q2
How accurate is the Top-K result and what are the guarantees?
Accuracy guarantees: (1) No false negatives in the Top-K — if an item is truly in the top-K by frequency, it will appear in the result with probability ≥1-δ (where δ=0.007 for the configured sketch). (2) False positives are possible — an item not truly in the Top-K may appear if its estimated count is inflated by hash collisions. At K=100 and ε=0.001, empirical false positive rate is <0.1% per 30-second window. (3) Count error bound: reported_count ≤ true_count + ε×N with probability 1-δ. In practice, median error is <0.0001×N — much better than the worst-case bound. These guarantees hold even under adversarial inputs designed to maximise hash collisions.
Q3
How is the leaderboard kept consistent across read replicas?
Read replicas for the Redis ZADD leaderboard use asynchronous replication with replica lag typically <50ms. For exact consistency, clients read from the primary Redis node. For trending displays (where slight staleness is acceptable), clients round-robin across read replicas. Consistency window is bounded by the merge interval (30 seconds) — since the leaderboard is only updated every 30 seconds, a 50ms replica lag is imperceptible. For use cases requiring stronger consistency (e.g., prize competitions where top-K rank determines winners), the merge job writes to primary with WAIT 1 0 (wait for at least 1 replica to confirm) before making the update visible.
System Design Series · Every Tuesday & Thursday
Level up your system design interviews
Each post covers Gist, Functional & Non-Functional Requirements, Key Metrics, System Diagram, User Journey, HLD, LLD, and Trade-offs & FAQs.
Subscribe to never miss a post →
Previous Articles
Categories: System Design
Tags: count-min sketch, heavy hitters, interview prep, leaderboard, stream processing, system design, top-k
Leave a Reply