On this page
System Design — Ad Click Aggregation
Last reviewed 11 Sept 2026
Part of the system design series. See the framework and building blocks first if you haven’t.
1. Requirements
Functional
- Ingest ad click events (
ad_id,user_id,campaign_id,timestamp,ip) at high volume. - Aggregate clicks per ad, per minute, so advertisers/dashboards can query “clicks for ad X between t1 and t2.”
- Support ad-hoc filters (by campaign, by geo) for a reporting UI.
- Detect and drop obviously fraudulent clicks (bot farms, click storms) before they pollute the aggregate.
Non-functional
- High write throughput — tens to hundreds of thousands of events/sec at peak (a campaign launch, a Super Bowl ad).
- Near-real-time aggregate freshness for dashboards (seconds to low minutes), not millisecond query latency.
- Billing/reporting numbers must be exact, not approximate — this is money. Real-time dashboards can tolerate approximation; the end-of-day invoice cannot.
- Must tolerate late-arriving events (client retries, mobile devices reconnecting) without silently losing or double-counting them.
2. High-level architecture
flowchart TD C[Ad-serving clients] -->|click/conversion events| GW[Ingest API] GW --> K[[Kafka: raw click events]] K --> F[Flink job: window + aggregate] F --> AGG[(Aggregate store<br/>per ad/minute counts)] K --> S3[(Raw event log, S3)] S3 --> BATCH[Nightly batch reconciliation job] BATCH --> LEDGER[(Billing ledger - source of truth)] AGG --> API[Query API] API --> DASH[Advertiser dashboard]
This is a Lambda-architecture shape: a fast streaming path (Kafka → Flink → aggregate store) for near-real-time dashboards, and a slower, exact batch path (raw log → nightly reconciliation) that produces the number that actually gets billed. The streaming path is allowed to be slightly wrong; the batch path is not.
- Ingest API does minimal validation (schema, rate-limit per publisher) and writes straight to Kafka — decouples ingestion rate from processing rate.
- Kafka is the durable, replayable buffer. Partition by
ad_id(orcampaign_id) so all events for one ad land on the same partition, giving per-key ordering for free. - Flink consumes, deduplicates (client-generated event id), windows, and aggregates. Chosen over a hand-rolled consumer because it gives you checkpointed state, exactly-once sinks, and watermark-based late-event handling out of the box — reinventing this is the classic mistake.
3. Windowing and dedup trade-offs
| Approach | What it does | Pros | Cons |
|---|---|---|---|
| Tumbling window (1 min) | Fixed, non-overlapping buckets | Simple, cheap, matches “clicks per minute” query | An event 2ms late for its window is dropped or needs a lateness allowance |
| Sliding window | Overlapping windows, recomputed per event | Smooth trend lines | Higher compute — every event touches multiple windows |
| Session window | Groups by activity gaps | Good for “user session” metrics, not raw click counts | Not the right shape for billing aggregates |
| Watermarks + allowed lateness | Tracks event-time progress, holds windows open briefly past the boundary | Handles clock skew and network delay correctly | Have to pick a lateness bound — too short drops real events, too long delays results |
Dedup: mobile SDKs and flaky networks mean the same click can be sent twice. Every click event carries a client-generated idempotency key; Flink’s keyed state (or a Redis/RocksDB-backed dedup set with a TTL matching the allowed-lateness window) drops duplicates before they hit the counter. This is the same idempotency principle as everywhere else in distributed systems — you can’t get exactly-once delivery, so you get exactly-once processing by deduplicating at the consumer.
4. Deep dive — exactly-once counting and the batch/stream reconciliation
The hardest part isn’t the windowing, it’s making sure the number that gets billed is exact even though the streaming path is best-effort:
- Streaming path writes running per-minute counts to the aggregate store using Flink’s checkpointing + a transactional/idempotent sink (e.g. upsert by
(ad_id, minute)key) so a job restart after a crash doesn’t double-count already-processed events — Flink replays from the last checkpoint offset in Kafka, and the idempotent upsert makes replay safe. - Raw event log (every click, unaggregated) is durably written to S3/cold storage straight from Kafka via a sink connector, kept for the full billing dispute window (30–90 days is typical).
- Nightly batch job re-reads the full raw log for the day, does an exact groupby/count (Spark or similar), and writes that as the billing ledger — overwriting or reconciling against the streaming aggregate. If the two disagree beyond a small tolerance, alert — that’s usually a sign of a dedup bug or a watermark misconfiguration in the streaming job.
This is exactly Zalando’s approach for ad event joins: a near-real-time stream path for ad-server logic with events joined inside a bounded window (they use 15 minutes — interactions arriving later are dropped from the real-time path but still land in the batch/billing pipeline), separate from the batch pipeline that’s the actual source of truth for billing and campaign reporting.
5. What real systems do today
- Uber built its real-time ad event pipeline on Apache Flink + Kafka + Pinot, consuming raw ad interaction events from Kafka, using Flink’s windowing to bucket and aggregate by minute with exactly-once guarantees end-to-end, and serving the aggregates from Apache Pinot for low-latency OLAP-style dashboard queries.
- Zalando migrated their ad event processing from a homegrown stateful join service to Flink, describing their system as a “near-classical lambda architecture”: a near-real-time stream for ad-serving logic joined within a bounded (15-minute) window, and a separate batch pipeline that is the actual source of truth for billing and campaign reporting.
- Cardinality estimation: for “unique users reached” (as opposed to raw click counts, which must be exact), HyperLogLog is the standard tool — a probabilistic structure that estimates distinct counts in a small fixed memory footprint (a few KB) with roughly 1–2% error, which is fine for a “reach” metric but never acceptable for a billed click count.
- Kafka partitioning by
ad_id/campaign_idand Flink keyed state are the near-universal building blocks across these write-ups; the differentiator between vendors is mostly the serving layer (Pinot vs. Druid vs. ClickHouse) for the dashboard query path.
6. Scaling & failure
- Hot ad/campaign (a viral ad gets 100x normal clicks) → a single Kafka partition becomes a bottleneck. Fix: partition key includes a salted sub-key (
ad_id + hash(user_id) % N) for hot keys, with a merge step in the aggregation job — the same hot-partition fix used in sharded counters generally. - Flink job falls behind (backpressure from a downstream slow sink) → checkpoint intervals lengthen, watermark lag grows, dashboard staleness increases. Fix: scale out Flink task parallelism, and make sure the aggregate-store sink is the actual bottleneck target for scaling, not the Flink compute itself.
- Aggregate store hot key (one popular ad’s minute-counter updated by many parallel tasks) → contention on a single row. Fix: local pre-aggregation inside each Flink task (combine before shuffle) so only one update per task per window hits the store, not one per event.
What happens when Kafka (the ingest buffer) dies: this is the single point of failure to name proactively. Kafka itself is normally replicated (RF=3) so a broker loss doesn’t lose data, but if the whole Kafka cluster is unreachable: the ingest API should buffer client-side or in a local durable queue for a short window and retry, then start shedding load (return 503, let ad-serving clients retry with backoff) rather than blocking — an ad click pipeline can tolerate a few minutes of ingestion delay far better than it can tolerate dropped clicks. If Kafka data is genuinely lost (rare, misconfigured retention), the batch billing numbers for that window are simply short — this is why the raw click should also be logged synchronously at the edge/CDN layer as a second durable copy in high-value pipelines.
Interview follow-ups
- “Why not just increment a counter in Redis on every click?” — Fine for a rough real-time dashboard, wrong for billing: no replay, no exactly-once guarantee, no audit trail if a customer disputes an invoice.
- “How do you handle a click that arrives 20 minutes late because of a flaky mobile network?” — Watermark + allowed lateness in the streaming job for near-real-time views; the nightly batch job over the immutable raw log is exact regardless of arrival time, so it’s the true source of billing numbers.
- “How do you avoid double-counting a retried click?” — Idempotency key generated client-side, deduplicated in Flink’s keyed state (or a TTL’d dedup store) before the count increments.
- “Streaming aggregate says 10,412 clicks, batch reconciliation says 10,398. What do you do?” — Trust the batch number for billing, alert on the delta, and investigate — usually a watermark/lateness misconfiguration or a dedup TTL that expired too early.
- “How would you detect click fraud in this pipeline?” — A separate stream-processing stage (or the same Flink job) flags anomalous patterns (many clicks from one IP/device in a short window) and either drops them pre-aggregation or tags them for the batch job to exclude — fraud filtering has to happen before the number is billed, not after.
- “What’s the difference in how you’d compute ‘total clicks’ vs. ‘unique users reached’?” — Total clicks is an exact count (must reconcile to the raw log); unique reach is a cardinality estimation problem where HyperLogLog’s small error is acceptable.
Sources: Real-Time Exactly-Once Ad Event Processing with Apache Flink, Kafka, and Pinot — Uber Engineering · From Homegrown to Flink: Migrating a Stateful Ad Event Join at Scale — Zalando Engineering · How to Build a Real-Time Advertising Platform with Apache Kafka and Flink — Kai Waehner · Design an Ad Click Aggregator — Hello Interview