Mid to Senior Engineer

System Design Interview Prep

A structured path from the interview framework through core concepts, key technologies and patterns to eighteen full problem breakdowns, each with diagrams and weak, solid and excellent answers to every deep dive.

Chapter 34 of 36Problem breakdowns · Design Ad Click Aggregation

Design Ad Click Aggregation

An advertising platform needs to know, within seconds to minutes, how many times each ad was clicked, by campaign, by region and by hour. The numbers drive billing, dashboards and bidding, so they must be accurate. It looks like counting, and it is the standard interview problem for stream processing: very high event volume, exactly-once effects, late and duplicate events, time windows, and the need to reconcile a fast approximate answer with a slower exact one.

The chapter follows the usual shape: understand the problem, set up the interface, build the high-level design, then go deep on the questions interviewers use to separate levels.

1. Understanding the problem

Click events flow in from servers that handle ad clicks. We aggregate them into counts per ad over time, and serve queries such as "clicks for ad 42 in the last hour, per minute" and "top 100 ads by clicks in the last 5 minutes".

Functional requirements

Core:

  1. Ingest ad click events at very high volume.
  2. Aggregate clicks per ad over time windows (per minute, hour, day), with filters such as campaign and region.
  3. Serve queries on the aggregates with low latency.

Confirm in or out: impressions as well as clicks, conversion attribution, fraud detection, real-time bidding, and long-term reporting. A sensible opening: "I will design click ingestion, deduplicated aggregation by ad and minute, and a query API, with a reconciliation path, and treat attribution and fraud scoring as extensions."

Non-functional requirements

  • Correctness. Counts drive money, so duplicates and losses must be controlled, and the numbers must be reconcilable and auditable.
  • Timeliness. Dashboards update within seconds to a minute or two.
  • Scale. A billion clicks a day is plausible for a large platform.
  • Fault tolerance. A failure must not lose or double count events.
  • Late data. Events arrive out of order and sometimes late, and the system must define how to handle them.

Estimation

QuantityCalculationResult
Clicks per dayabout 11,500 per second average
Peak5 times averageabout 58,000 per second
Event sizeabout 100 bytes (ad id, user, time, region, click id)5.8 MB per second at peak
Raw storage Babout 100 GB per day, 36 TB per year
Aggregates1 million active ads 1,440 minutesabout 1.4 billion rows per day, if every ad had every minute

What the numbers say. The ingest rate is high but well within a partitioned log. Raw storage is large enough to need object storage with compression. The aggregate table is sparse (most ads are not clicked every minute), so real volume is much lower, and what matters is the design of the pipeline: partitioning by ad, windowing by event time, and exactness. Treat these as assumptions you state and adjust.

2. The set up

Core entities

  • Click event: unique click identifier, ad identifier, campaign, user or device, timestamp, region, source.
  • Aggregate: ad identifier, window start, window length, count (and other measures).

Interfaces

Ingestion is a write-only path from click servers into a log. Query:

GET /v1/ads/{ad_id}/clicks?from=...&to=...&granularity=minute          -> [{ window_start, count }, ...]
GET /v1/campaigns/{id}/clicks?group_by=region&from=...&to=...
GET /v1/top-ads?window=5m&n=100

3. High-level design

Split the system into stages with a durable log between them:

  1. Click servers emit each click as an event into a partitioned raw click log, keyed by ad identifier.
  2. A raw archive copies every event to object storage, for replay, audit and batch recomputation.
  3. A stream processor consumes the log, deduplicates and enriches events (bot filtering, region lookup), and aggregates them in time windows per ad.
  4. Aggregates are written to an aggregates store, which a query API reads.
  5. A batch reconciliation job recomputes aggregates from the raw archive and corrects the store.
<!--fig:pipeline-->
correct Clickservers Raw click log Dedupe + enrich click id, bot filter Stream aggregator windows by ad id Aggregates store per ad per minute Raw archive object storage Batch reconcile recompute, compare Query API Figure 1. Click aggregation: dedupe and enrich, aggregate in time windows, serve fast queries, and reconcile against raw data.

The log between stages decouples producers from consumers, absorbs bursts and allows replay. The raw archive is the source of truth, and the stream path is a fast derived view of it.

4. Potential deep dives

Deep dive 1: How do you count without losing or duplicating events?

The challenge. Servers retry, networks redeliver, and consumers crash and restart. Each can cause a click to be counted twice or not at all, and billing depends on the count.

Weak: increment a counter in a database for every event. A retry or a restart double counts. A crash between consuming and incrementing loses or repeats work. A single hot ad becomes a hot row.

Solid: at-least-once delivery with deduplication. Give every click a unique identifier at the source. The processor remembers recently seen identifiers (in a store with expiry, sized to the retry window) and drops repeats. Acknowledge the log position only after the aggregate is safely written.

Excellent: make the effect exactly-once end to end. State the guarantee precisely: the log gives at-least-once delivery, and you build exactly-once effects. Options:

  • Transactional processing. The stream processor checkpoints its position in the log and its aggregate state together, atomically, so after a crash it restores both and resumes without double counting. This is the model of modern stream processing frameworks.
  • Idempotent writes. Write each aggregate as an upsert keyed by (ad, window), computing the count for the window from the events, so rewriting the same window with the same events yields the same value. Use the click identifier set, or recompute the window from the log for the window.
  • Deduplicate at the source of identity. Generate the click identifier on the server that records the click, derived deterministically from the request, so the same physical click always has the same identifier, even if recorded twice.

Add the batch path as a safety net, which makes any residual streaming error correctable.

Deep dive 2: Windows, event time and late events

The challenge. "Clicks per minute" requires deciding which minute a click belongs to, and clicks do not arrive in order.

Weak: bucket by arrival time. A click that happened at 00:59 but arrives at 01:01 is counted in the wrong minute, and delays or backlogs shift counts between buckets, so the numbers are wrong whenever the system is slow.

Solid: bucket by event time. Use the timestamp recorded when the click happened. A click belongs to the window containing its event time, whenever it arrives.

Excellent: watermarks and a lateness policy. Event-time windows must decide when to close. A watermark is the processor's estimate of how far event time has progressed: for example, the maximum event time seen minus an allowed lateness (say five minutes). A window is closed and emitted when the watermark passes its end.

<!--fig:windows-->
Tumbling window (1 min): events counted once, in the window of their event time window 0 window 1 window 2 window 3 Late event: happened at 00:59, arrives at 01:04. Watermark decides whether window 0 is still open. event time 00:59 (belongs to window 0) Watermark = max event time seen minus allowed lateness (say 5 minutes).Close a window when the watermark passes its end. Later events go to a correction path or are counted by the batch job. Figure 2. Windows are defined on event time, with a watermark that trades latency against completeness.

Then decide what to do with events that arrive after their window closed:

  • Drop them (simple, and loses data).
  • Update the already emitted aggregate as a correction, for as long as lateness is allowed (accurate, and downstream consumers must handle updates).
  • Send them to a side output for the batch job to include, which resolves the discrepancy in the next reconciliation.

State the trade-off: a longer allowed lateness gives more complete counts and delays the final result, and the right setting depends on how late events really arrive in your system, which you measure.

Know the window types: tumbling (fixed, non-overlapping), sliding (overlapping, for moving averages) and session (gap-based). For per-minute clicks, tumbling windows are the answer.

Deep dive 3: How do you scale the aggregation?

The challenge. Tens of thousands of events a second, with some ads far more popular than others.

Weak: one aggregator. A single process cannot keep up, and it is a single point of failure.

Solid: partition by ad identifier. Key the log by ad id, so all clicks for one ad go to one partition and one processor instance, which keeps that ad's window state locally. Scale by adding partitions and instances.

Excellent: handle hot ads with two-stage aggregation. A very popular ad can overload one partition. Use two stages: first, aggregate partially in many parallel instances, keyed by (ad, a random shard number), producing partial counts per window. Second, merge the partials per ad into the final count. The first stage spreads the load, and the second handles a small number of partial results. Pre-aggregation in the click servers themselves (counting locally for a second before sending) can cut event volume further, at the cost of losing per-click detail, so keep the raw events for the archive. Monitor lag per partition and skew across keys.

Deep dive 4: Storing and serving aggregates

The challenge. Queries ask for a time range for an ad or group, and must return quickly.

Weak: query the raw events at read time. Scanning billions of events per query is far too slow.

Solid: a store of precomputed aggregates, keyed by ad and time. A time-series or wide-column store keyed by (ad id, window start) makes a range read a single partition scan. Create separate rollups per granularity: minute, hour, day.

Excellent: tiered rollups, dimensional aggregates and caching. Compute rollups hierarchically: hours from minutes, days from hours, so long-range queries read few rows. For dimensions such as region and campaign, precompute the combinations that dashboards use, and compute the rest from raw data on demand in an analytical store. Choose the retention for each tier: minute data for days, hour data for months, day data for years. Cache popular queries. Pick a storage engine designed for fast aggregation over many rows for ad-hoc analysis, and use the key-value store for fixed dashboards.

For approximate measures such as the number of unique users who clicked, use a probabilistic sketch that estimates distinct counts in a few kilobytes, and that can be merged across windows and shards. Exact distinct counting needs memory proportional to the number of users, which does not scale.

Deep dive 5: Reconciliation and the batch path

The challenge. The stream path can still be slightly wrong: a late event beyond the lateness limit, a bug, or a processor failure. Billing needs exact numbers.

Weak: trust the streaming result. Small errors accumulate unnoticed and show up as billing disputes.

Solid: a daily batch recomputation. Recompute the aggregates from the raw archive for the previous day, and overwrite or compare with the streaming result.

Excellent: stream for speed, batch for truth, with a comparison. This is the lambda architecture: a speed layer gives low-latency approximate results, and a batch layer gives the authoritative result later. The alternative, the kappa architecture, uses a single streaming path and fixes errors by replaying the log through the same code, which avoids maintaining two implementations. Describe both, and say that the key is to keep the raw events immutable and replayable. Automate a comparison between stream and batch results, alert when they diverge beyond a threshold, and treat the batch result as final for billing. Make recomputation idempotent, partitioned by time, and re-runnable.

Deep dive 6: Fraud and invalid traffic

Not every click is a real user. Bots, click farms and accidental double clicks inflate counts, and advertisers should not pay for them. Filter in the pipeline: rule-based checks (too many clicks per device in a short time, known bot signatures) and statistical or model-based scoring, marking clicks as valid or invalid rather than deleting them, so the decision can be audited and revised. Run a slower, more thorough analysis in the batch path, and adjust billing afterwards. This is a policy-heavy area, so say you would work with specialists on the rules and on how disputes are handled.

Deep dive 7: Failure and operations

  • A processor instance dies. Another takes over its partitions and restores state from the last checkpoint, then reprocesses events since, with the exactly-once mechanism preventing double counts.
  • The log is unavailable. Click servers buffer locally for a short time and retry, and the system must not drop clicks silently. Alert on buffer growth.
  • Consumer lag grows. Scale processors and investigate hot keys. Set alerts on the lag, since lag is a direct delay in the dashboard.
  • Schema changes. Versioned event schemas with compatibility rules, so producers and consumers can change independently.
  • Backfill and replay. Re-run a time range from the archive after fixing a bug, writing to a new version of the aggregates and switching when verified.
  • Observability. Track events in, events deduplicated, late events, watermark delay, end-to-end latency and the stream versus batch difference.

5. What is expected at each level

Mid-level. You propose a queue for ingestion, workers that aggregate counts and a database for results, and notice that duplicates must be handled.

Senior. You partition by ad, window by event time with a watermark, achieve exactly-once effects through checkpoints or idempotent writes, add a raw archive with batch reconciliation and discuss storage and rollups.

Staff. You discuss lambda versus kappa, hot-key two-stage aggregation, approximate distinct counts, invalid traffic handling, the economics of lateness, schema evolution, replay and backfill, and how finance consumes the numbers with an audit trail.

6. Interview questions and model answers

Q: How do you avoid double counting? Each click has a unique identifier assigned at the source. The pipeline gives at-least-once delivery, and I make the effect exactly-once by deduplicating on the identifier and by checkpointing processor state together with its log position, so a restart resumes without recounting.

Q: What time do you use for windows? Event time, the time the click happened, not the time it arrived. A watermark decides when a window closes, with an allowed lateness based on measured delays. Late events beyond that go to a correction path or the batch job.

Q: How do you handle one extremely popular ad? Two-stage aggregation: partial counts keyed by ad and a random shard spread the load across instances, then a second stage merges the partials. I also monitor for key skew.

Q: How do you get exact numbers for billing? The stream gives fast results, and a batch job recomputes from the immutable raw archive. The two are compared automatically, and the batch result is authoritative for billing.

Q: How do you count unique users cheaply? A probabilistic sketch for distinct counts, which uses a few kilobytes per counter, can be merged across windows and shards, and gives a small bounded error.

Q: What if a processor crashes? Another instance restores the last checkpoint of state and log position and resumes. Idempotent writes and checkpointing together prevent duplicates.

7. Common mistakes

  • Incrementing a database row per event.
  • Bucketing by arrival time instead of event time.
  • No unique click identifiers, so duplicates cannot be removed.
  • No raw archive, so errors cannot be corrected.
  • Trusting the stream result for billing with no reconciliation.
  • One partition for all events, or a hot partition for a popular ad.
  • Deleting suspicious clicks instead of marking them.
Header Logo