Ad Click Aggregator
Worked solution

Ad Click Aggregator — a worked solution

Stream processing with watermarks and exactly-once.

Try it yourself first.

You will remember almost none of this if you read it cold. The workspace walks the same 10 stages and runs the architecture you draw through a simulator, so you find out where your version breaks before you see ours.

Open the Ad Click Aggregator workspace

The problem

Build the system that powers "clicks per ad per minute, per hour, per day" reports for advertisers. Inputs: billions of click events flowing in from edge servers, late-arriving by seconds to days. Output: dashboards advertisers refresh constantly, plus billing rollups they will sue you over if they're wrong. Every choice — windowing, watermarks, idempotency, lambda vs kappa — is about the trilemma: low latency, exact accuracy, late-event tolerance.

The reference architecture

Reference architecture for Ad Click Aggregator: 9 components — Edge / SDK, Ingest API, Kafka (clicks), Flink: dedup + window, ClickHouse, S3 / Iceberg, Nightly Spark, Billing DB, Reports API — connected by 9 flows.Edge / SDKclientIngest APIserviceKafka (clicks)KafkaFlink: dedup + windowserviceClickHouseClickHouseS3 / IcebergobjectNightly SparkserviceBilling DBPostgresReports APIservice
9 components, 9 flows. A dashed line is an asynchronous hop. This is the reference design, not the only one that works.

Stage by stage

The same 10 stages the workspace walks, answered.

01Clarifications

What would you ask before drawing a single box?

Worth asking:

  • What counts as a click? Server-side recorded? Browser-fired beacon? Defines dedup semantics.
  • Latency target on dashboards. 1 minute? 5? Defines the streaming window cadence.
  • Late events. How late do we accept? Mobile devices in tunnels can re-emit clicks hours later. Caps storage and retraction logic.
  • Exact billing. Yes — eventually. Approximate live, exact in nightly batch reconciliation.
  • Time zones. Reporting in advertiser-local TZ, but storage in UTC.

Assumptions:

  • 10B clicks/day → ~120K/s avg, 500K/s peak.
  • ~1% events arrive > 1 hour late; 0.1% arrive > 24 h late.
  • Live dashboard: 1-min freshness, ±1% accuracy.
  • Billing: exact, computed nightly with 7-day retention for late corrections.

02Functional reqs

What must this system actually do?

  • Ingest click events from edge servers.
  • Aggregate clicks per (ad_id, time bucket) at multiple resolutions: 1-min, 1-hour, 1-day.
  • Serve real-time dashboard queries (ad_id, time range → count).
  • Produce daily billing rollups that are exact (after late-event tolerance).
  • Deduplicate: a click delivered twice by retry must count once.

03Non-functional

What must it promise about speed, uptime and correctness?

  • Live latency: Click → dashboard < 1 minute P99.
  • Accuracy: ±1% live; exact within 7 days for billing.
  • Durability: Click events must survive consumer crashes, broker restarts. Replicated.
  • Exactly-once aggregate: No double-counting on retry, on crash, on partial failure.
  • Scalability: Horizontal at every layer; no single-writer bottleneck.

04Capacity estimation

How much load and data does this have to hold?

  • Click QPS: 10B/day = 115K avg, 500K peak.
  • Per-event size: ~200 B (ad_id, user_id, ts, geo, device). Daily volume = 2 TB/day raw.
  • 7-day retention raw: 14 TB on Kafka + cold storage.
  • Aggregates (1-min bucket × 100M ads × 7 days): worst case 10B rows of aggregate state. Compresses heavily — most ads don't get clicks every minute. Realistic: ~100M active aggregate rows.
  • Network: 500K events/s × 200 B = 100 MB/s peak ingress.

05API design

What does the outside world call, and what comes back?

POST /ingest/clicks
Authorization: Bearer <edge-token>

[
  { "adId": "...", "userId": "...", "ts": "...", "ctx": { "geo": "...", "device": "..." }, "clientId": "<idempotency-key>" },
  ...
]

202 { "accepted": 500, "batchId": "p3-off-9241188" }

The ingest endpoint only appends to Kafka (it can't dedup synchronously — dedup happens later in Flink, against a stateful key set over a 1-hour window). So the contract is honest about the async append: 202 Accepted with the count written to Kafka and a partition-batch id. Deduplication is an eventual property surfaced through the aggregates and reports, not something the write path can report. (A synchronous dedup count would require an idempotency store — e.g. a Redis set keyed on clientId — at the ingest tier; that's a different design from the Flink-window dedup here.)

GET /reports?adId=ad_123&from=2026-04-01T00:00&to=2026-04-29T00:00&granularity=hour
200 { "buckets": [{ "ts": "...", "clicks": 18234 }, ...], "watermark": "2026-04-29T05:30:00Z" }

The watermark field tells the caller: "all events up to time T are accounted for." Anything after may still be revised.

06Data model

What gets stored, and what is it looked up by?

Raw event log — Kafka, partitioned by ad_id (so all clicks for one ad land on one partition for ordered processing):

fieldtype
ad_idstring
user_idstring
event_tstimestamp
ingest_tstimestamp

| client_id | string | // idempotency key

| ctx | json |

Aggregate store — keyed by (ad_id, bucket_ts, granularity). Append-friendly, time-ordered:

| ad_id | bucket_ts | granularity | clicks | last_updated |

Database choice — recommended:

  • Ingestion / source of truth: Kafka (or Kinesis, Pulsar). 7-day retention. The replayable log is the kappa architecture foundation — you can rebuild any aggregate by replaying.
  • Streaming compute: Flink (preferred for exactly-once + watermarks) or Kafka Streams.
  • Aggregate store: ClickHouse or Druid. Columnar OLAP, time-partitioned, compresses click time-series 50–100×, query "sum clicks for ad over range" in milliseconds.
  • Billing batch: Spark or Trino over the same Kafka-archived events in S3/Iceberg/Delta. Recomputes exact totals nightly. The streaming numbers are advisory — billing is reconciled offline.

Why not Postgres for aggregates? It can't sustain 100K+ writes/s on time-series data without serious sharding. ClickHouse was built for this.

07High-level design

Which components handle a request, and in what order?

Architecture (kappa with batch reconciliation):

  1. Edge servers record clicks, push to Kafka with a client-generated idempotency key.
  2. Kafka is the source of truth, partitioned by ad_id, replicated across brokers.
  3. Dedup stage (Flink) uses a state-backed key set: drop duplicate client_ids within a 1-hour window.
  4. Windowed aggregation (Flink) with event-time watermarks:
  • Tumbling 1-min windows for live dashboard.
  • Sliding 1-hour, 1-day windows materialized incrementally.
  • Allowed lateness: 1 hour for live; events later than that go to a side-output for batch reconciliation.
  1. Sink to ClickHouse every minute. Flink→ClickHouse is at-least-once — ClickHouse can't participate in Flink's 2-phase-commit sink (no XA / transactional coordinator a TwoPhaseCommitSink can drive). Make it effectively-once with idempotent writes: emit each 1-min aggregate under a deterministic key (ad_id, bucket_ts, granularity) carrying a version = window/checkpoint id, and let ReplacingMergeTree (latest-version-wins) collapse a replayed window so re-emitting overwrites rather than double-counts. (If you want true 2PC, sink into a Kafka aggregates topic — Kafka transactions are 2PC-capable — and have ClickHouse consume read-committed.)
  2. Cold storage: Kafka also tees to S3 / Iceberg via a sink connector. 7-year retention for billing audit.
  3. Batch reconciliation: Nightly Spark job over S3 recomputes the previous day's totals exactly. If they disagree with streaming numbers, replace with batch.
  4. Query API reads from ClickHouse for live, from the reconciled tables for billing.
  5. Watermark service publishes "events up to T are committed" so dashboards can show "as of …".

Lambda vs kappa: Pure lambda (separate streaming and batch paths) duplicates code and is operationally hellish. Kappa (one streaming pipeline + batch reconciliation as a re-derivation) is the modern answer. We're using kappa.

08Deep dives

Which part breaks first, and what do you do about it?

1. Exactly-once is a lie unless you carefully define it.

  • At the source (edge → Kafka): at-least-once. Edge retries until ACK; idempotency key dedups downstream.
  • Within Flink: exactly-once state via aligned checkpoints; the external sink can't be 2PC against ClickHouse, so it's made effectively-once with idempotent writes (below).
  • Sink to ClickHouse: at-least-once writes made idempotent. Key each aggregate row deterministically by (ad_id, bucket_ts, granularity) with a version = the checkpoint/window id that is identical whenever the same window is re-emitted after a restart, and use ReplacingMergeTree (latest-version-wins) so a re-processed window overwrites its prior row. Do not key on a per-run id that changes on restart — the replayed rows would then carry a fresh key, never collide with the pre-crash rows, and both copies would survive (the opposite of dedup).
  • End-to-end: "Effectively once" — exactly-once aggregate output, even if individual events were delivered more than once.

2. Watermark calibration. Watermark = oldest unprocessed event time. If you set it at "now - 1 min" and accept up to 1 hour late, late events that arrive create retractions and side outputs. Trade-off:

  • Tight watermark (5 min): freshest dashboards, more late-event noise.
  • Loose watermark (1 hour): cleaner data, dashboards lag 1 hour.
  • Solution: dual watermarks — fast for dashboards (advisory), slow for billing (authoritative).

3. Hot ad — one ad gets 50K clicks/s. Kafka partition by ad_id puts all of those in one partition → one consumer → bottleneck. Solutions:

  • Two-level keying: partition by hash(ad_id, sub_bucket) for the first stage, re-shuffle by ad_id at aggregation.
  • Pre-aggregation at the edge: edge servers locally aggregate 1-second buckets and ship aggregates instead of raw events.

4. Click fraud / bot filtering. Don't try to do this in the hot streaming path — it's expensive and accuracy isn't critical for live dashboards. Filter at the daily reconciliation step, when you have the full event history and can run heavier ML.

5. Schema evolution. Click event schemas change (new context fields, deprecated fields). Use Avro / Protobuf with a schema registry. Forward + backward compatibility is mandatory because raw events sit in Kafka/S3 for years.

09Trade-offs

What did this design cost, and what breaks at 10×?

What breaks at 10× scale (100B clicks/day):

  • Kafka partition count. Need 1000+ partitions for the click topic. Becomes operationally expensive.
  • Flink state size. Per-key state for dedup (1 hour window × 100B keys/day) gets enormous. Use RocksDB-backed state with TTL.
  • ClickHouse ingest rate. 500K rows/s is fine for one ClickHouse cluster; 5M/s requires sharding or moving to Druid.
  • Cost. Storing 7 days of raw events at 10× volume is ~140 TB on hot storage. Tier aggressively.

Failure stories:

  • Flink job crash mid-window: Restart from last checkpoint. Window state restored, and the window is re-emitted under the same deterministic key — the ReplacingMergeTree sink collapses the replay so the re-processed window overwrites rather than double-counts.
  • ClickHouse disk full: Aggregates don't write. Kafka holds events; on recovery, replay from offset. Live dashboards lag during outage.
  • Edge → Kafka outage: Edge servers buffer locally (10-min disk buffer). On Kafka recovery, drain buffers. Some clicks lost if the edge dies during outage — accepted as < 0.001% of data.
  • Late event burst (mobile network restored at scale): Watermark stalls; allowed-lateness fires → updates buckets up to an hour old. Dashboard shows retraction; billing absorbs the change.

Primary sources

  • Kafka paper
  • Streaming 101 & 102
  • Flink exactly-once

Now defend it

Reading a design is not the same as being able to hold one under questioning. The workspace asks the same questions an interviewer would, and the simulator disagrees with you when the diagram does not support the claim.

Work Ad Click Aggregator yourself