Skip to content
BytePatterns

Design an Ad Click Aggregator: Dedup, Windows and Watermarks

9 min readBytePatterns

Design an ad click aggregator: an append-only log, dedup by event id, per-campaign minute windows, and a watermark that trades lost late clicks against delay.

Count clicks per campaign per minute. Stated like that, it is a GROUP BY. The interview version adds what makes it hard: a hundred thousand events a second, the same click delivered twice, clicks arriving late and out of order, and a number somebody is billed against. This article designs the counting path and measures the one trade-off that cannot be engineered away.

The problem it solves

  • Count exactly once. A click counted twice is an overcharge.
  • Count by event time. A click belongs to the minute it happened, not the minute it arrived.
  • Publish promptly. Dashboards and budget caps want the minute's number soon after the minute ends.
  • Answer new questions later. Finance will ask for counts by country next quarter.

With illustrative assumptions, 100,000 events a second is about 8.6 billion a day; at 100 bytes each, roughly 860 GB of raw events a day. The counts themselves are tiny by comparison: one number per campaign per minute.

The intuition

Four decisions carry the design:

  1. Append, don't count. Each click goes into a durable, partitioned log with a unique event id, assigned where the click is recorded. The log is the source of truth, and it is kept.
  2. Dedupe by id. Delivery is at least once, so redeliveries are normal. The consumer remembers the ids counted in each open minute and drops repeats; once a minute is sealed, its ids can be forgotten.
  3. Windows keyed by event time. Each click increments the bucket for its (campaign, minute). The bucket key decides what can be answered; anything else needs a replay.
  4. A watermark closes the minute. The consumer cannot know that no more 09:41 clicks are coming, so it declares it: here, the watermark is the newest event time seen minus an allowed lateness. Once it passes 09:42, the 09:41 bucket is sealed and published. A click for a sealed minute goes to a side output for a correction job.

The lateness is the trade-off. Small, and late clicks miss their minute; large, and every number waits that long. Billing often uses a batch recount over the raw log as the reconciled number, while the stream feeds dashboards (from memory, a common split).

To scale, partition the log by campaign; a very hot campaign can be split across partitions with a key suffix and summed afterwards.

Watch it run

The animation starts from the stakes: a hundred thousand clicks a second, and the final number is what somebody is billed. Nothing is counted on arrival; each click is appended to a log with its own id. The consumer reads forward and drops each click into its minute bucket: 09:41 for c7 reads 1. A second campaign lands in the same minute: one bucket per key, nothing else kept. Delivery is at least once, so e1 arrives again, and a second count would be a second charge. The id is already in the seen set, so the event is dropped instead of added. Then a gate reports late: its click belongs to 09:41 even though it is now 09:43, and c7 rises to 2. A watermark is what decides the minute is over; until it passes, the bucket stays open. Once it passes, the totals are sealed and published as that minute's answer. Close early and stragglers are lost; close late and every dashboard trails reality. And the bucket only answers what its key held: a new question replays the log.

Design Ad Click Aggregation

Step 1 of 11

A hundred thousand clicks a second, and the number at the end of it is what somebody is billed.

The same interactive animation as the lesson — step through it with the controls.

The code

A toy model, not a stream processor: a simulated firehose of 3,000 clicks over twenty minutes, 10% redelivered, 5% arriving 30 to 150 seconds late, consumed in arrival order by a counter that dedupes per open minute and seals minutes behind a watermark:

import random
from collections import Counter, defaultdict

def click_stream(n, seed):
    """Toy firehose: (arrival_s, event_id, campaign, country, event_s), at least once, out of order."""
    rng = random.Random(seed)
    out = []
    for i in range(n):
        t = rng.uniform(0, 1200)                                   # twenty minutes of clicks
        event = (f"e{i}", rng.choice(["c7", "c9", "c12"]), rng.choice(["DE", "TR", "US"]), t)
        delay = rng.uniform(0, 5) if rng.random() < 0.95 else rng.uniform(30, 150)
        out.append((t + delay, *event))
        if rng.random() < 0.1:                                     # redelivered after a retry
            out.append((t + delay + rng.uniform(0, 20), *event))
    return sorted(out)

class MinuteCounter:
    """Toy model: dedupe by id, count per (campaign, minute), seal when the watermark passes."""
    def __init__(self, lateness):
        self.lateness, self.watermark = lateness, float("-inf")
        self.open = defaultdict(Counter)            # minute -> campaign -> count
        self.seen = defaultdict(set)                # minute -> ids counted in it
        self.published, self.late = {}, set()

    def receive(self, event_id, campaign, event_s):
        minute, counted = int(event_s // 60), False
        if minute in self.published:                # its minute is sealed: side output
            self.late.add(event_id)
        elif event_id not in self.seen[minute]:     # a redelivery stops here
            self.seen[minute].add(event_id)
            self.open[minute][campaign] += 1
            counted = True
        self.watermark = max(self.watermark, event_s - self.lateness)
        for m in sorted(m for m in self.open if (m + 1) * 60 <= self.watermark):
            self.published[m] = self.open.pop(m)    # the answer for that minute
            del self.seen[m]                        # ids of a sealed minute can be forgotten
        return counted

    def flush(self):
        for m in sorted(self.open):
            self.published[m] = self.open.pop(m)

stream = click_stream(3000, seed=38)
clicks = {(e, c, t) for _, e, c, _, t in stream}
print(len(stream), "deliveries,", len(clicks), "distinct clicks")
for lateness in (0, 10, 60, 180):
    agg, waits = MinuteCounter(lateness), []
    for arrival, e, c, _, t in stream:
        before = len(agg.published)
        agg.receive(e, c, t)
        waits += [arrival - (m + 1) * 60 for m in list(agg.published)[before:]]
    agg.flush()
    counted = sum(sum(v.values()) for v in agg.published.values())
    print(f"lateness {lateness:>3}s: missed {len(clicks) - counted:>3},",
          f"typical wait after the minute {sorted(waits)[len(waits) // 2]:.0f}s")
# 3280 deliveries, 3000 distinct clicks
# lateness   0s: missed 176, typical wait after the minute 2s
# lateness  10s: missed 124, typical wait after the minute 11s
# lateness  60s: missed  62, typical wait after the minute 62s
# lateness 180s: missed   0, typical wait after the minute 182s

The whole trade-off in four lines: zero lateness publishes two seconds after each minute and misses 176 clicks, about 6%; 180 seconds of lateness misses none and makes every number that much older. The animation's sequence, with a minute of lateness: a duplicate dropped, a late click still counted, one after the seal sent to the side output, then a replay answering a question the buckets never stored:

def hhmm(s):
    return f"{int(s // 3600):02}:{int(s % 3600 // 60):02}:{int(s % 60):02}"

agg = MinuteCounter(lateness=60)
nine41 = 9 * 3600 + 41 * 60
for e, c, t in [("e1", "c7", nine41 + 5), ("e2", "c9", nine41 + 20), ("e1", "c7", nine41 + 5),
                ("e4", "c9", nine41 + 110), ("e3", "c7", nine41 + 30), ("e5", "c7", nine41 + 130),
                ("e6", "c7", nine41 + 40)]:
    print(e, hhmm(t), agg.receive(e, c, t), "watermark", hhmm(agg.watermark))
print(dict(agg.published[nine41 // 60]), agg.late)
# e1 09:41:05 True watermark 09:40:05
# e2 09:41:20 True watermark 09:40:20
# e1 09:41:05 False watermark 09:40:20
# e4 09:42:50 True watermark 09:41:50
# e3 09:41:30 True watermark 09:41:50
# e5 09:43:10 True watermark 09:42:10
# e6 09:41:40 False watermark 09:42:10
# {'c7': 2, 'c9': 1} {'e6'}

by_country = Counter((k, int(t // 60)) for _, k, t in {(e, k, t) for _, e, _, k, t in stream})
print(by_country[("TR", 5)], sum(by_country.values()))             # 47 3000

The seeded check runs 200 random streams. No click may be counted twice, every click must be either counted or in the side output, published buckets must equal the counted events, and with lateness longer than any delay the result must equal a brute-force distinct count:

ok = True
for seed in range(200):
    r = random.Random(seed)
    s = click_stream(r.randint(0, 300), seed)
    lateness = r.choice([0, 5, 30, 200])
    agg, counted = MinuteCounter(lateness), []
    for _, e, c, _, t in s:
        if agg.receive(e, c, t):
            counted.append((e, c, int(t // 60)))
    agg.flush()
    ids = [e for e, _, _ in counted]
    ok &= len(ids) == len(set(ids))                           # nothing counted twice
    ok &= set(ids) | agg.late == {e for _, e, _, _, _ in s}   # nothing vanishes without a trace
    got = Counter({(c, m): k for m, v in agg.published.items() for c, k in v.items()})
    ok &= got == Counter((c, m) for _, c, m in counted)
    truth = Counter((c, int(t // 60)) for _, c, t in {(e, c, t) for _, e, c, _, t in s})
    ok &= all(got[k] <= truth[k] for k in got)                # late clicks only ever subtract
    if lateness == 200:                                       # longer than any delay in the stream
        ok &= got == truth and not agg.late
print(ok)                                                     # True

The complexity

  • Per event: one log append, one set lookup and one counter increment, all O(1).
  • Memory: the ids of open minutes only, about event rate × (window + lateness).
  • Storage: raw events grow forever unless aged out; counts are campaigns × minutes.

Where it goes wrong

  • Counting by arrival time. Late clicks land in the wrong minute.
  • No event id. Retries become double charges, with no way to tell.
  • Remembering every id forever. Forget them when their minute seals.
  • Dropping the raw log. Then a new question has no answer.
  • One hot campaign on one partition. Split its key and sum.

When it shows up in interviews

As "design an ad click aggregator", "count events per minute at scale" or "real-time analytics", and inside billing and fraud designs. Follow-ups: duplicates, late data, watermarks, hot keys, and reconciling the stream with a batch recount. As of October 2026, mainstream stream processors offer event-time windows, watermarks and allowed lateness, and some offer exactly-once state updates through checkpoints (from memory); deduplicating by event id is still needed for producer retries. The neighbouring message queue design covers at-least-once delivery itself.

How to say it in an interview

"Every click is appended to a partitioned log with a unique event id, and the raw log is kept. A stream consumer partitioned by campaign dedupes by id and counts into (campaign, minute) buckets by event time. A watermark, newest event time minus an allowed lateness, decides when a minute is sealed and published; anything later goes to a side output and a correction job. The lateness is the trade-off between missed clicks and delay. Billing uses a batch recount from the log, which also answers any new grouping by replay."