Ad Click Aggregator (Real-Time Billing)
Sharpened prompt. Design a system ingesting 1M ad events/sec (impressions, clicks, conversions), producing advertiser-facing dashboards under 1 minute stale, enforcing daily budgets so an advertiser never overspends by more than a fraction of a percent, detecting click fraud before it is billed, correctly attributing conversions to clicks across a 7-day window, and reconciling to the cent because the output is an invoice.
Two constraints collide: billing demands exactness and the stream demands approximation. Where you draw the line between them — and how you reconcile the two — is the interview.
1. Problem framing
Functional requirements
- Ingest impressions, clicks, and conversions with campaign/ad/placement dimensions.
- Real-time aggregation by campaign, ad, geography, device, hour.
- Enforce daily and lifetime budgets; stop serving when exhausted.
- Attribute conversions back to the click (and impression) that caused them.
- Detect and exclude fraudulent traffic from billing.
Non-functional requirements
| Property | Target | Consequence |
|---|---|---|
| Billing accuracy | Exactly-once, reconcilable | Two-phase commit sinks plus a batch ground truth |
| Dashboard freshness | < 1 min | Streaming aggregation, not batch |
| Budget overspend | < 0.5% | Distributed budget with local reservations |
| Late data | Accept up to 24h | Event time, watermarks, and window revision |
| Query latency | P95 < 2s on arbitrary slices | Columnar OLAP with pre-aggregation |
Back-of-the-envelope
Events: 1M/sec = 86B/day. Impressions ≈ 95%, clicks ≈ 5%, conversions ≈ 0.05%
Raw size: 86B × 400 B = 34 TB/day raw -> ~4 TB/day compressed columnar
Dimensions: campaign × ad × placement × country × device × hour
10k × 100k × 1k × 200 × 5 × 24 = astronomically large in principle,
but the SPARSE realised set is ~500M rows/day. Pre-aggregate on that.
Budget: 500k active campaigns, each needing a spend check per auction
-> budget checks happen at ad-serving rate (millions/sec), not click rate
Attribution: 7-day click window -> ~30B clicks retained in a joinable state store
The attribution number is the one that shapes the design: you must keep 30 billion clicks queryable by user for a week, which is a state-management problem, not an aggregation problem.
2. High-level architecture
Note the two paths from Kafka: the streaming path for dashboards and budgets (fast, approximate-then-corrected) and the lake path for billing (slow, exact, reconciled). Making that split explicit early is the strongest structural move on this problem.
3. Component inventory
| Component | Concrete choice | Why this one |
|---|---|---|
| Collector | Lightweight edge service writing straight to Kafka | Must never be the bottleneck; no business logic at ingest |
| Bus | Kafka with transactional producers | Exactly-once requires transactional semantics end to end |
| Stream engine | Flink with RocksDB state + S3 checkpoints | Event time, large keyed state, and 2PC sinks — all three are needed |
| OLAP | ClickHouse (or Pinot for lower-latency ingest) | Sub-second GROUP BY over billions of rows with sparse dimensions |
| Budget store | Redis with sharded counters | Read at ad-serving rate, which is far above click rate |
| Raw lake | Parquet on S3, hour-partitioned | Ground truth for reconciliation and model training |
| Fraud | Flink CEP for patterns + an ONNX model for scoring | Rules catch the known, models catch the novel |
4. The toughest parts
4.1 Exactly-once, when "exactly-once" does not exist
Why it's hard. Advertisers are billed from these counts. Double-counting means overbilling — refunds, disputes, and in aggregate a lawsuit. Under-counting means lost revenue. But every layer is at-least-once by nature: the browser retries the click beacon, Kafka producers retry on timeout, and a Flink job restarting from a checkpoint reprocesses everything since that checkpoint.
Solution — exactly-once effect through three cooperating mechanisms.
// 1. Deduplicate on a client-generated event ID before anything else.
events.keyBy(e -> e.eventId)
.process(new KeyedProcessFunction<>() {
ValueState<Boolean> seen; // RocksDB-backed, TTL 24h
public void processElement(Event e, Context ctx, Collector<Event> out) {
if (seen.value() != null) { metrics.inc("dedup.dropped"); return; }
seen.update(true);
ctx.timerService().registerEventTimeTimer(e.ts + DAY); // TTL cleanup
out.collect(e);
}
})
// 2. Checkpointed state + transactional sink = atomic "state advanced AND output
// committed". On restart, uncommitted output is rolled back with the state.
env.enableCheckpointing(60_000, CheckpointingMode.EXACTLY_ONCE);
env.getCheckpointConfig().setCheckpointStorage("s3://ckpt/adagg");
// 3. Two-phase commit into the sink: pre-commit on checkpoint, commit on
// checkpoint-complete. Kafka and JDBC sinks implement TwoPhaseCommitSinkFunction.
aggregated.sinkTo(KafkaSink.<Row>builder()
.setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
.setTransactionalIdPrefix("adagg-")
.build());
Say the honest version out loud: "exactly-once delivery is impossible; exactly-once processing semantics is achievable, and it means at-least-once delivery combined with idempotent or transactional state updates." Flink's guarantee is that state and output advance atomically with the checkpoint, so a replay after failure does not double-apply.
Where a 2PC sink is unavailable, fall back to an idempotent upsert keyed by (window_start, campaign_id, ad_id, country, device). Replaying a window overwrites rather than adds. ClickHouse's ReplacingMergeTree with a version column does exactly this, and naming the specific table engine is a good concrete detail.
The real safety net is 4.7: reconcile the stream against a batch recomputation from the raw lake, and bill from the reconciled number.
4.2 Budget enforcement that cannot overspend
Why it's hard. An advertiser sets a $10,000 daily budget. Ad servers across many regions each decide, millions of times per second, whether to show the ad. Checking a central counter on every auction adds a network round trip to a sub-10ms serving path. Not checking means overspend — and overspend is money you cannot bill, so it comes directly out of your margin. Worse, spend is only known after clicks arrive, which lag the impression that caused them.
Solution — hierarchical budget reservation: a central authority leases slices to local enforcers.
# Each ad server leases a small slice of budget and spends it locally.
class BudgetLease:
def __init__(self, campaign_id, slice_cents=5_00):
self.remaining = 0
self.slice = slice_cents
async def try_spend(self, cents: int) -> bool:
if self.remaining >= cents:
self.remaining -= cents # local, zero-latency decision
return True
if not await self._refill():
return False # budget exhausted globally
return await self.try_spend(cents)
async def _refill(self) -> bool:
# Atomic lease from the central counter. Returns 0 if the budget is spent.
got = await redis.eval("""
local rem = tonumber(redis.call('GET', KEYS[1]) or '0')
local want = tonumber(ARGV[1])
local give = math.min(rem, want)
if give > 0 then redis.call('DECRBY', KEYS[1], give) end
return give
""", keys=[f"budget:{self.campaign_id}:{today()}"], args=[self.slice])
self.remaining += got
return got > 0
Three properties to defend. Slice size controls the overspend bound: worst case is slice × number_of_servers if every server holds an unreturned slice when the budget hits zero. With a $5 slice and 200 servers that is $1,000 on a $10,000 budget — too much, so make the slice adaptive: large when far from the budget, shrinking to cents as it nears exhaustion.
def slice_for(remaining_cents, daily_budget_cents):
frac = remaining_cents / daily_budget_cents
if frac > 0.20: return 5_00 # plenty left: big slices, few round trips
if frac > 0.05: return 50 # getting close: tighten
return 5 # nearly done: near-exact enforcement
Return unused slices when a server drains or shuts down, so budget is not stranded. And handle the click lag: an impression's cost may only be known when the click arrives seconds later, so reserve the expected cost at serve time (bid × predicted CTR) and true it up when the actual event lands. Mentioning that spend is probabilistic at decision time is a detail that shows you understand ad serving rather than just counters.
Finally, pacing: without it, a campaign spends its entire daily budget in the first hour on cheap morning traffic. Spread spend against a forecast of the day's traffic curve, which is a smoothing problem on top of the enforcement problem.
4.3 Click fraud in the stream
Why it's hard. Botnets click competitors' ads to drain their budgets, publishers generate fake clicks on their own inventory, and click farms produce traffic that looks human at the individual level. Detection must happen before billing (refunding is expensive and damages trust) yet cannot add latency to ingestion, and it must catch novel patterns that no rule anticipated.
Solution — layered detection: cheap deterministic rules in the stream, a model for the grey zone, and a batch pass for patterns only visible in aggregate.
// Layer 1: Flink CEP for velocity and impossible-behaviour patterns.
Pattern<Click, ?> rapidFire = Pattern.<Click>begin("first")
.followedBy("burst").times(20)
.within(Time.seconds(10)); // 20 clicks in 10s from one identity
// Layer 2: per-event features -> gradient-boosted model, ~1ms on CPU.
clicks.map(c -> {
Features f = new Features(
/* time from impression to click */ c.ts - c.impressionTs,
/* click coordinates within the ad */ c.x, c.y,
/* IP reputation, ASN type */ ipIntel.score(c.ip),
/* user agent entropy, JA4 hash */ fingerprint.score(c.ua, c.ja4),
/* prior activity for this identity */ history.summary(c.userId),
/* conversion rate of this source */ sourceStats.cvr(c.publisherId));
c.fraudScore = model.score(f);
return c;
})
.process(new Classify()); // < 0.3 bill, > 0.8 drop, in between: HOLD
// Layer 3 (batch, nightly): graph clustering over (ip, device, publisher) to find
// coordinated rings that no single event reveals.
The three-way outcome is what makes this workable: bill, drop, and hold. Held clicks are excluded from billing pending the nightly batch pass, which sees patterns — a publisher whose traffic converts at 1/50th the network average, a cluster of devices sharing an improbable fingerprint — that no per-event scorer can detect. Held clicks are then either released and billed, or credited.
Signals worth naming because they are the ones that actually work: time from impression to click (humans take at least a few hundred milliseconds; bots often click in under 50 ms), click coordinate distribution (real clicks cluster on the visual call-to-action; synthetic clicks are uniform or dead-centre), downstream conversion rate (fraud produces clicks and never conversions — this is the strongest single signal), and cross-advertiser correlation (the same identity clicking many unrelated advertisers).
Publish a fraud rate metric per publisher and use it in payout decisions. That closes the economic loop, which is ultimately how fraud is reduced — detection alone just moves it.
4.4 Late data and windows that must be revised
Why it's hard. A mobile app buffers events offline and uploads them hours later. A tracking pixel fires on a page that stays open overnight. Meanwhile the hour's billing window closed, the invoice line was computed, and the advertiser's dashboard shows a number they have already screenshotted. Dropping late events under-bills and under-reports; accepting them indefinitely means no number is ever final.
Solution — event time with watermarks, a bounded lateness window that revises, and a hard cut-off with a compensating ledger entry beyond it.
events
.assignTimestampsAndWatermarks(
WatermarkStrategy.<Event>forBoundedOutOfOrderness(Duration.ofMinutes(5))
.withIdleness(Duration.ofMinutes(1))
.withTimestampAssigner((e, ts) -> e.clientEventTimeMs))
.keyBy(Event::dimensionKey)
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
.allowedLateness(Time.hours(24)) // revise the window for a full day
.sideOutputLateData(veryLate) // beyond 24h: not a window update
.aggregate(new CountSpend())
.sinkTo(upsertSink); // keyed upsert: revision, not double-count
// Beyond the lateness window, do NOT rewrite history. Book a correction instead.
veryLate.addSink(new CorrectionLedgerSink());
Two ideas to articulate. Within the lateness window, revise; beyond it, compensate. Rewriting a closed billing period breaks audit and invalidates invoices already sent; issuing a dated correction (a credit or an additional charge in the current period) is what finance actually expects, and it mirrors the payment system's rule that you never mutate history.
And trust client timestamps only after validating them. A malicious or broken client can send an event dated last month, forcing a rewrite of a settled period. Clamp client timestamps to a plausible range around server receipt time, and route anything outside it to the fraud path rather than the billing path — a subtle attack vector worth mentioning unprompted.
4.5 Attribution: joining conversions to clicks across seven days
Why it's hard. A conversion today might be attributable to a click seven days ago. That means keeping every click joinable by user for a week — around 30 billion records. It is a stateful stream join with an enormous window, and the attribution rule itself is a business policy (last click? first click? multi-touch?) that must be recomputable when the policy changes.
Solution — a keyed state store holding a bounded click history per user, with attribution as a pure function over it.
// Key by user identity so clicks and conversions for one user land on one operator.
clicks.union(conversions)
.keyBy(Event::userId)
.process(new KeyedProcessFunction<String, Event, Attribution>() {
// Bounded per-user history: the most recent N clicks within the window.
ListState<Click> recentClicks; // RocksDB, TTL 7 days
public void processElement(Event e, Context ctx, Collector<Attribution> out) {
if (e instanceof Click c) {
recentClicks.add(c);
ctx.timerService().registerEventTimeTimer(c.ts + SEVEN_DAYS);
return;
}
Conversion v = (Conversion) e;
List<Click> eligible = filterWindow(recentClicks.get(), v.ts, SEVEN_DAYS);
// Attribution model as a pure function -> swappable, testable, replayable.
out.collect(ATTRIBUTION_MODEL.attribute(v, eligible));
}
});
// Last-touch: 100% credit to the most recent eligible click.
// Linear: equal split across all eligible clicks.
// Time-decay: weight ∝ exp(-(t_conv - t_click) / halfLife).
Three practical points. State size is the design constraint — 30B records in RocksDB means careful key design, TTL timers to bound growth, and incremental checkpoints (Flink's changelog-based checkpointing) so a 60-second checkpoint does not have to write terabytes.
Identity is the hard part, not the join. Cookies expire, users switch devices, and privacy changes (ATT, third-party cookie deprecation) break the join outright. Model identity resolution as a separate service producing a stable identity_id, and be explicit that a growing share of conversions will be unattributable — handled by modelled/aggregate attribution rather than pretending the join succeeded.
Make the attribution model swappable and re-runnable. Advertisers change models and ask "what would last quarter look like under time-decay?" Because the raw events are in the lake, you re-run the pure function over history rather than needing the stream to have anticipated it.
4.6 Serving arbitrary slices under two seconds
Why it's hard. An advertiser asks for spend by campaign × country × device × hour for the last 30 days, then pivots to placement, then filters to mobile. Each query touches billions of rows across sparse high-cardinality dimensions. Pre-aggregating every combination is a combinatorial explosion; computing from raw on demand is too slow.
Solution — a columnar store with a small number of well-chosen pre-aggregations, plus rollup tiers.
-- Base fact table: one row per minute per realised dimension combination.
CREATE TABLE ad_events_1m (
minute DateTime,
campaign_id UInt32,
ad_id UInt32,
country LowCardinality(String),
device LowCardinality(String),
placement_id UInt32,
impressions UInt64,
clicks UInt64,
spend_micros UInt64,
conversions UInt64
) ENGINE = SummingMergeTree()
ORDER BY (campaign_id, minute, country, device, ad_id, placement_id);
-- SummingMergeTree merges duplicate keys by summing -> idempotent upserts are free.
-- Rollups: hourly and daily materialised views, built by the same engine.
CREATE MATERIALIZED VIEW ad_events_1h TO ad_events_1h_table AS
SELECT toStartOfHour(minute) AS hour, campaign_id, country, device,
sum(impressions) AS impressions, sum(clicks) AS clicks,
sum(spend_micros) AS spend_micros, sum(conversions) AS conversions
FROM ad_events_1m GROUP BY hour, campaign_id, country, device;
The choices that matter: ORDER BY leading with campaign_id because every advertiser query filters on it, giving enormous data skipping; LowCardinality for country and device, turning string comparisons into dictionary lookups; SummingMergeTree so late revisions are just another insert that merges; and rollup tiers so a 30-day query reads hourly rows rather than 43,000 minute rows per dimension combination.
Then the same query-planning trick as metrics monitoring: pick the coarsest tier that still fills the chart's pixels. A 30-day chart at 1,200 pixels needs hourly at most.
4.7 The number on the invoice
Why it's hard. The streaming pipeline is fast but has a failure surface: a bad deploy, a checkpoint corruption, a fraud model misfire, a partition that lagged for six hours. Any of those can produce a number that is wrong in a way nobody notices until an advertiser disputes their bill. Streaming is optimised for freshness, and billing needs to be optimised for being right.
Solution — dual pipelines with reconciliation, and bill from the batch.
# Nightly: recompute the day from the immutable raw event lake.
def batch_truth(day):
raw = spark.read.parquet(f"s3://lake/events/dt={day}/")
return (raw
.dropDuplicates(["event_id"]) # exact dedup, not probabilistic
.filter(col("fraud_verdict") != "drop")
.groupBy("campaign_id", "country", "device", "hour")
.agg(sum("spend_micros").alias("spend"), count("*").alias("events")))
def reconcile(day):
truth = batch_truth(day)
stream = clickhouse.read_day(day)
diff = truth.join(stream, on=KEYS, how="full_outer") \
.withColumn("delta", col("truth_spend") - col("stream_spend"))
total_delta = diff.agg(sum(abs(col("delta")))).collect()[0][0]
if total_delta / truth_total > 0.001: # > 0.1% divergence
alert.page("adagg reconciliation divergence", detail=top_offenders(diff))
# Billing always uses the batch number. The stream is for dashboards.
billing.write_invoice_lines(truth, day)
clickhouse.correct(diff) # heal the dashboard too
The principle to state plainly: the stream is for decisions that must be fast (budget, dashboards); the batch is for numbers that must be right (invoices). They should agree within a fraction of a percent, and the reconciliation job's purpose is to prove it daily and to alarm the moment they diverge.
This is the lambda architecture, and it is worth naming while also naming its cost: two implementations of the same logic that can drift apart. Mitigate by sharing the aggregation logic as a library used by both paths, so a change to fraud rules or dimension handling applies to both — the drift risk is real and the mitigation is what separates the answer from the textbook.
5. What breaks first
| Event | First failure | Mitigation |
|---|---|---|
| Traffic spike (a major event) | Kafka partition lag on hot campaigns | Partition by a composite key; scale consumers on lag |
| Flink checkpoint slow | Backpressure to ingest | Incremental checkpoints, RocksDB tuning, larger state backend nodes |
| Budget Redis unavailable | Serve without limits → overspend | Local leases keep serving briefly, then fail closed on budget (money at risk) |
| Fraud model false positives | Legitimate traffic unbilled | Three-way verdict with a hold bucket; nightly batch releases |
| Mass late data after an outage | Windows revised, dashboards shift | 24h allowed lateness; corrections beyond it, never history rewrites |
| Stream/batch divergence | Wrong invoices | Daily reconciliation gates billing; alarm above 0.1% |
| Advertiser queries 90 days raw | OLAP cluster saturation | Rollup tiers + query limits + resolution selection |
Note the budget row's asymmetry against the rate limiter: rate limiting fails open because it protects capacity, budget enforcement fails closed because it protects money. Making that contrast explicitly is a strong way to show the principle rather than the rule.
6. Cheat sheet
- Exactly-once effect: event-ID dedup in keyed state + Flink checkpoints + 2PC/idempotent-upsert sink. Delivery is at-least-once; semantics are exactly-once.
- Budget: hierarchical leases with an adaptive slice size (large far from the cap, cents near it); reserve expected cost at serve time; fail closed.
- Fraud: rules + model + a hold verdict resolved by a nightly graph pass. Strongest signal is clicks that never convert.
- Late data: event time, 24h allowed lateness with revision, corrections beyond it, and clamp untrusted client timestamps.
- Attribution: keyed 7-day click state, attribution as a pure swappable function, identity resolution as a separate concern with an unattributable bucket.
- Serving: ClickHouse
SummingMergeTreeordered by campaign,LowCardinalitydimensions, rollup tiers chosen by pixel count. - Billing: dual pipeline — stream for speed, batch from the raw lake for truth, reconcile daily, invoice from the batch.
- The one-liner: "Two pipelines over one event log: a streaming one that must be fast enough to stop spending, and a batch one that must be right enough to bill — and a daily reconciliation that proves they agree."