Course Content
System Design Interview
31 sections · 71 lessons
Ad Click Aggregation: requirements, scale and the pipeline
"Design an ad click event aggregation system." Clicks arrive; counts come out. The difficulty is entirely in one word that the prompt does not contain: billing.
This lesson scopes the problem, sizes it, and builds the pipeline. The second lesson goes deep on the two hard parts — windowing with late events, and exactly-once counting — and the third covers reconciliation, scaling and the follow-up questions.
Why this is not a counting problem
Every click has a price. Aggregated click counts are what advertisers are invoiced from and what publishers are paid from. A 1% over-count is fraud; a 1% under-count is lost revenue. So the correctness bar is far above what a dashboard needs, and that single fact drives Exactly-once aggregation and Reconciliation and recovery.
At the same time, the product team wants near-real-time numbers, because an advertiser running a campaign wants to see spend within a minute, not tomorrow. Fast and exact pull against each other, and the resolution — a fast path plus a slow correcting path — is the section's central idea.
The questions that shape everything after
- What aggregation windows are needed? Per-minute counts, per-hour, per-day? The finest window drives the state size in the stream processor.
- How fresh must the numbers be? Under a minute, or is five acceptable? A one-second requirement changes the architecture; a five-minute one relaxes it a lot.
- Is exact correctness required? For billing, yes. Ask anyway, because if the answer is "approximate is fine for the dashboard" you can serve the dashboard from a cheap path.
- How late can events arrive? A phone offline in a lift buffers clicks and sends them twenty minutes later. Mobile apps can be far worse — hours or days.
- What dimensions must the counts be broken down by? Ad ID alone, or ad × country × device × hour? Each added dimension multiplies the output rows.
- How long are raw events retained? This is the recomputation window, and Reconciliation and recovery shows why it is the most important retention decision in the system.
The assumptions this section uses
| Question | Assumption |
|---|---|
| Windows | Per-minute counts, rolled up to hour and day |
| Freshness | Under one minute for the dashboard |
| Correctness | Exact for billing; the dashboard may be provisional |
| Late events | Up to 1 hour normal; up to 24 hours tolerated |
| Dimensions | Ad ID, plus country and device type |
| Raw retention | 30 days, so any bug within a month can be recomputed |
Requirements and scale
With the assumptions agreed, write the requirements down and put numbers on them.
Functional requirements
- Ingest click events, each carrying at minimum an event ID, an ad ID, a user or device ID, a timestamp, and context such as country and device type.
- Aggregate clicks per ad per minute, and roll up to hour and day.
- Serve queries: clicks for ad X over a range, top N ads by clicks in the last M minutes.
- Support recomputation of any window from raw events.
Non-functional requirements
- Correct counts for billing — the hard requirement.
- Dashboard results within about one minute of the click.
- Survive a stream-processor crash without losing or duplicating counts.
- Handle a peak several times the average without dropping events.
The arithmetic
Assume 1 billion clicks per day. Invented, but the right order of magnitude for a large ad network.
Average rate. 1,000,000,000 ÷ 100,000 seconds per day (the rounding from Section 3) = 10,000 clicks per second.
Peak rate. Traffic is not flat. Evening peaks in a large market run roughly 5× the daily average, so design for 50,000 clicks per second.
Event size. Event ID (16 bytes), ad ID (8), user ID (8), timestamp (8), country (2), device (1), plus JSON or Avro framing overhead. Call it 0.5 KB on the wire.
Ingest bandwidth. 50,000/s × 0.5 KB = 25 MB/s at peak; 5 MB/s average.
Raw storage. 1 billion × 0.5 KB = 500 GB per day. Over 30 days of retention that is 15 TB, and at replication factor three, 45 TB. That is the price of being able to recompute, and it is affordable — which is the point.
Aggregated storage. Suppose 1 million active ads. Per-minute counts across three dimensions: 1 million ads × 1,440 minutes = 1.44 billion rows per day if every ad were active every minute. Realistically only a few per cent are, so call it 50 million rows per day at roughly 50 bytes = 2.5 GB per day. Aggregates are 200× smaller than raw. Also worth noticing: the aggregated store is small enough that a relational or wide-column database handles it comfortably, which is not true of the raw stream.
Query load
Advertiser dashboards: 100,000 advertisers, a few checking per minute, call it 500 queries per second, each hitting the aggregated store over a bounded range. Trivial next to the write path. As in Section 22 (Design a Metrics Monitoring and Alerting System), the write side dominates.
The pipeline
Now the architecture. Build up from the simplest thing that could work, and break it.
The naive version, and why it fails
The click endpoint receives a click and executes UPDATE ad_counts SET clicks = clicks + 1 WHERE ad_id = ? AND minute = ?.
At 50,000 clicks per second this fails three ways at once. It is 50,000 row updates per second against a hot row per popular ad, so lock contention serialises everything. There is no raw record, so a bug is unrecoverable. And a database outage drops clicks on the floor, because the web tier has nowhere to put them.
Each failure points at a component.
The four stages
1. Ingestion into a distributed log. The click endpoint does one thing: validate the event, assign an event ID if the client did not, and append it to a log topic. Kafka is the common implementation. The append is a few milliseconds and the log absorbs the peak, so the endpoint stays fast and the downstream never applies backpressure to users. Retention on this topic is short — a few days — because the raw archive is separate.
2. Raw event archive. A consumer copies every event from the log to cheap object storage, partitioned by hour. This is the 500 GB/day, 30-day store, and it exists solely so reconciliation can recompute.
3. Stream aggregation. A stream processor consumes the log, groups events into per-minute windows keyed by (ad ID, country, device), counts them, and writes the results to the aggregated store when the window closes. This is where windowing and exactly-once aggregation live.
4. Aggregated store and query service. A database holding minute, hour, and day rollups, fronted by an API. Because the rows are small and the access pattern is a range scan by ad ID and time, a wide-column or time-partitioned relational store both work; recommend whichever your team already runs, and say so, because operational familiarity is a real criterion.
Why keep raw events when aggregates are what get read
Because every number the fast path produces is provisional until something independent confirms it. The slow path is that something. It also gives you: replay after an aggregation bug, the ability to add a new dimension retroactively, and an audit trail when an advertiser disputes an invoice. None of those are possible from aggregates alone.