Course Content
System Design Interview
31 sections · 71 lessons
Ad Click Aggregation: reconciliation, hot ads and follow-ups
Every mechanism in Exactly-once aggregation is a defence against a known failure. Reconciliation is the defence against the failures you did not anticipate, and it is what makes the numbers trustworthy.
This lesson finishes the design: the batch path that makes the invoice right, how to recover from a bug, how to scale when one advertiser dominates, and the follow-up questions interviewers use to close the round.
The nightly recomputation
Each night, a batch job reads the raw events for the previous day from object storage and recomputes every aggregate from scratch, using logic written independently of the streaming path where practical. It then compares its output against what the stream produced.
Three outcomes:
- Match within tolerance. Log the comparison and move on. The stream numbers stand.
- Small mismatch, typically from late events routed to the side output. Overwrite the aggregated rows with the batch figures and record the delta. This is expected and routine.
- Large mismatch. Alert a human. A discrepancy above a threshold — say 0.5% on any ad with meaningful volume — means a bug, not a straggler, and no invoice should go out until someone has looked.
The batch figures, not the stream figures, are what billing uses. The stream exists to make the dashboard fast; the batch exists to make the invoice right.
This is the lambda architecture, described plainly
The pattern of a fast approximate streaming path beside a slow authoritative batch path over the same raw data has a name. The name matters less than the reason: the two paths fail differently. A bug in windowing logic affects the stream; a bug in a batch join affects the batch. When two independent implementations over the same input agree, the number is very likely right. When they disagree, you have found something.
Be honest about the cost. Two implementations of the same aggregation means two codebases to maintain and a real risk that they drift apart because someone changed one. The "kappa" alternative — a single streaming implementation, with recovery done by replaying the same stream code over historical data — removes the duplication and keeps the replay capability, at the cost of losing the independent cross-check.
Recommendation: for billing, keep both paths and accept the duplication, because the cross-check is the entire point. For non-financial analytics, a single streaming implementation with replay is enough and much cheaper to run. Say which regime you are in.
Replay after a bug
The aggregation logic had a defect for six hours — a country code was mapped wrong, so counts were attributed to the wrong region. The fix:
- Deploy the corrected aggregation code.
- Delete or mark superseded the aggregated rows for the affected window.
- Run the batch job over raw events for those six hours with the corrected logic.
- Write corrected rows, and emit a correction record to any downstream system that already consumed the wrong numbers.
Step 4 is the one candidates forget. If an invoice already went out, correcting the database does not correct the invoice — a credit note has to follow. Naming that consequence demonstrates you have thought past the pipeline into the business.
Partitioning, and the hot ad
With correctness settled, scale. Partition the click log by ad ID so every event for an ad lands on one partition and one aggregator instance owns its window state. Clean, until one ad dominates.
A brand's launch campaign takes 20% of all clicks. At 50,000 clicks per second peak that is 10,000 per second into a single partition, while sibling partitions handle a few hundred. One aggregator instance saturates, its lag grows, and its windows close late — while the cluster as a whole is idle.
Three repairs, in increasing order of complexity:
Salting. Partition by (ad_id, random(0..N-1)) for keys known to be hot, splitting one ad across N partitions. Each partition produces a partial count; a second aggregation stage sums the N partials per window. Costs one extra stage and a little latency, and it is the standard answer.
Two-stage aggregation for everything. Pre-aggregate locally on each ingestion node — count clicks per ad per second in memory — and emit counts rather than events. A 10,000-per-second ad becomes one record per second per node instead of 10,000. This reduces the downstream volume by orders of magnitude for hot keys and does nothing for cold ones, which is exactly the right shape.
Dedicated capacity. Route known-hot advertisers to their own topic and cluster. Operational rather than elegant, and real systems do it.
This is the same hot-key problem as the celebrity fan-out in Section 13 (Design a News Feed System), the hot partition in Section 21 (Design a Distributed Message Queue), and the global leaderboard in Section 27 (Design a Real-Time Gaming Leaderboard). The patterns behind the twenty-five systems collects them.
Click fraud
A meaningful fraction of raw click traffic is not human. Filtering is a system of its own, but name the layers:
- Cheap synchronous filters at ingestion: known bot user agents, data-centre IP ranges, impossible timing such as a click before the impression rendered.
- Near-real-time rules: more than N clicks on one ad from one device in a minute.
- Offline models scoring sessions after the fact, feeding a nightly adjustment.
The architectural consequence is that "clicks" is not one number. Keep raw, filtered, and billable counts as separate columns, because advertisers dispute the difference and you must be able to show your work.
Multi-dimensional aggregation
Advertisers want clicks by ad × country × device × hour. Three dimensions plus time is a small cube; the danger is unbounded growth as dimensions are added, since the row count is roughly the product of the cardinalities.
Two options. Precompute the common combinations — ad alone, ad × country, ad × device — which is fast to query and grows combinatorially with dimension count. Or store the finest grain and roll up at query time, which is flexible and slower. Recommend precomputing the handful of combinations the dashboard actually shows and keeping the finest grain for ad-hoc queries against a columnar store, which is what the query pattern justifies.