Ad Click Aggregator System Design: Streams, Late Events, and the Billing-Grade Recount
Quick Overview
A senior engineer's walkthrough of the ad click aggregator system design interview, grounded in real questions asked at Meta, Amazon, Rippling, and Robinhood. Covers the redirect ingest path, event-time windowing with watermarks, the three sources of duplicate clicks, stream enrichment joins, tiered OLAP storage, and the nightly reconciliation layer that makes billing exact.
An ad click aggregator sounds like a counting problem: users click ads, advertisers want to know how many. The interview version is anything but. At hundreds of thousands of clicks per second, every event is money. An undercount means refunding an advertiser; an overcount means billing for clicks that never happened. The question is really a test of whether you can reason about windowed stream aggregation, duplicate delivery, late events, and the batch reconciliation layer that quietly fixes everything the stream got wrong. This guide walks the design end to end, grounded in versions of the question asked at Meta, Amazon, Rippling, and Robinhood.
Key Takeaways
- Split the system into a write path (append-only raw click log) and a read path (pre-aggregated minute buckets). Advertiser dashboards read the stream's output; invoices read a nightly batch recount. Saying this split out loud in the first five minutes reframes the whole interview.
- Exactly-once is a property of the entire pipeline, not a Kafka checkbox. Kafka transactions plus a two-phase-commit sink cover the middle of the pipeline; you still need an impression-scoped idempotency key to kill client retries, and reconciliation to catch what slips through.
- Late clicks get three tiers of handling: allowed lateness inside the window operator, a side output topic for stragglers, and the batch recount as final authority. Know the destination of a click that arrives 5 seconds, 5 minutes, and 5 hours late.
- Hot ads break naive partitioning. A single viral ad can pin one Kafka partition and turn the consumer behind it into the straggler that gates the whole job. Key by
ad_idplus a small random suffix and merge in a second aggregation stage. - Raw events are the source of truth. Every aggregate is disposable and rebuildable from the log; if you can only defend one storage decision in the interview, defend keeping the raw events.
Pin down what a "click" is worth before you draw a single box
Asked at Meta — Design an ad click aggregation service This problem anchored a Meta candidate's loop. The canonical form asks you to build the service that records every click on an ad and lets advertisers query per-ad click totals in near real time. The candidate's offer retrospective called out that grinding low-level detail from books was less useful than practicing the structured walk-through this question demands.
Start with the two questions that shape everything downstream: who reads the counts, and how fresh and how correct do they need to be?
The functional surface is small. Advertisers query "clicks for ad N over time range T," typically at minute granularity for the last day or two and at hourly or daily granularity for history. There is also a second, quieter consumer: billing, which charges per click and cannot tolerate drift.
The scale math is where the design forks. A large ad network sees on the order of 10 billion clicks a day. That averages to roughly 115,000 clicks per second, and peaks (product launches, live sports) can run 3 to 5 times the average, so plan for 500,000 events/second. At about 100 bytes per event, peak ingest is around 50 MB/s and the raw log grows by roughly 1 TB a day. None of these numbers require exotic hardware. What they rule out is anything that does a synchronous database write per click.
Now the correctness requirement, which is the part most candidates state too casually. Dashboards can be approximately right within a minute or two. Invoices must be exactly right, and "exactly" here includes retroactively removing fraudulent clicks discovered hours later. Those are two different consistency contracts, and no single path satisfies both cheaply. The design that follows is really one pipeline wearing two hats: a stream that is fast and slightly wrong, and a batch layer that is slow and authoritative.
The common mistake is skipping this framing and jumping straight to "Kafka, then Flink." The interviewer has seen that opening a hundred times. What distinguishes candidates is naming the dashboard/billing split unprompted, because every later trade-off (delivery semantics, late-event policy, storage tiering) resolves differently for each consumer.
The ingest path: own the click, then get out of the way
Asked at Amazon — Design ad clickstream analytics pipeline This version asks for the full ingestion platform: click events flowing through Kafka into raw and curated S3 zones, queryable with Presto. It wants explicit choices about how the topic is keyed and partitioned, how the event format is serialized and evolved through a registry, and which ordering and delivery guarantees the pipeline makes — a strong hint about where the interviewer expects depth.
There are two ways to capture a click, and the choice is a genuine trade-off, not a formality.
The first is the redirect pattern. The ad's link points at your click server. The user's click hits it, the server validates and enqueues an event, then answers with a 302 to the advertiser's landing page. You are now on the critical path of a human navigating the web, so this server must be deliberately boring: verify a signature, append to Kafka, redirect. Budget under 10 ms of server time. In exchange, you capture essentially every click, because the user cannot reach the destination without passing through you.
The second is a client-side beacon: the browser or SDK fires an event asynchronously and navigates immediately. Lower perceived latency, but lossy. Ad blockers eat beacons, and events queued during page unload are dropped at meaningful rates. Ad networks that bill per click use the redirect; analytics products that merely observe clicks often accept the beacon. Say which one you're building and why.
For the event itself, use a compact binary schema (Avro or Protobuf) registered in a schema registry so producers and consumers can evolve independently:
message ClickEvent {
string impression_id = 1; // minted and signed when the ad was SERVED
string ad_id = 2;
int64 event_time_ms = 3; // client/server click timestamp, not arrival time
string device_class = 4; // web, ios, android
string ip_hash = 5; // for fraud scoring, not identity
}
The field that earns you points is impression_id. It is created when the ad is rendered, signed by the ad server, and echoed back with the click. It gives you an idempotency key (a click counts at most once per impression), a replay-attack defense (an unsigned or expired ID is discarded at the edge), and a natural join key back to the impression log.
Partitioning is the next explicit decision. Keying the Kafka topic by ad_id gives you per-ad ordering and lands all of one ad's clicks on one partition, which makes downstream windowed aggregation local. It also creates the hot-key problem: one Super Bowl ad taking 10% of global traffic pins a single partition, and the consumer behind it becomes the straggler that lags the whole job. The standard fix is a two-stage aggregate. Produce with key ad_id + ":" + rand(0..9), run a first aggregation over the salted keys, then a second, much smaller aggregation that merges the ten partial counts per window. The second stage handles ten records per ad per minute; it is never the bottleneck.

Note the fan-out at Kafka: the raw log lands in S3 in parallel with stream processing, not downstream of it. If the aggregation job has a bad day, the source of truth is untouched.
Windowed aggregation and delivery semantics: watermarks, late clicks, and where duplicates die
Asked at Robinhood — Design Real-Time Analytics Pipeline with Kafka and Flink The prompt hands you the stack (click events in via Kafka, processing in Flink, aggregates out to a warehouse) and grades you on the parts inside: topic partitioning, window semantics, state management, and how the pipeline behaves under failure and backpressure. Being told the tools and still having plenty to design is the tell that windowing and delivery details are the actual content of this question.
The core computation is a tumbling one-minute window keyed by ad_id: for each ad, count clicks whose event time falls in [12:04:00, 12:05:00). The two words doing the work there are "event time."
A click that happened at 11:59:58 might arrive at 12:00:25 because a phone was in a tunnel. If you window by processing time (when the event arrived), that click lands in the wrong bucket, and worse, your buckets smear whenever the pipeline lags. During a backlog replay after a deploy, processing-time windows produce garbage: an hour of clicks crammed into two minutes of windows. Event-time windowing with watermarks is the fix, and interviewers probe it because it separates people who have run a streaming job from people who have read about one.
A watermark is the pipeline's rolling claim: "I believe no event older than T is still in flight." Flink advances it from observed event times minus a bound you choose, say 30 seconds. When the watermark passes a window's end, the window fires and emits its count.
So where does a late click go? Three tiers:
- Within the watermark bound (seconds late): it simply arrives before the window fires. No one notices.
- Within allowed lateness (you configure, say 1 minute): the window has fired but its state is retained. The late event triggers an updated emission, and the sink overwrites the old count by key. This is why the OLAP sink must support overwrite-by-key rather than blind append — with the caveat that in some stores the overwrite is eventual rather than synchronous, which the storage section returns to.
- Beyond allowed lateness: the window state is gone. The event is routed to a side output, a
late-clickstopic that the batch layer consumes. The dashboard stays slightly wrong for that minute; the nightly recount makes it right.

Late events are half of what the Robinhood prompt calls failure handling. The other half is duplicates, and the first thing to say about them is that exactly-once is a pipeline property, not a Kafka setting. Duplicates enter this system in three distinct places, and each needs its own defense. Candidates who answer "we'll use Kafka exactly-once" have covered exactly one of the three.
Client retries. A user double-clicks, or a flaky mobile network resends the request. This duplicate is born before Kafka ever sees it, so no broker setting can help. The defense is the signed impression_id: the aggregation job (or a dedup operator ahead of it) keeps recently-seen impression IDs in keyed state and drops repeats. State for this is bounded because an impression ID is only valid for a short click window anyway, an hour or two, after which the signature is expired and the edge rejects it. Deduplicating "at most one click per impression" also happens to match the billing model, which is a nice thing to point out.
Broker redelivery. A consumer crashes after processing but before committing offsets, restarts, and reprocesses a batch. Kafka transactions plus Flink's checkpoint-aligned two-phase-commit sink close this gap: offsets, state, and sink writes commit atomically. This is the part vendors mean by "exactly-once," and it is real, but it only spans source-to-sink of the streaming job.
Sink double-writes. The window fires, the write lands, the job crashes before the checkpoint completes, and after recovery the window fires again. If the sink is transactional you're covered by the previous mechanism; if it isn't (many OLAP stores aren't), make the write idempotent instead: upsert keyed by (ad_id, window_start). Writing the same count twice is harmless. This is usually the cheapest correct answer, and saying "idempotent upsert beats a distributed transaction here" is a senior move.
| Strategy | Kills which duplicates | Cost | Where it still leaks |
|---|---|---|---|
At-least-once + idempotent upsert by (ad_id, window) | Broker redelivery, sink retry | Cheap; sink needs upsert | Client-born duplicates pass straight through |
| Kafka transactions + 2PC sink ("exactly-once") | Broker redelivery, sink double-write | Latency tied to checkpoint cadence; sink must participate | Duplicates created before ingestion |
| Impression-ID dedup + nightly reconciliation | Everything, eventually | Dedup state in stream; results final only next day | Dashboard can overcount intraday |
The honest architecture uses all three rows: idempotent upserts because they're nearly free, transactional checkpointing because Flink gives it to you, and impression-ID dedup plus reconciliation because the first two cannot see client-born duplicates. If the interviewer asks which one you'd drop under time pressure, drop the middle row; the outer two cover its failure modes a day later.
State management ties both threads together. Per-key window counts and the dedup set live in Flink's keyed state, backed by RocksDB and checkpointed to S3 every, say, 60 seconds. On failure, the job restores the last checkpoint and rewinds Kafka to the matching offsets. Checkpoint interval is a real dial: shorter means less reprocessing after a crash but more sustained I/O and more frequent alignment pauses under backpressure. Sizing the state aloud is worth doing, and the trap is sizing the wrong state. Window counters for 10 million active ads come to single-digit gigabytes. The impression-ID dedup set dominates: at 10 billion clicks a day and a one-to-two-hour ID validity window, that is roughly 400 to 800 million live keys, tens of gigabytes. Still comfortable for RocksDB on local SSD, but budget for the biggest state in the job, not the smallest.
The common mistake in this section is hand-waving "Flink handles late data" or "Kafka is exactly-once." Name the allowed-lateness setting, name the side output, name who consumes it, and name which duplicates each mechanism cannot see. That third destination for a late click existing at all is the insight.
Enrichment: join the stream without letting the join stall it
Asked at Rippling — Design an ad-click aggregation and enrichment pipeline A phone-screen variant of the classic: the click stream must be joined against other data sources to enrich each event before aggregation. The interviewer also pushed on how events should be sent from mobile apps versus browsers, and on the trade-off between the number of requests and delivery latency.
Raw click events are deliberately skinny; a 100-byte event doesn't carry campaign budgets or advertiser account tiers. Somebody downstream wants clicks grouped by campaign, geo, or account segment, which means joining the stream against reference data. The design question is where that join runs, and the rule that should govern your answer: never let a slow lookup service apply backpressure to click ingestion.
Three patterns, by size of the reference data:
- Broadcast state for small dimensions (campaign metadata, a few GB): replicate the whole table into every Flink task's memory, kept fresh by consuming a change-data-capture topic. Joins are local memory lookups, effectively free. This is the default and covers most of what an ad pipeline needs.
- Async lookup with a cache for large dimensions (per-user profiles): fire a non-blocking request to the profile store with a strict timeout, cache hot keys, and, critically, on timeout emit the event un-enriched with a flag rather than holding up the stream. Enrichment is best-effort; counting is not.
- Defer to batch for anything only analysts need: land the skinny event, join it in the warehouse at query time. The TikTok-style conversion analysis in the storage section below is exactly this shape.
The mobile question Rippling appended is worth engaging on its own terms. An SDK that fires one HTTP request per click drains batteries and hammers your edge with tiny payloads. Batching (flush every 30 events or 5 seconds, whichever comes first, plus a forced flush when the app backgrounds) cuts request volume by an order of magnitude at the cost of bounded staleness. For an aggregator with one-minute windows and 30 seconds of watermark slack, five seconds of client-side delay is invisible. The trade-off has a right answer given the latency budget you established in section one, which is why establishing it early pays off.
Storage and the read path: pre-aggregate for dashboards, keep raw for questions you haven't thought of yet
Asked at TikTok — Calculate Conversion Rate from Ad Clicks to Page Visits A SQL round built on exactly this system's output: given a user event log containing
ad_clickandpage_visitrows with timestamps, compute the rate at which clicks convert to subsequent visits. It's the consumer's view of the aggregator, and it needs event-level data that no pre-aggregated rollup can answer.
Two stores, two jobs.
The OLAP store (ClickHouse, Druid, Pinot; name one and move on) holds the stream's output. In ClickHouse terms:
CREATE TABLE ad_clicks_minutely (
ad_id UInt64,
window_start DateTime, -- minute boundary
clicks UInt64,
uniques_state AggregateFunction(uniqHLL12, UInt64) -- HLL sketch, not a count
)
ENGINE = ReplacingMergeTree
ORDER BY (ad_id, window_start);
A dashboard query is a range sum:
SELECT toStartOfHour(window_start) AS hr, sum(clicks)
FROM ad_clicks_minutely
WHERE ad_id = 4217 AND window_start >= now() - INTERVAL 7 DAY
GROUP BY hr ORDER BY hr;
Three details in that schema deserve narration. The engine is where the "upsert" from the windowing section actually lives, and it is an eventual one: ReplacingMergeTree keeps the last row per (ad_id, window_start) only after background merges run, so a re-emitted window briefly coexists with its older version unless the query pays for FINAL. A dashboard tolerates that wobble; it is one more reason billing never reads this table. Unique users is stored as an HLL sketch state rather than a number, read back with the uniqHLL12Merge combinator, because sketches merge across windows and exact distinct counts don't; you cannot sum minute-level unique counts into an hourly unique count. And the table is tiered for retention: minute rows for 7 days, hourly for 90 days, daily forever. The first rollup cuts row count 60x and the second a further 24x, and nobody examines minute-level data from last quarter.
If the read pattern is narrower, "current spend for this ad, right now," a wide-column KV store with atomic counters is a leaner fit than an OLAP engine; the modeling trade-offs there are the subject of our DynamoDB interview guide. And whether the dashboard polls the store or the pipeline pushes deltas to it over a socket is the classic freshness-versus-fan-out decision covered in Push vs Pull Architecture.
The raw log in S3 (Parquet, partitioned by event hour) is the other store, and it is not a backup. It answers every question the rollup schema didn't anticipate: the TikTok conversion-rate analysis above, a fraud investigation over one IP range, a rebuild of the OLAP table after a bug in the aggregation job. That last one is the argument to make in the interview: aggregates are a cache of the log. Caches get invalidated; logs don't.
The common mistake is proposing one store to serve both jobs, usually "we'll just query the raw events with Presto for the dashboard too." At 1 TB/day, an interactive per-ad range scan over raw Parquet is a 30-second query on a good day. Pre-aggregation is not an optimization here; it's what makes the read path exist.
Reconciliation and fraud: why billing never reads the stream
Asked at Meta — Identify Probability Distributions for Modeling Ad Clicks A data science round on the same domain: name probability distributions suited to modeling click behavior, state each one's core assumption and moments, and reason about a population that mixes high-intent and low-intent users. It is not a design question, but it drills the kind of distributional reasoning that fraud models are built on.
Everything so far produces numbers that are fast and almost right. The batch layer produces numbers that are slow and actually right, and money only ever flows against the second kind.
A nightly job reads the raw S3 log for the closed day and recomputes every aggregate from scratch: exact dedup by impression_id (trivial in a batch GROUP BY, no streaming state gymnastics), inclusion of everything from the late-clicks topic, and application of the day's final fraud verdicts. It writes a billing-grade table, then diffs its totals against what the stream emitted. Drift below a threshold, say 0.1%, is logged; drift above it pages someone, because it means the stream lost or double-counted events in a way the design didn't predict.
Fraud is the reason this layer can't be optional. Click farms and bot traffic are detected by models that compare observed click behavior against expected distributions, per user segment, per placement, per time of day; the Meta data science question above drills the same kind of distributional reasoning. The operational consequence for your design is that a click's fraud verdict can change hours after the click: a device flagged at 6 PM invalidates its clicks from that morning. A stream cannot retract an emission it made eight hours ago, but a batch recount simply excludes those clicks. So invoices wait for the recount, and the dashboard carries a footnote that today's numbers are preliminary.
This resolves the exactly-once anxiety from earlier, and it's worth closing the loop explicitly in the interview: the stream doesn't have to be perfect, it has to be close enough that advertisers trust the dashboard, and detectably wrong when it isn't. Perfection is the batch layer's job, and it gets a full day and an immutable log to deliver it.
Practice these on PracHub
Work these in roughly this order; each drills one layer of the design above.
- Design an ad click aggregation service (Meta) — the canonical end-to-end prompt; practice the dashboard/billing framing until it's reflexive.
- Design ad clickstream analytics pipeline (Amazon) — depth on ingestion: partition keys, schema registry, delivery semantics, raw/curated zones.
- Design Real-Time Analytics Pipeline with Kafka and Flink (Robinhood) — windowing, watermarks, state, and backpressure with the stack fixed for you.
- Design an ad-click aggregation and enrichment pipeline (Rippling) — stream joins and the mobile batching trade-off.
- Calculate Conversion Rate from Ad Clicks to Page Visits (TikTok) — the analyst's side of the pipeline, straight SQL over the event log.
- Identify Probability Distributions for Modeling Ad Clicks (Meta) — the statistics under the fraud layer, if you're interviewing for DS-adjacent roles.
More system design prompts with company tags are in the interview questions bank, SQL and coding drills live under coding questions, and the rest of our deep-dives are on the resources page.
Comments (0)