System Design Interview

Course Content

System Design Interview

31 sections · 71 lessons

Ad Click Aggregation deep dive: windowing, late events and exactly-once counting


The pipeline from the first lesson has a stream processor in the middle that "counts clicks per minute". That phrase hides the two hardest problems in this design. This lesson takes them one at a time: deciding which minute a click belongs to, and making sure a crash never counts it twice.

"Count clicks per minute" contains a hidden question: which minute? The answer is harder than it looks, and most candidates have never had to think about it.

Why processing time is the wrong choice, mostly

Counting by processing time is trivial: whatever arrives in this wall-clock minute goes in this bucket. It is also wrong for billing, because a network hiccup moves clicks between minutes and a backlog after an outage dumps an hour of clicks into a single bucket, producing a spike that never happened.

Counting by event time is correct and requires a mechanism for deciding when to stop waiting. That mechanism is the watermark.

How a watermark works

The aggregator tracks the maximum event time it has seen and subtracts an allowed lateness — say two minutes. That is the watermark. When the watermark passes 12:01:00, the window 12:00:00–12:00:59 is declared closed, its count is emitted, and its state is dropped from memory.

The allowed lateness is a direct trade. Two minutes means results are two minutes behind but almost all events are included. Ten seconds means fast results and more stragglers excluded. Our requirement of "under one minute" points at an allowed lateness of about 30 seconds for the dashboard path.

The three options for a late event

An event for 12:00 arrives at 12:04, after the window closed. Three choices, and the answer is not one of them — it is a combination.

Drop it. Cheap and clean. Acceptable for a dashboard, unacceptable for billing, because you are discarding revenue.

Emit an update. Reopen the window, re-emit a corrected count, and require every downstream consumer to handle a restatement of a number it already read. This is correct but pushes real complexity into billing systems that may not tolerate it.

Route it to a side output. Write late events to a separate "late" topic or table, keyed by their true window. Nothing downstream is disturbed, and the nightly reconciliation in Reconciliation and recovery folds them in.

Recommendation: drop from the fast dashboard path with a metric counting how many were dropped, and route to a side output so the slow path can include them. The dashboard is approximate by design and labelled as such; the invoice is computed from the reconciled figures. Being explicit that the dashboard and the invoice are allowed to differ is the mark of a candidate who has thought about this rather than recited it.

An event created at 12:00 that arrives at 12:0712:0012:0212:0412:0612:0812:010event time 12:00processing time 12:077 minutes late12:00 windowwatermark"nothing older than 12:00 will arrive after this"this event arrives after the watermark —it is late data, and the policy for it is adesign decisionThree options for late data: drop it, hold the window open longer, or emit a correction. Say which one you are choosing and why.
The watermark is a bet about lateness, not a fact — which is why every streaming system needs a late-data policy.

Exactly-once aggregation

Windows decide where a click is counted. The second problem is how many times. At-least-once delivery is the default everywhere in this course. Here it is a revenue defect.

Three mechanisms, one guaranteeEventcarriesits own IDDedup storedrops repeatsWindowstate updatedOffset andresult commitDownstreamwrite is keyedThe offset and the aggregate must land atomically, or a crash double-counts.
No single mechanism is sufficient; exactly-once here is the product of dedup, atomic checkpoints and idempotent writes.

What at-least-once costs, in money

The aggregator has consumed 40,000 events for the 12:00 window and crashes before committing its offset. On restart it resumes from the last committed offset, replays those 40,000 events, and adds them to a count that already includes them. The window reports 80,000 clicks.

At an average cost-per-click of ₹5, a 40,000-click over-count is ₹200,000 billed for traffic that did not happen — from one crash, in one window, for one ad. Crashes are routine. This is not an edge case.

Mechanism 1: idempotent writes keyed by event ID

Make the effect of processing an event depend on the event's identity, not on how many times it is processed.

The client — or the ingestion service, if the client cannot be trusted — assigns a globally unique event ID. The aggregator maintains a set of event IDs already counted for the open window. On replay, an ID already in the set is skipped.

The cost is memory: 40,000 events per minute per key, at 16 bytes per ID, is manageable, but across a million ads it is not. Two practical reductions: keep the set only for the open window and its allowed-lateness tail, dropping it when the window closes; and use a Bloom filter, a probabilistic set that answers "definitely not present" or "probably present", to cut memory roughly 10× at the cost of occasionally skipping a legitimate event. For billing, prefer the exact set over the Bloom filter and pay the memory — the failure mode of a Bloom filter here is under-counting, which is a revenue loss you cannot detect.

Mechanism 2: checkpointing with an atomic offset commit

The aggregator's in-memory window state and its consumed offset must move together. If the state is saved but the offset is not, replay double-counts. If the offset is saved but the state is not, restart loses counts.

The fix is a checkpoint: periodically write the window state and the consumed offset in one atomic operation, either as a single transaction into a store that supports it, or into the log itself using a transactional producer. On restart, load the state and resume from the offset in the same checkpoint. Whatever was consumed after the checkpoint is replayed, and Mechanism 1 makes that replay harmless.

Checkpoint frequency is a tuning dial: every 10 seconds means at most 10 seconds of replay after a crash, and a pause of tens of milliseconds every 10 seconds to write state.

Mechanism 3: end-to-end deduplication at the write

The final write into the aggregated store should also be idempotent. Write the count for a window as INSERT ... ON CONFLICT (window_start, ad_id, country, device) DO UPDATE SET clicks = EXCLUDED.clicks — an overwrite of an absolute value, not an increment. Re-emitting the same closed window produces the same row. Incrementing would not.