Problem & requirements
Real-time bidding auctions buy ad slots in milliseconds, and click counts drive both billing and bidding decisions, so aggregated click data must be timely and correct. Inputs are click events with an ad ID, timestamp, user ID, IP and country, arriving as log files on ad servers. There are about 1 billion clicks per day across 2 million ads, growing 30% yearly.
The system must report the number of clicks for an ad in the last M minutes, return the top 100 most clicked ads every minute, and support filtering by attributes such as country or IP. Non-functionally, correctness is the top priority because money depends on it; duplicate and late events must be handled; the pipeline must tolerate partial failures; and end-to-end latency should be a few minutes at most.
Back-of-the-envelope
1 billion clicks per day divided by 86,400 seconds is about 11,600 per second; round to 10,000 QPS average. Assuming peaks are five times the average gives roughly 50,000 QPS. Check: 10,000 x 100,000 seconds would be 1 billion, and a day is a bit shorter than 100,000 seconds, so the true average is slightly above 10K, consistent with 11.6K.
If each event is about 0.1 KB, a day's raw data is 1B x 0.1 KB = 100 GB, and a month is about 3 TB. That fits comfortably on a scalable distributed store. Aggregated data is far smaller, since 2M ads times 1,440 minutes per day is under 3 billion small rows at the finest grain, and most queries only need recent windows.
Raw versus aggregated data and storage
Store both. Raw events are the source of truth: if a bug corrupts aggregates, they can be recomputed, and auditors can trace a bill back to individual clicks. But querying raw events for every dashboard is far too slow. Aggregated data, such as clicks per ad per minute, makes queries fast and is what serving systems read. Raw data can move to cold storage after a while; aggregated data stays hot.
The workload is write-heavy for raw data and read-heavy with time-range access for aggregates. A wide-column store such as Cassandra handles the raw write rate and scales horizontally; aggregated rows keyed by ad ID and minute, plus filter dimensions, also fit a column store or a time-series-friendly database. A relational database would struggle to absorb the raw write rate at peak.
High-level design
A log watcher on each ad server tails click logs and publishes events to a Kafka topic. The aggregation service consumes them and computes per-minute counts and top-N results. It publishes results to a second Kafka topic, from which a database writer stores them. Raw events are also written to the raw database by a separate consumer.
Why a second queue rather than writing directly? It decouples aggregation from storage speed, lets results be consumed by several systems, and, crucially, gives a clean place to achieve end-to-end exactly-once semantics with transactional writes. The aggregation service itself is modelled as a MapReduce-style DAG: map nodes read and clean events and route them by ad ID, aggregate nodes count clicks per ad per minute in memory, and reduce nodes merge partial results, for example combining each aggregator's local top 100 into a global top 100.
Deep dive: event time, watermarks and windows
Timestamps can come from when the click happened on the client (event time) or when the aggregator processed it (processing time). Processing time is simple but wrong whenever events are delayed by network issues or backlogs, which would misattribute clicks to the wrong minute and hence the wrong bill. Event time is accurate but means events can arrive late, after their window has been reported.
A watermark extends each window's wait by a grace period, say 15 seconds, before closing it. A longer watermark catches more late events but adds latency. Events later still are handled by an end-of-day reconciliation rather than holding windows open forever. Aggregation uses tumbling windows (fixed, non-overlapping one-minute buckets) for per-minute counts, and sliding windows for questions such as 'top N ads in the last M minutes', where the window moves forward continuously.
Deep dive: exactly-once delivery
Duplicates in a billing pipeline mean overcharging advertisers, and losses mean undercharging, so the system needs effectively exactly-once processing. The danger point is the gap between computing a result and recording the consumed offset. If an aggregator commits its upstream offset before the result is safely downstream and then crashes, events are lost; if it sends the result first and crashes before committing, the restarted node reprocesses and duplicates.
The fix is to make the result write and the offset commit atomic. One approach stores the offset in external storage alongside results and only acknowledges upstream once the downstream write succeeds; another uses a distributed transaction, such as Kafka's transactional producer, so the output messages and the consumer offset commit succeed or fail together. Deduplication by event ID at the edge also guards against client-side duplicates and click fraud.
Scaling, hotspots and fault tolerance
Scale the queues by adding partitions up front, keyed by ad ID so all of an ad's events reach the same aggregator; adding partitions later reshuffles keys, so do it at quiet times. Aggregation nodes scale horizontally, and within a node, work can be split across threads or delegated to a resource manager such as YARN. A hotspot ad with far more clicks than others can be given extra aggregation nodes: split its events across several workers and merge their partial counts in the reduce step.
Aggregators keep state in memory, such as counts and top-N heaps, so recovery by replaying from the earliest uncommitted Kafka offset could take a long time. Periodic snapshots of node state along with the upstream offset let a replacement node load the snapshot and replay only the events after it.
Monitoring, reconciliation & alternatives
Monitor latency at each stage, Kafka consumer lag, and resource use on aggregators. Because correctness is paramount, run a nightly reconciliation: a batch job recomputes aggregates from raw events, sorts them and compares them with the streaming results, flagging mismatches. A recalculation service can replay raw data through a separate aggregation path to repair historical results after a bug fix without disturbing live traffic.
Architecturally, keeping a streaming path and a batch path is the lambda architecture; using only a replayable stream for both live results and recomputation is the kappa architecture, which this design leans toward, since it maintains one code path. An alternative is to store raw clicks in an OLAP engine such as Apache Druid, ClickHouse or Elasticsearch, or in Hive for batch, and aggregate at query time. That is simpler to build but less predictable at this scale and harder to make exactly-once for billing.
Key numbers
Key terms
- Event time
- The timestamp at which the click actually occurred on the client.
- Processing time
- The timestamp at which the aggregation service handled the event.
- Watermark
- A grace period after a window ends during which late events are still accepted.
- Tumbling window
- Fixed-length, non-overlapping time buckets, each event belonging to exactly one.
- Sliding window
- A fixed-length window that advances continuously, so windows overlap.
- Kappa architecture
- A design that uses a single replayable stream pipeline for both real-time and historical processing.
- Lambda architecture
- A design that runs separate batch and streaming pipelines and merges their results.
- Reconciliation
- Comparing streaming aggregates with a batch recomputation from raw data to detect discrepancies.
Common mistakes
- Using processing time and silently attributing delayed clicks to the wrong minute.
- Storing only aggregated data, leaving no way to recompute after a bug.
- Committing upstream offsets before downstream results are durable, losing or duplicating events.
- Ignoring hot ads that overwhelm a single aggregation partition.
- Holding windows open forever for very late events instead of reconciling later.
Further study
- The Dataflow Model (Akidau et al., VLDB 2015)
- MapReduce: Simplified Data Processing on Large Clusters (2004)
- Apache Flink
- Questioning the Lambda Architecture (Jay Kreps, 2014)
- Apache Druid