Skip to content
Reliable Data Engineering
Practice problem hard streaminglambdaexactly-onceolapbilling
Practise with timer, notes and rubric

Design an Ad Click and Impression Aggregation System

Problem

An ads platform serves impressions and records clicks. Advertisers need near-real-time dashboards of impressions, clicks, CTR and spend per campaign. Finance bills advertisers daily based on valid clicks. Design the aggregation system.


Clarifying questions

QuestionAssumed answer
Volume?10 B impressions/day, 1 B clicks/day, 5× peak
Dashboard freshness?≤ 1 minute, approximate OK
Billing?Daily, exact, auditable, after click-fraud filtering
Late events?Up to 24 h (mobile SDK batching)
Query patterns?Per campaign / ad / day / hour / minute, by country, device; last 90 days interactive
Number of campaigns?~1 M active ads, 100k advertisers

1. Requirements

2. Estimates

Impressions 10B/day → 100k/s avg, 500k/s peak; clicks 1B/day → 10k/s avg, 50k/s peak
~300 B/event Avro → ~165 MB/s peak → Kafka 128+ partitions
Raw storage: 11B × 300 B ≈ 3.3 TB/day → ~0.5 TB/day Parquet → 2 years ≈ 365 TB
Aggregates: 1M ads × 1,440 minutes = 1.44B minute-rows/day worst case;
            realistically ~100M (sparse) → OLAP store with rollups to hour/day after 7 days

3. Architecture

Requirements diverge (fast/approximate vs slow/exact), so this is one of the rare places where a Lambda-style split is justified, but with shared code.

flowchart LR
    ADS[Ad servers] -->|impression| K1[[Kafka: impressions]]
    CLK[Click redirect service] -->|click, click_id| K2[[Kafka: clicks]]
    K1 --> F[Flink / Spark SS<br/>dedup + 1-min aggregates]
    K2 --> F
    F --> AGG[[Kafka: agg_minute]]
    AGG --> OLAP[(Pinot / Druid<br/>real-time table)]
    OLAP --> DASH[Advertiser dashboards]
    K1 --> BR[(Bronze Delta)]
    K2 --> BR
    BR --> IVT[Batch IVT / fraud filtering<br/>bots, click farms, rules + ML]
    IVT --> SIL[(silver.valid_clicks<br/>exact, deduped)]
    SIL --> BILL[(gold.billing_daily<br/>immutable per day)]
    BILL --> INV[Invoicing]
    SIL --> OFF[(OLAP offline table<br/>replaces real-time segments)]
    OFF --> OLAP

The OLAP store’s hybrid table (Pinot real-time + offline segments) lets the batch-corrected numbers replace the streaming numbers for completed days, so advertisers converge to billing-grade numbers.

4. Data model

-- click event
click_id STRING (UUID generated by redirect service), impression_id STRING, ad_id, campaign_id,
advertiser_id, user_id_hash, ip_hash, user_agent, country, device, click_ts, received_ts, cost_micros

-- real-time aggregate (grain: ad × minute × country × device)
minute_ts, ad_id, campaign_id, country, device, impressions BIGINT, clicks BIGINT, spend_micros BIGINT

-- billing (grain: advertiser × campaign × day), immutable once closed
bill_date, advertiser_id, campaign_id, valid_clicks, invalid_clicks, amount_micros, run_id, closed_at

Money in integer micros (never floats).

5. Deep dives

5.1 Exactly-once counting

sequenceDiagram
    participant C as Click service
    participant K as Kafka
    participant S as Stream job
    participant O as Sink (agg topic / Delta)
    C->>K: produce (idempotent producer, key=click_id)
    K->>S: read offsets 100-199
    S->>S: dedup by click_id (state, 24h TTL)
    S->>O: write aggregates in a transaction
    S->>K: commit offsets in the same transaction
    Note over S,O: crash before commit means abort and re-read, no double count

5.2 Hot campaigns (skew)

A Super Bowl campaign can be 5% of all traffic. Keyed aggregation by ad_id overloads one task.

5.3 Late data

5.4 Invalid traffic filtering

5.5 Reconciliation

Daily job compares: streaming totals vs batch totals per campaign (expect small diffs from late data/IVT), batch clicks vs click-service logs, billing totals vs invoices. Alert on diffs above thresholds.

6. Trade-offs

DecisionChoiceAlternative
Two pathsStreaming (approx) + batch (exact) with shared dedup logicPure Kappa: possible, but heavy IVT and 24h lateness push billing to batch anyway
ServingPinot/Druid hybrid tablesClickHouse (great, simpler ops for some); BigQuery/Delta SQL too slow/costly for thousands of advertiser QPS
Dedup stateKeyed state with TTLBloom filter (smaller, false positives drop legit clicks, bad for billing)
Money typeinteger microsdecimal; floats are never OK

7. Failure modes

FailureHandling
Stream job down 30 minDashboards stale (show “data delayed” banner); billing unaffected (batch from bronze)
Duplicate click storm from buggy SDKclick_id dedup; volume anomaly alert
Batch IVT job failsBilling close delayed; SLA alert; rerun idempotently (overwrite day partition)
OLAP node lossReplicated segments; rebuild from Kafka/deep store

8. What separates a senior answer

9. Follow-up questions

How would you show "unique users reached" per campaign over arbitrary date ranges?

Store HLL sketches per (campaign, day) in the OLAP store or a Delta table; union sketches over the selected range at query time. Exact uniques over arbitrary ranges would require scanning raw data. Approximate (±1–2%) is standard for reach.

An advertiser disputes yesterday's bill. How do you answer?

Billing rows carry run_id and link to the exact silver.valid_clicks snapshot (Delta version/time travel). Pull click-level evidence: valid vs invalid with reasons, timestamps, dedup decisions. Because billing tables are immutable and versioned, you can reproduce exactly what was billed.

Why not just count clicks directly in Postgres?

50k writes/s peak with hot rows (counter contention on popular ads), plus 1B rows/day of raw data, is well beyond a single OLTP database. You also lose replay and fan-out. Fine for a startup with 100 clicks/s, not here.


Self-assessment rubric