Design an Ad Click Aggregator
DifficultyHard | HelloInterview: problem breakdown
Problem Statement
Design a system that counts ad clicks in real-time and provides aggregated click metrics (clicks per ad, per campaign, per time window). Used by ad networks for billing, fraud detection, and optimization.
๐Real-world & security angle: Ad-click counting is a security problem wearing a data-engineering hat โ because clicks are money, every click is adversarial. Click fraud (bots, click farms, competitors draining your budget, publishers inflating their own numbers) is a multi-billion-dollar problem, and Google/Meta run enormous detection pipelines to filter invalid traffic before billing. So this design has two layers people forget: the stream-processing path (Kafka โ Flink/Spark, windowed aggregation, often approximate counts via Count-Min Sketch or HyperLogLog for hot keys) and a fraud-detection path that scores each click. The classic interview tension is fast-but-approximate (real-time dashboard) vs slow-but-exact (the billing ledger) โ most designs do both: an instant approximate count for optimization and a reconciled exact count for money. Idempotency matters too: a retried click event must not be counted twice, or you overcharge advertisers.
Requirements
Functional
- Record every ad click event
- Query: total clicks for an ad in a time range
- Query: top N ads by clicks in last 1h, 24h, 7d
- Near-real-time aggregation (< 1 minute lag)
- Billing-grade accuracy (cannot over or under-count)
Non-Functional
- 1B clicks/day โ 11,500 clicks/s average; 100k clicks/s peak (viral events)
- Query latency: < 1 second for pre-aggregated results
- Deduplication: same user double-clicking should count as 1
Core Design
Stream Processing Architecture
At 11,500 clicks/s, batch processing is too slow for near-real-time. Use stream processing.
Ad Click Event (browser/app)
โ
Ingestion API (load balanced, stateless)
โ
Message Queue (Kafka) โ partitioned by ad_id
โ
Stream Processor (Flink / Kafka Streams)
โ
โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
โ Aggregation Store (Redis/TSDB) โ โ sliding window counts
โ Raw Event Store (S3/HDFS) โ โ exact replay
โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
โ
Query Service โ Dashboards / BillingKafka Partitioning Strategy
Partition Kafka topic by ad_id:
- All clicks for the same ad go to the same partition โ same stream processor node
- Each processor maintains counts for its assigned ad IDs in memory
- No cross-partition coordination for per-ad aggregation
Stream Aggregation (Flink)
Flink job:
Input: click events (ad_id, user_id, timestamp, ...)
Tumbling window (1-minute buckets):
count(clicks) GROUP BY ad_id, window(1 min)
โ emit: {ad_id, minute, count} to aggregation store
Sliding window (last 1h, last 24h):
maintained in Redis:
ZINCRBY clicks:hourly:{ad_id} count {minute_bucket}
ZREMRANGEBYSCORE (remove buckets older than 1h)Exact vs Approximate Counts
For billing, you need exact counts. For dashboards/trending, approximate is fine.
| Use case | Approach |
|---|---|
| Billing | Exact: aggregate from raw event store (S3) with batch job |
| Real-time dashboard | Approximate: stream aggregation in Redis |
| Top K ads | Count-Min Sketch + heap (approximate heavy hitters) |
| Fraud detection | Exact per-user click counts in Redis |
Count-Min Sketch for Top K
2D array of counters: hash(ad_id) โ increment multiple rows
Estimate of count for ad_id: min of hashed counters
Maintains Top K heap: if estimated count > min in heap โ add
Memory: O(K + sketch_size) regardless of total unique adsDeduplication
Same user clicking the same ad multiple times (double-click, refresh) = 1 click for billing.
On click event arrival:
key = hash(user_id + ad_id + time_bucket_5min)
if Redis.SETNX(key, 1, TTL=10min) == 1:
count this click (first occurrence)
else:
discard (duplicate)Key Design Decisions
| Decision | Choice | Reason |
|---|---|---|
| Ingest | Kafka (partitioned by ad_id) | High throughput; ordered per partition |
| Processing | Flink stream | Sub-minute aggregation; windowing built-in |
| Billing counts | Batch from raw S3 events | Exact; replayable; source of truth |
| Dashboard counts | Redis counters | Fast reads; approximate is OK |
| Top K | Count-Min Sketch | O(1) memory; acceptable error rate |
| Dedup | Redis SETNX with TTL | Per-user, per-ad, per time window |
Security Considerations
| Threat | Mitigation |
|---|---|
| Click fraud (bots inflating click counts) | Deduplication; IP rate limiting; device fingerprinting; ML bot detection |
| Invalid traffic (IVT) | Compare click patterns against IAB invalid traffic standards; filter bot signatures |
| Click injection (mobile fraud) | Attribute only clicks within valid time window after ad display |
| Billing manipulation | Billing from immutable raw events in S3 (not mutable Redis counters) |
| Data exfiltration | User-level click data is PII; anonymize or aggregate before exposing in dashboards |
| Latency attack | Inject delayed clicks with backdated timestamps; reject clicks > N seconds old |
Interview Tips
is the central trade-off โ billing needs exact (batch from raw events), dashboards need fast (stream + approximate).
for Top K is a high-signal answer that shows you know probabilistic data structures.
is not obvious but critical โ it means aggregation is single-threaded per ad, no cross-partition coordination.
mention that you can't keep dedup state forever; use a time-bucketed key with TTL.