Technologies referenced in this case study: Kafka · Flink · Time-Series Databases · Redis · Cassandra · DynamoDB
Related case studies: Stream Processing · Idempotency · Leaderboards & Counting · Payment Processing · Rate Limiting
How to Use This Case Study#
Organized for interview use first, reference second. Stream-engine internals (checkpoint barriers, watermark propagation, state backends) live in Stream Processing and Flink; this case study is about counting money correctly while also counting fast.
| Mode | Time | What to Read |
|---|---|---|
| Quick Review | 15 min | Executive Summary → Interview Walkthrough → Fault Lines table → Drills 1–3 |
| Targeted Study | 1–2 hrs | Executive Summary → Walkthrough → Section 3 (Fault Lines) → Section 4 (Failures) → Deep Dives 1, 2 and 4 |
| Deep Dive | 3+ hrs | Everything, including the Principal Lens (Section 11) and the appendices |
What is Ad Click Aggregation? — Why interviewers pick this topic
An ad click aggregator receives every click on every ad, redirects the user to the advertiser's landing page, and turns the raw click stream into counts: clicks per ad per minute, spend per campaign per hour, click-through rate per creative. Those counts feed three consumers that want very different things — the billing system (charge advertisers per click), the advertiser dashboard (see campaign performance now), and the budget pacer / fraud system (stop spending when the budget is gone or the traffic is fake).
Before vs After — the double-billing incident:
Without a correctness design:
t=0: Flink job restarts after a TaskManager OOM
t=+40s: Job resumes from last checkpoint (taken 4 min ago)
t=+40s: Sink is append-only INSERT into the OLAP store
t=+5min: 4 minutes of clicks replayed → counted twice
t=+1day: Invoices go out. 3 large advertisers overbilled ~1.8%
t=+9day: Advertiser's third-party tracker flags the discrepancy
t=+30day: Credits issued, finance restates revenue, trust damaged
With a correctness design:
t=0: Same OOM, same restart
t=+40s: Replay from checkpoint — sink is an idempotent upsert
keyed by (ad_id, window_start, job_epoch)
t=+40s: Replayed windows overwrite themselves. Zero double count.
t=+1day: Nightly batch recount agrees with streaming within 0.02%
t=+1day: Invoices generated from the reconciled ledger, not the stream
Why interviewers reach for this question: it is the cleanest prompt in the library for testing whether a candidate understands that "real-time" and "correct" are two different products. Almost every candidate can draw Kafka → Flink → OLAP. The interview is decided by whether they can say which number is allowed to be wrong, by how much, for how long, and who signs off on it.
Mechanics Refresher: Aggregation Architectures
| Architecture | How It Works | Pros | Cons |
|---|---|---|---|
| Batch only (Spark/Hive hourly) | Land raw clicks in object storage; recount every hour/day | Simple, exact, cheap, replayable | 1–24h latency — useless for pacing and fraud |
| Streaming only, at-least-once | Kafka → stream job → increment counters | Seconds of latency, simple | Double counts on every restart; no audit trail |
| Streaming exactly-once (kappa) | Kafka → Flink with checkpoints → idempotent/transactional sink; replay from Kafka for corrections | One codebase, seconds of latency, replayable within retention | Correctness depends on every sink being idempotent; replay window bounded by Kafka retention (7–30 days) |
| Lambda | Streaming path for speed + batch path for truth; serving layer merges | Batch is the audited source of truth; streaming can be approximate | Two codebases computing "the same" number, which drift |
| Kappa + batch reconciliation | Streaming is primary; nightly batch recount over the immutable raw log checks streaming and issues corrections | One business-logic definition, an independent audit, a clean correction story | Requires a correction protocol and a shared logic library |
For most production systems: kappa with batch reconciliation. Streaming produces the numbers people look at; an independent batch recount over the raw log produces the numbers finance signs.
Executive Summary
If you only read one section, read this. The whole case study is one idea: the dashboard number and the invoice number are different products with different correctness bars, and the design must make that explicit.
What This Interview Actually Tests#
Ad click aggregation is not a streaming question. Everyone can draw Kafka → Flink → Druid.
This is a financial correctness under time pressure question that tests:
- Whether you separate billing-grade counts from real-time counts before drawing anything
- Whether you can define "exactly-once" end to end — source, processing, sink — and say where it breaks
- Whether you know that late, duplicate and fraudulent clicks are the normal case, not the edge case
- Whether you have a correction story: when the number was wrong yesterday, how does it get fixed, and who is told?
The key insight: every click is money, but not every click count needs to be money-grade. Staff engineers build two products on one log — a fast approximate one and a slow audited one — and own the reconciliation between them.
The L5 → L6 → L7 Contrast — Start Here#
| Behavior | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| First move | Draws click service → Kafka → Flink → OLAP store | Asks "Is this number for invoices, dashboards, or pacing/fraud?" and commits to billing-grade as the anchor | Asks who signs the invoice number today (finance? ads-quality?) and whether a reconciliation process with advertisers already exists, because that contract constrains the design more than any technology |
| Exactly-once | "Flink has exactly-once" | Defines it end to end: dedup by click_id, checkpointed state, idempotent upsert sink keyed by window — and names the gap (redirect logging can still lose clicks) | Defines an org-level count-accuracy SLO (e.g., streaming within 0.5% at T+5min, ledger within 0.01% at T+1d) and makes it a contract between the ads platform and finance |
| Late data | Uses a 1-minute tumbling window | Sets watermark + allowed lateness by intent: 30s for dashboards, 24–72h for billing via batch recount | Decides the billing close policy: when a day is final, how post-close corrections are invoiced, and who approves restatements |
| Fraud | "Filter bots with a blocklist" | Separates real-time invalid-traffic filtering from post-hoc fraud adjustments; counts are versioned, not overwritten | Treats invalid-traffic rate as a business metric with an owner; sets refund policy with sales and legal; funds the fraud team as a cost-of-revenue line |
| Failure | "Kafka is replicated" | Names the restart double-count, the stalled watermark, and the lost-click-at-redirect failure; each has a metric and an owner | Designs the org's correction posture: every counted number is reproducible from the raw log for N days, and N is a finance/legal decision, not a Kafka config |
Why "first move" separates levels
L5: Starts with the pipeline because the pipeline is what the prompt seems to ask for. The result is one count that is expected to be simultaneously real-time, exact, fraud-filtered and final — a set of requirements no single path can meet.
L6: "Three different consumers want this number. Billing wants it exact and auditable and can wait a day. Advertiser dashboards want it in under a minute and tolerate ±1%. Budget pacing wants it in seconds and would rather over-count than under-count. I'll anchor on billing-grade correctness and show how the real-time path hangs off the same log."
L7: "Before I design the count, who owns the invoice number, and what's our contractual discrepancy tolerance with advertisers? If we've promised 'we bill from our logs, discrepancies above 10% are investigated', that's my correctness envelope — and the design is a reconciliation system, not a pipeline."
Why "exactly-once" separates levels
L5: Cites the engine feature. Flink's exactly-once covers state inside the job. It does not cover the redirect server that dropped a log line, the client that fired the click beacon twice, or the sink that appends instead of upserting.
L6: Walks the chain: "A click is generated once at the edge with a signed click_id. Ingest is at-least-once. Dedup by click_id in keyed state with a 24h TTL. Aggregates are written with an idempotent upsert keyed by (ad_id, window_start) so replay overwrites rather than adds. Exactly-once is a property of the whole path, and the weakest link is the edge log, so that's where I'd spend durability."
L7: Prices the guarantee. Exactly-once across the whole path costs checkpoint latency, dedup state (tens of GB), and engineering complexity. The L7 move is deciding which consumers actually need it — billing yes, a CTR heatmap no — and refusing to pay the exactly-once tax on every derived metric.
Why "late data" separates levels
L5: Picks a window size and moves on. Late clicks are silently dropped or silently added to the wrong window.
L6: Treats lateness as a per-intent budget. Dashboards close a minute window after 30s of lateness and accept a small undercount. Billing never closes in the stream; it closes in the nightly batch, which sees everything that arrived within 72h. Mobile SDKs that batch offline clicks get a separate late-arrival metric.
L7: Turns "when is a number final?" into a published policy: daily billing close at T+3d 00:00 UTC, corrections after close go into the next invoice as an adjustment line, restatements above a threshold require finance approval. The engineering design implements a policy; it doesn't invent one.
The Staff Positions#
| Position | Rationale |
|---|---|
| Separate billing-grade from real-time counts | Different latency, correctness and failure postures; one path cannot serve both without compromising each |
| Immutable raw click log is the source of truth | Every aggregate must be reproducible; aggregates are caches of the log, never the system of record |
Dedup on a server-minted click_id | Client-generated IDs can be forged or reused; mint and sign at the redirect edge |
| Idempotent upsert sinks, never append | Replay after restart must overwrite, not add; this is where most double-billing comes from |
| Kappa + nightly batch reconciliation | One logic definition; an independent recount catches what streaming got wrong |
| Redirect first, log asynchronously — but durably | The user's redirect must not wait on Kafka; the click must not be lost if Kafka is slow (local disk spool) |
| Fraud filtering versions counts, never deletes | Invalid clicks are reclassified, not erased; billing reads the latest classification, audit reads all of them |
The Three Intents#
| Intent | Constraint | Strategy | Failure Mode | Correctness Bar |
|---|---|---|---|---|
| Billing-grade counting | Exact, auditable, reproducible | Raw log → dedup → batch recount → ledger; streaming is only a preview | Over-billing (refunds, trust) or under-billing (lost revenue) | ≤0.01% vs recount; every number traceable to click IDs |
| Real-time advertiser dashboards | Fresh (<1 min), cheap to query | Stream aggregation into an OLAP store, pre-rolled by minute/hour | Stale or slightly wrong numbers; advertisers notice drift vs invoice | ±0.5–1% at T+5min; converges to ledger by T+1d |
| Budget pacing & fraud response | Seconds, err on the safe side | In-memory spend counters per campaign; fraud scores in the hot path | Overspend past budget (we eat the cost) or false-positive blocks | Directionally right within seconds; over-count preferred to under-count |
🎯 Staff Move: "I'll anchor on billing-grade counting, because that's where the money and the trust live. The real-time dashboard and the pacer hang off the same raw log with weaker guarantees. If I design for the dashboard first, I'll end up with a fast number that finance can't sign."
The Five Fault Lines#
| # | Fault Line | The Tension |
|---|---|---|
| 1 | Freshness vs Correctness | Show the number now (and correct it later) or show it only when it's final? |
| 2 | Exactly-once vs Throughput & Cost | Pay checkpoint latency and dedup state for exact counts, or accept at-least-once and reconcile? |
| 3 | Lambda vs Kappa | Two codebases (stream + batch) with drift, or one codebase whose correctness depends on replay? |
| 4 | Fraud Filtering Inline vs Post-hoc | Block invalid clicks before counting (fast, false positives) or reclassify later (accurate, requires corrections)? |
| 5 | Pre-aggregation vs Query Flexibility | Roll up at ingest (cheap queries, fixed dimensions) or keep raw events (any query, 10–100× cost)? |
In the Wild: Real Production Systems#
Why this section belongs here: Ad counting is one of the few domains where the big players published their designs. Citing them shows you know the problem is about correctness, not plumbing.
Google — Photon (joining queries and clicks exactly once)#
Google's Photon paper (SIGMOD 2013) describes joining the ad query log with the click log in real time across multiple datacenters with exactly-once semantics. The key mechanism is the IdRegistry: a Paxos-replicated store of event IDs that have already been joined, so a click processed in two datacenters is only emitted once. Photon tolerates datacenter loss while keeping the output at-most-once and eventually-exactly-once.
Staff insight: Exactly-once at Google's scale was solved with a global dedup registry keyed by event ID, not by trusting the stream engine. In an interview, the speakable version is: "dedup is a first-class component with its own SLO, not a side effect of checkpointing."
Google — Mesa (geo-replicated ads reporting warehouse)#
Mesa (VLDB 2014) is the store behind Google's ads reporting. It ingests updates in versioned batches every few minutes, keeps pre-aggregated tables keyed by dimensions, and makes each version atomically visible across datacenters. Corrections are just new versions; queries read a consistent version.
Staff insight: Versioned, atomically-published aggregates are the clean answer to "how do corrections and backfills not corrupt live dashboards?" You never mutate the number in place; you publish a new version.
Uber — Exactly-once ad events with Flink, Kafka and Pinot#
Uber Engineering described its ad event pipeline for Uber Eats: Flink jobs with exactly-once checkpoints, Kafka transactions between stages, deduplication on a record UUID, and Apache Pinot as the serving layer with upsert support so replays overwrite rather than add. A separate reconciliation path compares against offline Hive data.
Staff insight: This is the modern kappa-plus-reconciliation shape. Even with end-to-end exactly-once in the stream, they still kept an offline comparison — because the stream proving itself correct is not an audit.
Metamarkets → Apache Druid#
Druid was built at Metamarkets specifically for interactive analytics over programmatic ad-exchange data: roll-up at ingestion, columnar segments, bitmap indexes, and sub-second slice-and-dice over billions of events. It is now a common choice for real-time ad dashboards.
Staff insight: The OLAP store choice is driven by the query shape (group-by over a few dimensions with time filters), not by the write rate. Roll-up at ingest is the cost lever: 100× fewer rows if dimensions are bounded.
What Interviewers Probe#
| After You Say... | They Will Ask... | (What They're Evaluating) |
|---|---|---|
| "Flink gives us exactly-once" | "Your job restarts from a 3-minute-old checkpoint. What happens to the sink?" | Do you know exactly-once ends at the sink boundary? |
| "1-minute tumbling windows" | "A mobile user clicks offline and the event arrives 6 hours later. Is it billed?" | Lateness as a business policy, not a window parameter |
| "We'll dedup by click ID" | "Who generates the ID? How long do you remember it? How much memory is that?" | Can you size dedup state and reason about forged IDs? |
| "We'll filter bots" | "Fraud team finds a click farm 3 days later. What happens to the invoice already sent?" | Correction and restatement design |
| "Druid / Pinot / ClickHouse for queries" | "An advertiser wants clicks by zip code by hour for 90 days. Can your rollup answer that?" | Pre-aggregation tradeoff and dimension cardinality |
| "Shard by ad_id" | "One Super Bowl ad gets 40% of all clicks for 10 minutes." | Hot-key handling without wrong counts |
System Architecture Overview#
Reading the diagram: The redirect edge mints the
click_idand never waits on Kafka — a local spool absorbs Kafka hiccups. Kafka plus the raw archive is the system of record. The streaming path serves dashboards and pacing in seconds; the batch path recounts from raw and writes the ledger that invoices read. The reconciler is the component most designs forget: it compares the two and publishes corrections as new versions instead of mutating in place.
Quick-Reference: The 30-Second Cheat Sheet#
| Topic | The L5 Answer | The L6 Answer — Say This | The L7 Answer — Say This |
|---|---|---|---|
| Correctness | "Flink exactly-once" | "Exactly-once is a path property: signed click_id, dedup state, idempotent sink. Billing closes in batch." | "We publish a count-accuracy SLO per consumer and finance signs the billing one." |
| Late events | "Drop events after the window" | "Dashboards: 30s lateness. Billing: 72h via batch. Late-arrival rate is a metric." | "Billing close is a policy: T+3d final, adjustments on the next invoice." |
| Duplicates | "Use a set of seen IDs" | "Server-minted click_id, 24h TTL dedup in keyed state, ~1% dup rate expected from retries." | "Dedup is a platform service with its own SLO; every ad surface uses it." |
| Hot ads | "Shard by ad_id" | "Two-stage aggregation: salt the key, pre-aggregate, then merge. Counts stay exact." | "Super Bowl-class events get a capacity reservation and a named on-call." |
| Fraud | "Blocklist bots" | "Inline IVT filter for obvious cases; post-hoc reclassification versions the counts." | "Invalid-traffic rate is a revenue metric reviewed with sales; refunds are a policy." |
| Queries | "Store raw events and query" | "Roll up at ingest by bounded dimensions; raw goes to the lake for ad-hoc." | "OLAP cost is per dimension: every new breakdown is priced before it's approved." |
Key Numbers Worth Memorizing#
| Metric | Value | Why It Matters |
|---|---|---|
| Clicks/day (large ad network) | ~1–10B | 1B/day ≈ 11.6K/s average, 3–5× at peak |
| Impressions per click | ~100–1,000× | CTR of 0.1–1%; impressions dominate volume if you count them too |
| Redirect latency budget | p99 < 50ms | Every 100ms of redirect delay measurably loses landing-page sessions |
| Client retry duplicate rate | ~0.5–2% | Why dedup is mandatory, not optional |
| Dedup state (24h, 1B clicks) | ~16 bytes/ID × 1B ≈ 16 GB + overhead ≈ 30–50 GB | Fits in RocksDB-backed keyed state across the cluster |
| Invalid traffic filtered | ~5–20% of clicks | Varies by surface; fraud is a normal-case volume, not an edge case |
| Dashboard freshness target | < 1 min | What advertisers perceive as "real-time" |
| Streaming vs batch drift budget | ≤0.5% at T+5min, ≤0.01% at T+1d | The reconciliation SLO |
| Late-arrival tail (mobile) | 1–3% of clicks after 5 min; long tail to 24–72h | Why billing closes in batch |
| Rollup compression | 50–500× fewer rows | Minute × ad × few dimensions vs raw events |
| Commonly tolerated advertiser discrepancy | ~5–10% vs third-party trackers | Above this, accounts escalate and credits follow |
Interview Walkthrough
The most common mistake: Candidates spend 20 minutes drawing the click → Kafka → Flink → database pipeline, then discover at minute 35 that the interviewer wanted to know what happens to yesterday's invoice when a fraud ring is found. The phases below compress the pipeline to under 10 minutes so the rest goes to correctness, lateness, fraud and corrections — the things that actually level you.
Phase 1: Requirements & Framing (2–3 min)#
State functional requirements in 30 seconds:
"Users click ads, we redirect them to the advertiser, and we aggregate clicks by ad, campaign and time so advertisers can see performance and we can bill them."
Then spend the remaining time on the non-functional split, which is the whole interview:
"This count has three consumers with incompatible needs. Billing needs it exact and auditable, and can wait a day. Advertiser dashboards need it within a minute and can be off by a percent. Budget pacing needs it within seconds and should over-count rather than under-count. I'll anchor on billing-grade correctness and hang the real-time paths off the same log."
Then commit to scale and constraints:
"I'll assume 1B clicks a day — about 12K/s average, 50K/s peak — 10M active ads, redirect p99 under 50ms, dashboards fresh within a minute, and billing final at T+3 days. I'll treat duplicates, late events and invalid traffic as the normal case."
🎯 Staff Move: Saying "billing closes in batch, dashboards close in the stream" in the first three minutes tells the interviewer you know where the bodies are buried. It also licenses you to make the streaming path simpler later.
Phase 2: Core Entities & API (1–2 min)#
Name the nouns, not the schema:
- ClickEvent:
click_id(server-minted, signed),impression_id,ad_id,campaign_id,advertiser_id,user_key(hashed),ts_event,ts_ingest,ip,ua,geo,ivt_score - Aggregate:
(ad_id, window_start, dims…) → clicks, valid_clicks, spend_micros, version - LedgerEntry:
(advertiser_id, day, campaign_id) → billable_clicks, amount, version, closed_at
Click path (hot, user-facing):
GET /c/{signed_token}
→ 302 Location: https://advertiser.example/landing?...&gclid-like param
(token encodes impression_id, ad_id, landing URL, expiry, HMAC)
Query path (advertiser-facing):
GET /v1/reports?advertiser=A&campaign=C&from=...&to=...&granularity=hour&dims=geo,device
→ { rows: [...], data_freshness_ts, is_final: false }
🎯 Staff Move: "Every report response carries
data_freshness_tsandis_final. That one field is the contract that lets the dashboard be fast and the invoice be right without anyone calling support about the difference."
Phase 3: High-Level Architecture (≤5 min)#
Staff candidates spend under 5 minutes here. Draw the two paths and say one sentence per box.
"The redirect edge is the only synchronous component. Everything after Kafka is asynchronous. The raw log is the source of truth; the OLAP store and the ledger are both derived from it, at different speeds and with different guarantees."
Phase 4: Transition to Depth#
The sentence that steers the interviewer:
"The pipeline is the easy part. The three places this design actually breaks are: double-counting on replay, late and offline clicks crossing the billing boundary, and fraud discovered after the invoice goes out. I'd like to go deep on exactly-once end to end first, because that's where the money is — then late data, then corrections."
Phase 5: Deep Dives (25–30 min)#
Pick 3–4 based on interviewer signals. Have all of these ready:
| Deep Dive | The 60-Second Version | Go Here |
|---|---|---|
| Exactly-once end to end | Signed click_id at edge → at-least-once Kafka → keyed dedup (24h TTL) → windowed aggregate → idempotent upsert keyed by (ad_id, window_start). Checkpoint every 30–60s. | Fault Line 2, Appendix B |
| Late events and windows | Event time, not ingest time. Watermark = max event time − 30s for dashboards. Late arrivals after window close go to a side output and a correction stream. Billing waits 72h in batch. | Fault Line 1, Appendix A |
| Hot ads | Two-stage aggregation: key by (ad_id, salt 0–15), pre-aggregate per minute, then merge by ad_id. Exact counts, no single-task bottleneck. | Section 4.3, Drill 5 |
| Lambda vs kappa + reconciliation | One logic library compiled into both stream and batch; nightly recount diffs against stream; corrections published as new versions. | Fault Line 3 |
| Fraud | Inline rules for the obvious (datacenter IPs, rate per user_key, invalid signatures); post-hoc ML reclassifies within 72h; billing reads latest classification. | Fault Line 4, Deep Dive 4 |
| Query serving | Druid/Pinot/ClickHouse with ingest-time rollup on bounded dimensions; raw in the lake for ad-hoc. | Fault Line 5, Appendix D |
Phase 6: Wrap-Up (2–3 min)#
Close with constraints, evolution and what you'd build later:
"To summarize: one immutable click log, two derived products. Streaming gives dashboards and pacing in seconds with a 0.5% accuracy budget. Batch recount gives the ledger finance signs, with a 72-hour lateness window and versioned corrections. The things I'd build next are impression-click joins for CTR and attribution, a self-serve backfill tool with approval gates, and a per-advertiser discrepancy report so sales sees problems before the advertiser does. What I wouldn't build on day one is multi-region active-active counting — regional pipelines that merge at the ledger are enough until a region-loss SLA demands more."
Common Timing Mistakes#
| Mistake | Time Lost | Fix |
|---|---|---|
| Designing the redirect service's load balancer in detail | 5–8 min | One sentence: stateless, behind anycast/L7 LB, p99 < 50ms |
| Explaining Flink checkpoint barriers | 5 min | Link the concept, say what it guarantees and where it stops |
| Debating Druid vs Pinot vs ClickHouse features | 5 min | Pick one, name the query shape that drives it, move on |
| Ignoring the billing boundary until asked | Whole level | Name billing-grade vs real-time in Phase 1 |
| Never mentioning corrections | Whole level | "What happens when yesterday's number was wrong?" — answer it before they ask |
1. The Staff Lens#
1.1 Why This Problem Exists in Staff Interviews#
Ad click aggregation is a counting problem where the count is revenue. That one fact turns a textbook streaming exercise into a test of correctness discipline, operational ownership and organizational judgment. At a large ad network, a 1% overcount on $10B of annual click revenue is $100M of refunds and a credibility problem with every agency that runs a third-party tracker against you. A 1% undercount is the same money left on the table. And the latency-sensitive consumers — pacing and fraud — cannot wait for the exact answer.
Interviewers pick it because the L5 answer (a correct streaming pipeline) is plausible and complete-looking, and the L6 answer differs almost entirely in judgment: which number is final, when, and who signs it.
1.2 The L5 → L6 → L7 Contrast — Visual#
The L5 path is not wrong — every box is necessary. It's that the L5 path produces one number with no stated correctness envelope, and the first incident (a restart, a fraud ring, an offline-click backlog) has no designed answer.
1.3 The Staff Question That Cuts Through Everything#
"When this number is wrong, who finds out first — us, or the advertiser?"
If the answer is "the advertiser", there is no reconciliation. If the answer is "us, within a day, with a diff report and a correction already queued", the design is Staff-grade. Every component in this case study exists to make that answer true.
2. Problem Framing & Intent#
2.1 The Three Intents — Explained#
Billing-grade counting. The output is a ledger: billable clicks and amounts per advertiser per day, with a version number and a closed flag. It must be reproducible from the raw click log, defensible in a dispute, and stable once closed. Latency is measured in hours to days. The correctness bar is roughly ≤0.01% disagreement with an independent recount, and every billed click must be traceable to a click_id. Failure is financial: refunds, revenue restatement, and lost advertiser trust. The owner is usually a billing or ads-finance engineering team with finance as the business signatory.
Real-time advertiser dashboards. The output is a set of rollups queried by advertisers and account managers: clicks, spend, CTR by campaign/ad/geo/device/hour. Freshness under one minute, p95 query latency under 1–2s, and an accuracy budget of ±0.5–1% while not final. Failure is reputational: an advertiser sees a number on the dashboard that does not match the invoice and opens a ticket. The owner is the ads reporting team.
Budget pacing & fraud response. The output is a control signal: "campaign C has spent 98% of its daily budget — stop serving" or "this publisher's click rate jumped 40× — quarantine." Freshness in seconds. Accuracy is asymmetric: over-counting causes slightly early budget stops (some lost revenue), under-counting causes overspend that the network usually eats because advertisers are not charged beyond budget. The owner is ads delivery / ads quality.
🎯 Staff Move: "Pacing and billing look like the same counter, but they fail in opposite directions. The pacer should over-count — stopping early costs us a little revenue. The ledger must not over-count — over-billing costs us trust and refunds. That's why they can't share a counter."
2.2 When NOT to Use a Real-Time Click Aggregator#
- Small volume, monthly invoicing. Under ~1M clicks/day with monthly billing and no live dashboard requirement, a nightly SQL job over a warehouse table is correct, cheap and operable by one engineer. Streaming adds an on-call rotation for no product benefit.
- Dashboards that nobody watches live. If advertisers check reports once a day, hourly micro-batch (Spark every 15 minutes) delivers 95% of the value at a fraction of the operational cost.
- Conversion-based billing (CPA). If you bill on conversions reported by the advertiser days later, the click count is an input to attribution, not the billing number. Design the attribution join first.
- Impression-level analytics with unbounded dimensions. If product wants arbitrary slicing over raw events, you need a lakehouse with a query engine, not a pre-aggregating stream.
🎯 Staff Move: "If billing is monthly and nobody needs sub-hour freshness, I'd ship a batch pipeline and revisit when pacing becomes a requirement. Streaming is a cost we take on for a consumer, not a default."
2.3 What the Interviewer Leaves Underspecified#
| Underspecified | Why It Matters | What to Assume Out Loud |
|---|---|---|
| Pricing model (CPC, CPM, CPA) | Determines which event is billable | CPC: the click is the billable event |
| Who the consumers are | Determines correctness bar per output | Billing, dashboards, pacing — name all three |
| Lateness tolerance | Determines when a number is final | Dashboards 30s, billing 72h |
| Duplicate sources | Client retries, double-taps, bots, replays | ~1% dups; dedup on server-minted ID |
| Fraud expectations | Inline vs post-hoc; refund policy | Inline rules + post-hoc model within 72h |
| Query dimensions | Drives rollup design and OLAP cost | ad, campaign, advertiser, hour, geo (country), device |
| Retention / audit | Drives raw storage cost | Raw 13 months (covers annual audit); rollups 3 years |
| Multi-region | Drives dedup and ledger design | Regional pipelines, global ledger merge |
2.4 Precise Terminology#
| Term | Precise Meaning |
|---|---|
| Event time | When the click happened on the client/edge (ts_event). The only time that matters for billing windows. |
| Processing time | When the stream job sees the event. Useful for lag metrics, wrong for counting. |
| Watermark | The job's assertion that no events older than T are still expected. Windows fire when the watermark passes their end. |
| Allowed lateness | How long after the watermark a window can still be updated. Beyond it, events go to a side output. |
| Exactly-once (effectively-once) | Each click affects each output exactly once, even across retries and restarts. Achieved by dedup + idempotent or transactional sinks — not by never processing twice. |
| IVT (invalid traffic) | Clicks classified as non-human or fraudulent (GIVT = general, e.g. known bots; SIVT = sophisticated). Not billable. |
| Billing close | The moment a day's ledger becomes final. Later changes are adjustments, not edits. |
| Reconciliation | Comparing two independently computed counts and resolving the difference with a recorded correction. |
| Backfill | Recomputing aggregates for a past time range from the raw log, typically after a bug fix. |
| Rollup | Pre-aggregation at ingest that collapses raw events into rows per (dimensions, time bucket). |
3. The Five Fault Lines#
Each fault line is a place where two competent engineers will disagree. The Staff move is to name the tension, pick a default, and say who pays.
3.1 Fault Line 1: Freshness vs Correctness#
The tension: A number shown at T+30s will be wrong — late clicks haven't arrived, fraud hasn't been scored, duplicates from slow retries haven't been seen. A number shown only when final is useless for pacing and frustrating for advertisers.
| Strategy | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Show only final numbers (T+1d) | One number, never revised | Dashboards are a day stale; pacing impossible | Advertisers (no live feedback), ads delivery (overspend) |
| Show streaming numbers as final | Fresh, simple | Numbers silently differ from invoices; disputes | Account managers fielding tickets; finance on refunds |
| Show streaming as provisional + converge to final (Staff default) | Fresh preview, audited invoice, explicit is_final flag | Requires a reconciliation path and UI that shows provisional state | Reporting team builds the correction pipeline |
Staff default: Two closes. The stream closes a minute window when the watermark passes window_end + 30s and marks it provisional. The batch close happens at T+3d and marks the day final. Any late click between the two lands in a correction stream that updates the provisional number in place (by version).
Watermark math, made concrete: if 97% of clicks arrive within 5s, 99% within 30s, and 99.9% within 10 min, a 30s watermark delay means dashboard minute windows are ~99% complete when first shown and the remaining ~1% arrives as corrections. Billing at T+72h captures >99.99%.
When to deviate: If the product has no live dashboards (monthly reporting), drop the stream close entirely. If pacing is the main consumer (e.g., tight daily budgets on small advertisers), make the pacer even fresher — processing-time counters with no watermark at all — and accept that it over-counts.
🎯 Staff Move: "I'll show two numbers to advertisers: today's provisional spend, updated every minute, and yesterday's final spend once billing closes. The UI labels which is which. That label saves more support tickets than any accuracy improvement."
3.2 Fault Line 2: Exactly-once vs Throughput & Cost#
The tension: End-to-end exactly-once requires dedup state, checkpointed operator state, and a sink that is either transactional or idempotent. Each adds latency, memory and failure surface. At-least-once is simpler and faster but double-counts on every restart and retry.
| Strategy | What Works | What Breaks | Who Pays |
|---|---|---|---|
| At-least-once + increment sink | Simplest; highest throughput | Every restart replays and double-adds; client retries double-count | Advertisers (overbilled), finance (refunds) |
| Engine exactly-once + append sink | Internal state is right | Sink sees replays as new rows — the classic double-bill | Same as above, harder to spot |
| Dedup + idempotent upsert sink (Staff default) | Replays overwrite; no transactional sink needed | Upsert key design must include everything that defines the row | Reporting team owns key design |
| Two-phase commit sink (Kafka transactions) | True transactional exactly-once between stages | Output visible only after checkpoint (adds 30–60s); read_committed consumers required | Latency budget; operational complexity |
Staff default: Dedup by server-minted click_id in keyed state with a 24h TTL. Aggregate by (ad_id, window_start). Write the full window value with an upsert keyed by (ad_id, window_start, dims) — never count = count + n. Kafka transactions only between internal stages where downstream is also a stream.
Dedup sizing: 1B clicks/day × 24h TTL × ~16 bytes key ≈ 16 GB raw; with RocksDB overhead ~30–50 GB spread across, say, 64 parallel tasks ≈ under 1 GB per task. A Bloom filter front can cut RocksDB lookups by ~95% for the non-duplicate majority.
Where exactly-once stops: the redirect edge. If the redirect server returns the 302 and crashes before the event reaches Kafka or the local spool, the click is lost. That is at-most-once at the edge — and it's the only place the system under-counts. Fix: write to local disk spool synchronously (≈0.1ms on NVMe with group commit) before the 302, then ship asynchronously.
When to deviate: For derived metrics nobody bills on (CTR heatmaps, creative experiments), at-least-once with HyperLogLog or approximate counts is fine — the exactly-once tax isn't worth it.
3.3 Fault Line 3: Lambda vs Kappa#
The tension: Lambda runs a streaming path and a batch path, merging at serving. It gives an independent truth but two implementations of the same business logic that inevitably drift. Kappa runs one streaming implementation and reprocesses from the log for corrections — one codebase, but no independent check, and replay is bounded by log retention.
| Strategy | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Classic lambda | Batch is authoritative; streaming can be sloppy | Two codebases; IVT rules implemented twice differ by 0.3%; every logic change ships twice | Reporting engineers (double work), on-call (which one is right?) |
| Pure kappa | One codebase; replay is the backfill story | No independent audit; Kafka retention (7–30 days) bounds replay; reprocessing a month means a parallel job at 30× speed | Finance (no independent check), infra (replay capacity) |
| Kappa + batch reconciliation (Staff default) | One logic library run in both modes; batch checks and corrects rather than serving | Need a shared logic library and a correction protocol | Platform team owns the shared library |
Staff default: Write dedup, IVT rules and attribution once as a library (e.g., a pure function over a click record plus lookups). Run it in the stream job and in the nightly batch job over the raw archive. The batch output is the ledger; the diff between batch and stream is a monitored metric (recon.drift_pct), and corrections are published to the OLAP store as a new version for the affected windows.
When to deviate: If you genuinely cannot share code between stream and batch (different languages, legacy Hive SQL), accept lambda but make the diff metric and an owner for it non-negotiable. If billing is low-stakes (internal chargeback), pure kappa is fine.
🎯 Staff Move: "The question isn't lambda versus kappa, it's whether you have an independent recount. I'll run the same logic library in both, and treat batch as the auditor, not as a second product."
3.4 Fault Line 4: Fraud Filtering Inline vs Post-hoc#
The tension: Inline filtering blocks invalid clicks before they're counted — fast and clean, but limited to cheap signals and false positives can't be undone easily. Post-hoc classification uses richer signals (cross-session patterns, publisher-level anomalies, ML models) but arrives hours or days later, after dashboards and possibly invoices have shown the number.
| Strategy | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Inline only | Simple; dashboard equals invoice | Sophisticated fraud passes; can't use 24h patterns | Advertisers (pay for fraud), trust |
| Post-hoc only | Accurate classification | Dashboards inflated by 5–20%; pacing overspends on fraud | Advertisers see big downward revisions |
| Inline rules + post-hoc reclassification, versioned (Staff default) | Obvious IVT removed in seconds; sophisticated IVT removed before billing close | Counts move after display; need versioning and an explanation in the UI | Ads quality team owns both classifiers; reporting owns versioning |
Staff default: Inline: signature validation on the click token, known-bot user agents, datacenter IP ranges, per-user_key and per-IP rate caps (e.g., >5 clicks on the same ad in 60s). Post-hoc within 72h: ML model over session and publisher features. Every click carries an ivt_version; aggregates store clicks_total and clicks_valid; billing reads clicks_valid at close.
Who pays when fraud is found after close: the policy question. Staff answer: credits on the next invoice, issued automatically below a threshold (e.g., <2% of the advertiser's monthly spend), with account-manager approval above it.
When to deviate: In-house ads with no external advertisers (e.g., internal cross-promotion) can skip post-hoc entirely.
3.5 Fault Line 5: Pre-aggregation vs Query Flexibility#
The tension: Rolling up at ingest (minute × ad × geo × device) cuts storage and query cost by 50–500×, but every dimension not in the rollup is unanswerable without going back to raw. Keeping raw events in the OLAP store answers anything but costs 10–100× more and queries slow down as data grows.
| Strategy | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Raw events in OLAP | Any slice, any time | 1B rows/day; storage and query cost scale linearly; p95 > 5s on 90-day ranges | Infra budget; advertisers (slow dashboards) |
| Fixed rollups only | Cheap, fast (sub-second) | New dimension = backfill from raw; product asks weekly | Reporting team (backfill toil) |
| Tiered: minute rollups (7d) → hour rollups (13mo) → raw in lake (Staff default) | Fast dashboards, cheap retention, ad-hoc still possible offline | Two query paths; lake queries take minutes | Analysts accept minutes for ad-hoc |
Staff default: Druid or Pinot with ingest-time rollup on bounded-cardinality dimensions (country, device type, placement, ad, campaign). High-cardinality dimensions (zip code, user, URL) only in the lake. Every new dimension request is priced: cardinality × time buckets × retention.
Cardinality math: 10M ads × 1,440 minutes = 14.4B potential rows/day — but only active (ad, minute) pairs exist, typically ~1–5% → ~150–700M rows/day before other dims. Adding country (×~20 effective) and device (×3) multiplies the active combinations by maybe 5–10×, not 60×. Adding zip code (×40K) makes the rollup close to raw — that's the line.
🎯 Staff Move: "I'll roll up on dimensions with cardinality under a few hundred. Anything finer goes to the lake. When product asks for zip code on the dashboard, the answer is a price, not a no."
4. Failure Modes & Operational Reality#
4.1 Replay Double-Count — The Restart That Overbills#
Scenario: A TaskManager OOMs during a traffic spike. The job restarts from a checkpoint taken 4 minutes earlier. The sink is an append-style insert.
t=0: TaskManager OOM (state backend memory pressure)
t=+30s: JobManager restarts job from checkpoint #18422 (t=-4min)
t=+40s: Source rewinds Kafka offsets 4 minutes
t=+40s: 4 min of windows re-emitted → appended to OLAP
t=+5min: Dashboard shows a 2x spike for 4 minutes across all ads
t=+20min: Nobody notices — spike looks like traffic
t=+1day: If billing reads the stream: overbill ~0.3% of daily volume
Detection: sink.rows_written vs source.records_consumed ratio jumps; recon.drift_pct exceeds 0.5% at the next reconciliation; job.restarts_total correlated with spike.
Mitigation: Make the sink idempotent: upsert keyed by (ad_id, window_start, dims) with the full window value. Druid/Pinot: use upsert tables or replace segments by interval. Never INCR.
Prevention: Contract test in CI: kill the job mid-window, restart, assert sink totals equal a clean run.
Owner: Reporting pipeline team. Finance is notified if the reconciler shows drift >0.1% on any billed day.
4.2 Stalled Watermark — Windows Never Close#
Scenario: One Kafka partition stops receiving events (a producer region drained for maintenance). The watermark is the minimum across partitions, so it stops advancing. No windows fire. Dashboards freeze; the pacer is fed nothing.
t=0: Region B redirect fleet drained; its partitions go idle
t=+1min: Global watermark stuck at t=0
t=+5min: Dashboards show flat zero for the last 5 min
t=+12min: Pacer sees no spend; campaigns overspend budgets
t=+20min: Advertiser support: "my campaign stopped reporting"
Detection: watermark.lag_seconds > 120s; window.fired_total rate drops to zero; per-partition last_event_age_seconds.
Mitigation: Idle-source detection (mark partitions idle after 30–60s so they don't hold back the watermark). Pacer falls back to processing-time counters when event-time lag exceeds 60s.
Owner: Stream platform team owns watermark config; ads delivery owns the pacer fallback. Details of watermark mechanics: Stream Processing.
4.3 Hot Ad — One Campaign at 40% of Clicks#
Scenario: A Super Bowl spot drives a single ad to 40% of all clicks for 10 minutes — ~20K clicks/s on one key.
t=0: Ad airs; clicks on ad X go from 50/s to 20,000/s
t=+10s: The task owning ad X hits 100% CPU
t=+30s: Backpressure propagates to the source; all ads lag
t=+2min: Checkpoints time out (barrier stuck behind hot task)
t=+5min: Job restarts — and replays straight back into the hot spot
Detection: task.busy_ratio skew (one task at 100%, median at 30%); checkpoint.duration > timeout; consumer.lag rising on all partitions.
Mitigation: Two-stage aggregation. Stage 1 keys by (ad_id, hash(click_id) % 16) and pre-aggregates per minute; stage 2 keys by ad_id and sums 16 partial counts. Dedup happens before stage 1 keyed by click_id, which is uniformly distributed. Counts remain exact; hot-key throughput scales 16×.
Prevention: Always run two-stage for the click aggregation (the cost is ~10% more CPU and one extra shuffle). Pre-announced big events get a capacity reservation.
Owner: Reporting pipeline team; ads sales gives 48h notice for tentpole events.
4.4 Lost Clicks at the Edge — The Silent Undercount#
Scenario: Kafka brokers in one AZ degrade; producer buffers fill; the redirect service is configured to drop events when its in-memory buffer is full (to keep redirect latency low).
t=0: Kafka produce p99 goes from 5ms to 2s
t=+20s: Producer buffer (64 MB) full on 30% of redirect pods
t=+20s: Pods drop events silently, still return 302s
t=+45min: Kafka recovers
Impact: ~3% of 45 min of clicks never logged — never billed
Detection: redirect.events_dropped_total (must exist and must page); compare redirect.302_total vs kafka.clicks_produced_total per minute — a gap >0.1% is an alert.
Mitigation: Local disk spool: append to a local log before returning 302 (≈0.1ms with group commit), ship asynchronously, drain on recovery. Size for ≥1h of the pod's traffic.
Owner: Click edge team. This is under-billing — revenue loss, not advertiser harm — which is exactly why it goes unnoticed without an explicit metric.
4.5 Late-Data Surge — Offline Mobile Clicks#
Scenario: A mobile SDK release batches clicks while offline. After a carrier outage, 2M clicks from 6 hours earlier arrive in 10 minutes.
Detection: events.late_total by lateness bucket (>5 min, >1h, >24h); side_output.late_records rate.
Mitigation: Late events beyond allowed lateness go to a side output that feeds a correction job, updating the affected provisional windows as new versions. Billing is unaffected — the batch close at T+72h includes them.
Owner: Reporting team for corrections; mobile SDK team for the batching behavior (which should cap offline buffering at 24h and stamp ts_event from a server-synced clock).
4.6 Reconciliation Drift — Stream and Batch Disagree#
Scenario: An IVT rule change ships to the stream job but not to the batch job (or vice versa). Nightly recon shows a 0.4% gap that grows each day.
Detection: recon.drift_pct per advertiser per day, alert at >0.1% for billed days, >0.5% for any day. logic_library.version mismatch between stream and batch.
Mitigation: Shared logic library with a version stamped on every output row; deployment pipeline refuses to ship a stream version that the batch job doesn't also have.
Owner: Platform team owns the shared library; reporting team owns the recon alert.
4.7 Operational Reality Matrix#
| Failure | Detection Signal | Blast Radius | Mitigation | Owner |
|---|---|---|---|---|
| Replay double-count | recon.drift_pct, job.restarts_total | All ads, restart window | Idempotent upsert sinks | Reporting pipeline |
| Stalled watermark | watermark.lag_seconds > 120s | All dashboards + pacing | Idle-source detection; pacer fallback | Stream platform / ads delivery |
| Hot ad | task.busy_ratio skew, checkpoint timeouts | Entire job (backpressure) | Two-stage salted aggregation | Reporting pipeline |
| Edge click loss | redirect.302_total vs kafka.produced_total gap | Under-billing, all ads in affected pods | Local disk spool before 302 | Click edge |
| Late surge | events.late_total by bucket | Provisional numbers | Side output → correction stream | Reporting |
| Recon drift | recon.drift_pct > 0.1% | Billing accuracy | Shared logic library, version gate | Platform + reporting |
| Fraud found post-close | ivt.reclassified_after_close_total | Specific advertisers' invoices | Automated credits below threshold | Ads quality + billing |
| OLAP ingestion lag | olap.ingest_lag_seconds > 60s | Dashboards only | Scale ingestion; queries show freshness ts | Reporting |
5. Evaluation Rubric#
5.1 Level-Based Signals#
| Dimension | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| Intent | Designs one count | Separates billing / dashboard / pacing with different correctness bars | Defines per-consumer accuracy SLOs signed by finance and ads delivery |
| Correctness | Cites engine exactly-once | Exactly-once as path property; names the edge as the at-most-once gap | Decides which outputs deserve exactly-once and which don't; prices the difference |
| Late data | Fixed window size | Watermark + allowed lateness per intent; billing closes in batch | Publishes billing-close and restatement policy |
| Fraud | Blocklist | Inline + post-hoc with versioned counts | IVT rate as revenue metric; refund policy with sales/legal |
| Operations | "Kafka is replicated" | Named failures with metrics and owners | Correction posture: reproducibility window, recon SLO, game days for replay |
| Organization | Not discussed | Names owning teams for edge, pipeline, ledger | Draws the boundary between ads platform and finance; one counting platform for all ad surfaces |
5.2 Strong Hire Signals#
| Signal | What It Sounds Like |
|---|---|
| Splits the number by consumer | "Billing can wait a day; dashboards can't wait a minute. Those are two products." |
| Knows where exactly-once ends | "The engine guarantees its state. The sink and the edge are mine to make idempotent and durable." |
| Has a correction story | "Corrections are new versions of a window, never in-place edits. The UI shows when a number is final." |
| Sizes the dedup state | "1B IDs × 16 bytes × 24h is ~16 GB raw — fits in RocksDB across 64 tasks." |
| Treats fraud as normal volume | "I expect 5–20% IVT. The pipeline is designed around reclassification, not around exceptions." |
5.3 Lean No-Hire Signals#
| Signal | Why It Misses the Bar |
|---|---|
Uses Redis INCR per click as the billing counter | No replay story, no dedup, no audit; one failover loses or doubles counts |
| "Exactly-once because Flink" with an append sink | Does not understand the sink boundary — this is the classic double-bill |
| Uses processing time for windows | Billing on the wrong day after any lag; numbers differ on every replay |
| No answer to "what if fraud is found after invoicing?" | Missing the most consequential product question in the domain |
| Stores raw events in the OLAP store "for flexibility" without a cost estimate | Ignores the dominant cost driver |
5.4 Common False Positives#
- Deep Flink internals ≠ correct counting. Explaining checkpoint barrier alignment fluently is not the same as designing an idempotent sink.
- Naming lambda architecture ≠ having a reconciliation design. Two paths without a diff metric and an owner is just two sources of disagreement.
- "We'll use Kafka transactions" ≠ exactly-once to the OLAP store. Transactions end at Kafka; the OLAP sink still needs idempotency.
- Knowing Druid features ≠ knowing the query shape. The OLAP choice should follow from dimensions and freshness, not feature lists.
6. Interview Flow & Pivots#
6.1 Typical 45-Minute Shape#
| Phase | Time | Goal |
|---|---|---|
| Framing | 0–4 min | Three consumers, anchor on billing, state scale |
| Entities + API | 4–6 min | click_id, aggregate, ledger; is_final flag |
| Architecture | 6–11 min | Edge → log → stream + batch → OLAP + ledger |
| Deep dive 1 | 11–20 min | Exactly-once end to end |
| Deep dive 2 | 20–28 min | Late data and billing close |
| Deep dive 3 | 28–36 min | Fraud and corrections, or hot ads |
| Operations | 36–41 min | Failure matrix, metrics, owners |
| Wrap-up | 41–45 min | Evolution, what you'd defer |
6.2 How Interviewers Pivot — And What They're Testing#
| Pivot | What They're Testing | Strong Response |
|---|---|---|
| "Now count impressions too" | Volume scaling (100–1,000×) | Impressions get sampled or HLL for reach; exact counts only where billed (CPM) |
| "Advertisers want CTR in real time" | Stream-stream join | Join clicks to impressions by impression_id within a bounded window (e.g., 1h); unmatched clicks go to a late-join path |
| "We're going multi-region" | Dedup across regions | click_id embeds region; dedup is regional because a click is minted in one region; ledger merges globally |
| "Finance says invoices were wrong last month" | Correction and governance | Versioned recount, diff report per advertiser, credit policy |
| "Cut the infra bill by 40%" | Cost levers | Drop minute rollups after 7d, tier raw to cold storage, sample impressions |
6.3 What to Deliberately Skip#
- Load balancer and redirect fleet internals beyond "stateless, p99 < 50ms"
- Stream engine internals (barriers, RocksDB compaction) — reference Flink
- Ad serving and auction design — a different system
- Detailed ML fraud models — name the features and the latency, not the model architecture
6.4 Follow-Up Questions to Expect#
- "How long do you keep click IDs for dedup, and why that long?"
- "Your job restarted. Walk me through why the counts are still correct."
- "An advertiser's third-party tracker shows 8% fewer clicks than your invoice. What do you do?"
- "How do you backfill three weeks of data after a bug in the IVT rule?"
- "One ad has 40% of traffic. What breaks and how do you fix it without wrong counts?"
- "Which OLAP store, and what does a typical dashboard query cost?"
- "When is a day's number final, and who decides?"
7. Active Drills#
Drill 1: The Opening (Intent + Constraints)#
Prompt: "Design a system that aggregates ad clicks in real time."
Staff Answer
"Before I draw anything: who consumes this count? I see three consumers. Billing needs exact, auditable numbers and can wait a day. Advertiser dashboards need minute-fresh numbers and tolerate about a percent of error. Budget pacing needs second-fresh numbers and should err toward over-counting. These have incompatible correctness and failure postures, so I'll build one immutable click log and derive two products from it: a streaming preview for dashboards and pacing, and a batch-closed ledger for billing.
I'll assume CPC billing, 1B clicks/day — ~12K/s average, ~50K/s peak — redirect p99 under 50ms, dashboards fresh within a minute, billing final at T+3 days. I'll go: click edge and ID minting → log → dedup and windowing → serving → reconciliation and corrections → fraud."
Why this is L6:
- Splits the problem by consumer before choosing technology — every later decision references this split
- States correctness envelopes as numbers (±1%, T+3d), not adjectives
- Orders the walkthrough by risk: the edge (where clicks can be lost) and reconciliation (where money is decided) are on the list from the start
What L7 adds:
- Asks who signs the invoice number today and what discrepancy tolerance has been promised to advertisers — that contract is the real requirement
- Asks whether other ad surfaces (video, search, shopping) already count clicks separately; if so, the project is a counting platform, not a pipeline
❌ Common L5 Trap
"Clicks go to a click service that writes to Kafka. Flink consumes, aggregates per ad per minute, and writes to Cassandra. The dashboard reads from Cassandra."
Why this misses: Every box is reasonable, but there is no correctness envelope. The interviewer asks "is this the number we bill on?" and the candidate has to either claim the stream is exact (it isn't) or bolt on a batch path late with no reconciliation design.
Drill 2: Exactly-Once, End to End#
Prompt: "You said exactly-once. Prove it — walk me from the click to the dashboard."
Staff Answer
"Five links, each with a different mechanism:
- Edge: the redirect server mints
click_id(region + timestamp + random, HMAC-signed with the impression token). Before returning the 302 it appends the event to a local disk spool. This is the at-most-once boundary — I make it durable rather than pretending it's exactly-once. - Transport: spool → Kafka with idempotent producers (
enable.idempotence=true). At-least-once overall; client retries and spool replays create duplicates by design. - Dedup: stream job keys by
click_id, keeps a seen-set in RocksDB state with a 24h TTL. First occurrence passes; later ones go to aduplicatesside output with a counter. - Aggregation: event-time tumbling 1-minute windows keyed by
ad_id(two-stage salted for hot ads). State is checkpointed every 30s. - Sink: each fired window writes the complete value with an upsert keyed by
(ad_id, window_start, dims). A restart replays from the last checkpoint and overwrites the same keys with the same values.
Exactly-once is the composition: durable edge + at-least-once transport + dedup + checkpointed state + idempotent sink. Break any one and you get double or lost counts."
Why this is L6:
- Locates the at-most-once boundary honestly (the edge) and spends durability there
- Distinguishes engine exactly-once (state) from output exactly-once (sink idempotency)
- Uses overwrite semantics so no transactional sink is required for the OLAP store
What L7 adds:
- Makes dedup a shared platform service with an SLO used by every ad surface, rather than re-implemented per pipeline
- Asks whether every downstream metric actually needs this — CTR experiments can run at-least-once with HLL and save ~30% of the job's cost
❌ Common L5 Trap
"Flink's checkpointing gives exactly-once, so every click is counted once."
Why this misses: The interviewer asks "what does your sink do when the job replays 3 minutes?" If the answer is "inserts rows", the design double-counts on every restart. Exactly-once inside the engine says nothing about side effects outside it.
Drill 3: Late and Offline Clicks — Make It Concrete#
Prompt: "A click happens at 23:59:50 on Monday on a phone in airplane mode. It arrives at 06:00 Tuesday. Which day is it billed on, and what does the dashboard show?"
Staff Answer
"It's billed on Monday — billing uses event time. Mechanically:
- Stream: Monday's 23:59 window fired at about 00:00:30 Tuesday when the watermark passed it. At 06:00 the event is ~6h late, beyond the stream's 5-minute allowed lateness, so it goes to the late side output. The correction job updates Monday 23:59 for that ad to version n+1 in the OLAP store. Monday's dashboard total ticks up by one click; it's still marked provisional.
- Batch: Monday's billing close runs Thursday 00:00 UTC (T+72h). The batch recount reads everything in the raw archive with
ts_eventon Monday, including this click. The ledger includes it. - Guardrail:
ts_eventfrom a phone can't be trusted blindly. The SDK stamps both device time and a server-synced offset; the redirect edge rejects event times more than 72h in the past or 5 minutes in the future, and those go to a quarantine counter."
Why this is L6:
- Event time for billing, with a precise mechanism for both paths
- Separates "late for the dashboard" (correction stream) from "late for billing" (batch close window)
- Treats client timestamps as untrusted input with explicit bounds
What L7 adds:
- Turns the 72h close into a published policy with finance and the advertiser terms of service — clicks later than close are dropped or billed as next-period adjustments, decided once, not per incident
- Tracks the late-arrival distribution as a product metric for the SDK team, because shrinking the tail is cheaper than widening the window
❌ Common L5 Trap
"We'll use processing time so we don't have to deal with late events."
Why this misses: Processing-time windows mean the click lands on Tuesday, a replay moves it again, and two runs over the same data give different invoices. Billing on processing time is non-reproducible by construction.
Drill 4: The Pipeline Is Down#
Prompt: "Kafka is unavailable for 20 minutes. What happens to clicks, redirects, dashboards and budgets?"
Staff Answer
"Each consumer degrades differently:
- Redirects: unaffected. The redirect path never waits on Kafka; it appends to the local disk spool (~0.1ms) and returns the 302.
- Clicks: not lost. Each redirect pod spools ~20 min × its share of 50K/s. At 300 bytes/click and 100 pods, that's ~180 MB per pod — trivially within a 50 GB spool. On recovery, spools drain at a capped rate (2× normal) to avoid a thundering herd on Kafka.
- Dashboards: freeze at the last fired window and show
data_freshness_ts20 minutes old. No wrong numbers, just stale ones. - Pacing: this is the dangerous one. With no spend signal, campaigns can overspend. The pacer switches to a conservative mode after 2 minutes without data: it throttles delivery to the campaign's recent spend rate × 0.8, and stops campaigns already above 90% of budget.
- Billing: unaffected — clicks reach the raw archive after the drain, before T+72h close."
Why this is L6:
- Decouples user-facing latency from the pipeline, and sizes the spool with arithmetic
- Recognizes pacing as the consumer that turns an outage into money lost, and designs its degraded mode
- Caps drain rate to avoid causing the next incident
What L7 adds:
- Pre-agrees with ads delivery and finance who absorbs overspend during pipeline outages (usually the network) and makes the conservative-pacing threshold their decision, not the pipeline team's
- Adds a quarterly game day: kill Kafka for 10 minutes in one region and verify spool drain, pacer fallback and dashboards labels
❌ Common L5 Trap
"Kafka is replicated across three brokers, so it won't be down."
Why this misses: The question was about behavior, not probability. Replication doesn't help when the cluster is unavailable for a config push, a controller bug or a network partition, and it says nothing about pacing.
Drill 5: Hot Ad#
Prompt: "A Super Bowl ad gets 40% of all clicks for 10 minutes. What breaks, and how do you keep counts exact?"
Staff Answer
"Keyed by ad_id, all 20K clicks/s land on one task. That task saturates, backpressure stalls the whole job, checkpoints time out, and a restart replays straight back into the hot spot.
Fix is two-stage aggregation:
- Dedup first, keyed by
click_id— uniform, no skew. - Stage 1 keyed by
(ad_id, salt)wheresalt = hash(click_id) % 16; tumbling 1-minute pre-aggregate. - Stage 2 keyed by
ad_id, sums the 16 partials per window and writes the upsert.
Counts stay exact because each click lands in exactly one salt bucket and stage 2 sums them deterministically. Cost: one extra shuffle of partial aggregates — tiny (16 rows/ad/minute) — and ~10% more CPU. I'd run two-stage always rather than detect hotness dynamically; adaptive salting is a source of subtle bugs during replay.
The pacer has the same hot key: I'd keep the campaign's spend counter in the pacer's local memory per shard and sum across 16 shards every second."
Why this is L6:
- Explains the failure cascade (backpressure → checkpoint timeout → replay into hotspot), not just "it's slow"
- Keeps exactness with deterministic salting tied to click_id
- Chooses static over adaptive salting for replay determinism
What L7 adds:
- Tentpole events get a capacity-reservation process with sales: 48h notice, pre-scaled cluster, named on-call
- Notes that the redirect edge and the fraud system see the same hot key and need the same review — the event is org-wide, not a pipeline problem
❌ Common L5 Trap
"Shard by ad_id across more partitions."
Why this misses: More partitions don't help one key — it still hashes to one partition and one task. The interviewer is checking whether you know the difference between many keys and one hot key.
Drill 6: Fraud Found After the Invoice#
Prompt: "Three days after billing closed, ads quality identifies a click farm that generated 4% of one advertiser's clicks last week. What happens?"
Staff Answer
"Three things, in order:
- Reclassify, don't delete. The fraud system writes new classifications for the affected
click_ids withivt_version = v+1. The raw log is unchanged. - Recompute affected aggregates as new versions. A targeted backfill reruns the batch logic for that advertiser and those days with the new classifications, producing ledger version n+1 and OLAP corrections for the affected windows. The dashboards now show the lower number with a 'revised' marker.
- Issue an adjustment, not a re-invoice. The billing system computes the delta between ledger versions n and n+1 and issues a credit on the next invoice. Below an automatic threshold — say 2% of monthly spend — it's automatic; above it, the account manager approves and contacts the advertiser first.
The audit trail shows: original count, reason for change, ivt_version, who approved, credit issued."
Why this is L6:
- Never mutates history; every change is a new version with a reason
- Uses targeted backfill, not a full recount
- Connects the data correction to the commercial process (credit, approval threshold)
What L7 adds:
- Writes the restatement policy with finance and legal: thresholds, approval chain, notice period, how it appears in revenue reporting
- Tracks post-close IVT rate as a KPI for ads quality — a rising rate means the inline filter is underperforming, which is cheaper to fix than to keep issuing credits
❌ Common L5 Trap
"Delete the fraudulent clicks from the database and rerun the aggregation."
Why this misses: Deleting destroys the audit trail. When the advertiser asks "why did my numbers change?", nobody can show what was removed and why.
Drill 7: Build vs Buy for the OLAP Store#
Prompt: "Druid, Pinot, ClickHouse, BigQuery, or build on Cassandra? Pick one and defend it."
Staff Answer
"The query shape decides: advertisers filter by advertiser_id and time range, group by one to three bounded dimensions, and expect sub-second p95. Ingest is streaming at ~50K events/s pre-rollup, and corrections must overwrite windows.
I'd pick Pinot or Druid for the dashboard: native Kafka ingestion, ingest-time rollup, time-partitioned segments, and support for upserts (Pinot) or segment replacement by interval (Druid) for corrections. ClickHouse is an equally good choice if the team already runs it — ReplacingMergeTree handles idempotent writes, though dedup is eventual until merge, so queries need FINAL or argMax patterns.
I wouldn't put dashboards on BigQuery/Snowflake directly — per-query cost and multi-second latency at thousands of advertiser queries per minute — but I'd use them for the lake and ad-hoc. I wouldn't build on Cassandra: no ad-hoc group-by, so every new breakdown becomes a new table.
Build vs buy: managed Pinot/Druid/ClickHouse if available in our cloud; self-host only if we have a team that already runs one."
Why this is L6:
- Derives the choice from the query shape and correction requirement
- Knows the idempotency mechanism in each candidate
- Separates dashboard serving from ad-hoc analytics
What L7 adds:
- Checks for an existing org OLAP standard — a second OLAP engine is a multi-year operational commitment
- Prices it: at ~500M rollup rows/day with 13-month retention, the dashboard store is ~$15–40K/month; the lake for raw is a separate line
❌ Common L5 Trap
"Store raw events in Elasticsearch and aggregate at query time."
Why this misses: 1B documents/day with aggregation at query time is the most expensive way to answer a fixed group-by. It also gives no clean story for idempotent corrections.
Drill 8: Changing an IVT Rule Without an Incident#
Prompt: "Ads quality wants to tighten the click-rate rule from 5 to 3 clicks per user per ad per minute. How do you ship it?"
Staff Answer
"An IVT rule change is a billing change — it moves revenue. So:
- Shadow: run the new rule alongside the old in the stream, writing
ivt_shadowwithout affecting counts, for 7 days. Report the delta per advertiser. - Review: if the global reduction is, say, 0.8% of billable clicks but 9% for two advertisers, the account managers for those two are told before enforcement.
- Enforce from a date boundary: the new rule applies to event times from 00:00 UTC on day D. Both stream and batch read the same rule version from the shared library, keyed by event-time effective date — so replaying day D−1 uses the old rule.
- Monitor:
ivt.rate_by_rule_version,recon.drift_pct, advertiser-level billable delta.
Effective-dated rules are the important part: a rule change must never retroactively change closed days unless that's an explicit restatement."
Why this is L6:
- Treats a classifier change as a financial change with shadow evaluation
- Uses effective dating so replays are reproducible
- Communicates to impacted advertisers before enforcement
What L7 adds:
- Sets a governance process: ads quality owns rules, finance must sign changes projected to move revenue by >0.5%, and every rule change appears in the monthly revenue bridge
- Publishes IVT methodology externally (as major ad platforms do) so advertisers trust the change
❌ Common L5 Trap
"Update the rule in config and redeploy the Flink job."
Why this misses: The next backfill uses the new rule on old data, silently changing closed invoices. And two advertisers lose 9% of their reported clicks with no warning.
Drill 9: Cost — Cut the Bill by 40%#
Prompt: "The aggregation platform costs $250K/month. Finance wants 40% off. Where do you cut?"
Staff Answer
"First, a breakdown. In a typical setup: raw storage and lake ~25%, OLAP serving ~30%, stream compute ~25%, Kafka ~10%, batch recount ~10%.
Biggest levers, in order of safety:
- Retention tiering (~$30–40K): minute-grain rollups only for 7 days, hour-grain for 13 months; raw clicks to infrequent-access storage after 30 days. Billing and audit only need hourly and raw.
- Kill unused dimensions (~$20–30K): query logs usually show 20–30% of rollup dimensions are never queried. Each removed dimension shrinks the rollup multiplicatively.
- Impressions sampling (~$20K if impressions are counted): exact counts only for CPM-billed placements; 10% sampling for CTR analytics.
- Right-size stream compute (~$15K): checkpoint interval from 10s to 60s, RocksDB tuning, autoscaling off-peak.
That's ~$100K, 40%, without touching billing correctness. I'd explicitly not cut the batch recount — it's 10% of cost and 100% of our audit story."
Why this is L6:
- Starts with a cost breakdown, not a guess
- Orders cuts by risk to correctness
- Names what not to cut and why
What L7 adds:
- Introduces chargeback: each product team that requested a dimension sees its cost, which stops dimension sprawl structurally
- Compares to the revenue at stake: $250K/month to count $1B/month of clicks is 0.025% — the case for not over-cutting
❌ Common L5 Trap
"Switch from exactly-once to at-least-once to save compute."
Why this misses: It saves maybe 10% of stream compute and puts billing correctness at risk. The interviewer wants to see you protect the expensive-to-lose property and cut where nobody notices.
Drill 10: Multi-Region#
Prompt: "We're launching in the EU with data residency. How does counting work across regions?"
Staff Answer
"A click is minted in exactly one region — the region whose redirect edge served it — and click_id encodes that region. So dedup is regional: no cross-region dedup needed except for the rare client that retries against a different region after failover, which I'd handle by having the redirect token bind to the serving region and rejecting cross-region replays of the same token for 24h via a small replicated deny-set.
Each region runs its own edge, Kafka, stream job and OLAP. EU raw click data stays in the EU. The ledger is global but only needs aggregates — billable clicks per advertiser per day per region — which are not personal data and can be replicated to the billing region.
Advertiser dashboards query a federation layer that fans out to regional OLAP stores; for an advertiser running in both regions, the API merges results. p95 goes up by one cross-region hop (~80–150ms), which is acceptable for a dashboard."
Why this is L6:
- Uses the minting region to avoid global dedup coordination
- Separates residency-scoped raw data from replicable aggregates
- Accepts and quantifies the dashboard latency cost
What L7 adds:
- Validates with legal that per-advertiser daily aggregates are outside residency scope before committing — a one-way door if wrong
- Designs regions as cells: a region's pipeline failure affects only that region's dashboards, and the billing close for that region can be delayed independently
❌ Common L5 Trap
"Replicate the Kafka topic globally and run one global Flink job."
Why this misses: It violates residency, adds cross-region latency to every event, and makes one job a global single point of failure.
8. Deep Dive Scenarios#
Deep Dive 1: Peak-Traffic Incident — The Super Bowl Spike#
Context: During a major televised event, click volume goes from 15K/s to 90K/s in under a minute. Consumer lag climbs to 8 minutes, dashboards freeze, and the budget pacer is starving. Three large advertisers' campaigns have overspent their daily budgets by 30%. The on-call escalates to you.
Questions to Surface First:
- Is the lag uniform across partitions, or concentrated on tasks owning a few hot ads?
- Is the pacer reading event-time windows (now 8 minutes stale) or processing-time counters?
- Are clicks being lost anywhere, or only delayed? Check
redirect.302_totalvskafka.produced_total. - Who decided the overspend policy — does the network absorb it, or is it billed?
Typical L5 Approach: Scales the stream job (more TaskManagers, more parallelism), waits for lag to drain, and confirms no data loss. Correct, but treats it as a throughput problem. If the hot key is the issue, more parallelism doesn't help; and the overspend continues while it drains.
Staff Approach: Separates the two problems: pacing is a control loop that must not depend on a lagging pipeline, so switch the pacer to its processing-time fallback immediately and apply conservative throttling to campaigns near budget. Then diagnose the pipeline: busy-ratio skew shows two hot ads; confirm two-stage aggregation is enabled (it was disabled in a config refactor). Re-enable, restart from checkpoint, lag drains in ~6 minutes. Overspend is credited per policy.
Principal Approach: Treats the incident as evidence that tentpole events have no owner. Creates a "tentpole readiness" process with sales: events are registered 48h ahead with expected volume, the pipeline is pre-scaled, a named incident commander is on call, and pacing switches to conservative mode during the event window by default. Prices the overspend: credits issued vs cost of pre-scaling (a few thousand dollars of compute) — the readiness process pays for itself in one event.
Staff Approach — Full Reasoning
| Phase | What to Do |
|---|---|
| Immediate (0–5 min) | Switch pacer to processing-time fallback; throttle campaigns >90% of budget. Verify no click loss at the edge. |
| Triage | Busy-ratio per task: uniform → capacity; skewed → hot key. Check checkpoint durations. |
| Quick fix | Re-enable two-stage aggregation; scale out; restart from checkpoint. Drain is safe because sinks are idempotent. |
| Guardrails | Cap Kafka fetch rate during drain to avoid OLAP ingestion overload; keep pacer on fallback until lag < 30s. |
| Post-mortem | Why was two-stage disabled without a test catching it? Why does the pacer depend on the windowed path at all? |
Metrics to Watch:
consumer.lag_seconds, task.busy_ratio (max vs median), checkpoint.duration_ms, pacer.mode (normal vs fallback), campaign.overspend_pct, redirect.events_dropped_total.
Organizational Follow-up:
Add a hot-key replay test to CI. Make pacer fallback automatic on watermark.lag_seconds > 60. Create the tentpole calendar with sales.
Ownership Question: "Who pays for the 30% overspend?" Staff answer: Per policy, the network — advertisers are never billed beyond their budget. That's why ads delivery, not the pipeline team, owns the pacer's fallback thresholds: they own the cost.
Key Takeaway: "Pacing is a control loop. Never let a control loop depend on a pipeline whose latency you don't control."
What clears the Staff bar:
- Splits the control-loop problem from the throughput problem
- Diagnoses hot key vs capacity with a specific metric
- Ties overspend to a named cost owner
Deep Dive 2: Silent Failure — The 3% Undercount Nobody Noticed#
Context: Finance notices that click revenue for the last six weeks is ~3% below forecast with no traffic change. Advertisers aren't complaining — they're being under-billed. You're asked to find out why.
Questions to Surface First:
- Do edge counts (
redirect.302_total) match logged clicks (kafka.produced_total)? - Did dedup start dropping legitimate clicks — e.g., a
click_idcollision? - Did an IVT rule change ship around six weeks ago?
- Does the batch recount agree with the stream? (If both agree, the loss is upstream of both.)
Typical L5 Approach: Checks the stream job for errors, finds none, checks the reconciliation, finds stream and batch agree. Concludes the counts are correct and traffic must have changed.
Staff Approach: Stream and batch agreeing means the loss is upstream of the log. Compares 302 count vs produced count per pod: a new redirect release six weeks ago changed the Kafka producer buffer config and drops events on buffer-full under load — about 3% at peak. No metric existed for dropped events. Fix: local disk spool and a
redirect.events_dropped_totalalert. Revenue is unrecoverable for clicks never logged — except that the 302 access logs on the load balancer still have them for 30 days; a one-time recovery job reconstructs the last 30 days from those logs.
Principal Approach: Recognizes that reconciliation only compares two things derived from the same log — it cannot detect losses before the log. Mandates an edge-to-ledger completeness check as a standard for every revenue pipeline: an independent count at the earliest point (load balancer access logs) reconciled daily to the ledger. Adds "revenue completeness" to the monthly business review so a 3% gap surfaces in days, not weeks.
Staff Approach — Full Reasoning
| Dimension | Staff Answer |
|---|---|
| Root cause | Producer drops events on buffer-full; no metric; reconciliation blind because both paths read the same log |
| Immediate action | Ship local disk spool; add dropped-events metric and alert |
| Recovery | Rebuild lost clicks from LB access logs (30-day retention) through the same dedup + IVT logic; bill as adjustment |
| System fix | Independent edge count reconciled daily to the ledger |
| Broader question | Which other revenue pipelines have no completeness check before their log? |
Metrics to Watch:
redirect.302_total, redirect.events_dropped_total, edge_to_log.gap_pct, ledger.vs_lb_logs_gap_pct.
Organizational Follow-up: The redirect team didn't know their producer config was a revenue control. Add revenue-path ownership tags to services and require review from billing engineering for config changes on tagged services.
Ownership Question: "Who owned the missing metric?" Staff answer: The click edge team owned the code; nobody owned the end-to-end completeness. The fix is assigning that to billing engineering explicitly, with the edge team accountable for emitting the signal.
Key Takeaway: "Reconciliation between two derived counts can't detect loss upstream of both. You need one count taken before the log."
What clears the Staff bar:
- Reasons about what reconciliation can and cannot detect
- Finds an independent source for recovery
- Turns a missing metric into an ownership fix
Deep Dive 3: Large Customer Onboarding — An Agency Brings 20% of Volume#
Context: Sales signs a holding company whose agencies will run ~20% of total click volume, starting in 3 weeks. They require hourly reports within 15 minutes, a contractual discrepancy cap of 3% vs their third-party tracker, and a monthly reconciliation file at click level.
Questions to Surface First:
- What does their third-party tracker count — clicks or landing-page arrivals? (The definitions differ by 5–10% naturally.)
- Do they need click-level data exports? What PII is in a click record?
- Is 20% of volume concentrated in a few campaigns (hot keys) or spread?
- Who agreed to 3%, and is it measurable with our definitions?
Typical L5 Approach: Capacity-plans for +20% throughput, adds an export job for click-level files, and adds a dashboard for the agency. Treats it as scaling work.
Staff Approach: The capacity is the easy part. The 3% cap is the risk: our click count and their tracker's landing-page count differ structurally (bounces before the landing page loads, IVT we filter that they don't). Builds a per-advertiser discrepancy report that compares our counts with theirs daily and alerts at 2%, so account managers see problems before the agency. The click-level export is hashed IDs only, through the existing data-sharing pipeline with legal sign-off.
Principal Approach: Pushes back on the contract term before it's signed: a 3% cap on a metric with a natural 5–10% definitional gap is a commitment engineering can't keep. Works with sales and legal to restate the term in terms of a shared, defined measurement (e.g., clicks with a matching landing-page beacon, using a shared click ID passed as a URL parameter). Turns the bespoke work into a standard "enterprise reporting tier" product with a price.
Staff Approach — Full Reasoning
| Phase | What to Do |
|---|---|
| Week 1 | Load-test at +30% with the agency's expected skew; confirm two-stage aggregation handles their top campaigns |
| Week 1 | Define the discrepancy metric with the agency: which events, which IDs, which time zone |
| Week 2 | Build daily discrepancy report; alert at 2% |
| Week 2 | Click-level export with hashed IDs; legal review |
| Week 3 | Shadow period: compute reports for 7 days before go-live |
Metrics to Watch:
advertiser.discrepancy_pct, report.freshness_minutes per advertiser, export.delivery_success.
Organizational Follow-up: Sales must involve reporting engineering before committing measurement SLAs. Add a contract-review checklist item.
Ownership Question: "Who owns the 3% commitment?" Staff answer: The account team owns the commitment; reporting engineering owns the measurement and the alert. If engineering wasn't in the room when it was signed, that's the root cause to fix.
Key Takeaway: "Discrepancy caps are measurement contracts. Define the measurement before you sign the cap."
What clears the Staff bar:
- Identifies the definitional gap as the real risk, not throughput
- Builds detection that reaches the account team first
- Connects engineering to the sales process
Deep Dive 4: Post-Mortem — Three Weeks of Wrong Invoices#
Context: A change to the attribution logic in the shared library shipped with a bug that counted clicks on video companion ads twice (once for the video, once for the companion). It went unnoticed for three weeks because both stream and batch used the same library. 140 advertisers were overbilled by an average of 1.2%. You're running the post-mortem.
Questions to Surface First:
- Why did reconciliation not catch it? (Same library in both paths — they agreed with each other.)
- Was there any independent signal? (Advertiser discrepancy reports? Edge counts?)
- How far back can we recompute? (Raw archive retention: 13 months — yes.)
- Which invoices are affected, and what's the credit process?
Typical L5 Approach: Fixes the bug, backfills three weeks, issues credits, adds a unit test.
Staff Approach: Fixes and backfills, but the core finding is that shared logic makes stream and batch agree on bugs. Adds two independent checks: (1) an invariant — billable clicks ≤ 302 redirects at the edge per ad per hour (a click can't be counted more times than it was redirected); (2) a canary comparison — new library versions run in shadow against the previous version for 48h, and any per-advertiser delta >0.2% blocks promotion.
Principal Approach: Recognizes the trade: a shared library eliminated drift but created correlated failure. Establishes that revenue-affecting logic changes follow a financial change process — shadow, delta report, finance sign-off above a threshold — and that the edge-count invariant is a standard for every revenue counter in the company. Owns the external communication with finance: proactive credits and a note to affected advertisers before they find it, which costs less trust than being caught.
Staff Approach — Full Reasoning
| Phase | What to Do |
|---|---|
| Immediate | Roll back the library version; freeze billing close for the current period |
| Scope | Targeted backfill of affected days and advertisers from the raw archive with the fixed library; diff ledger versions |
| Remediation | Credits on next invoice; proactive notice for advertisers above 1% impact |
| Guardrails | Edge-count invariant; 48h shadow for library versions; delta gate at 0.2% |
| Post-mortem | Correlated failure from shared logic; no invariant independent of the library |
Metrics to Watch:
invariant.billable_gt_redirects_total (must be zero), library.shadow_delta_pct, ledger.version_delta_by_advertiser.
Organizational Follow-up: Library changes touching attribution require review from billing engineering. Monthly revenue bridge includes "logic-change impact".
Ownership Question: "Who approves a backfill that changes closed invoices?" Staff answer: Finance, with billing engineering providing the delta report. Engineers can compute a restatement; they can't authorize one.
Key Takeaway: "Shared logic removes drift and adds correlated failure. Pair it with an invariant that doesn't use the shared logic."
What clears the Staff bar:
- Spots that agreement between paths isn't evidence of correctness
- Adds an independent invariant grounded in the edge
- Separates computing a restatement from approving one
Deep Dive 5: Multi-Region Expansion — Active-Active Redirects#
Context: Leadership wants the click edge in four regions, active-active, so that losing one region doesn't stop redirects. Today everything runs in one US region. Billing must remain correct during and after a regional failover.
Questions to Surface First:
- Does a region failover cause clients to retry the same click in another region? (Yes, if the client retries a failed redirect.)
- Does any data-residency requirement prevent centralizing raw clicks?
- What's the RPO for click logs if a region is lost — can we lose the last N seconds?
- Is the dashboard expected to be globally consistent?
Typical L5 Approach: Replicate Kafka across regions with MirrorMaker, run a global Flink job reading the merged topic, dedup globally.
Staff Approach: Regional cells. Each region mints
click_ids with a region prefix and runs its own edge, Kafka, stream and OLAP. Cross-region duplicates only occur on client retry after failover; the signed click token carries the original region, and a small globally-replicated deny-set (24h TTL, ~MBs) rejects the second mint. The ledger merges regional daily aggregates. RPO: the local spool is replicated to a second AZ, so a single-AZ loss loses nothing; a full region loss loses at most the unshipped spool (seconds).
Principal Approach: Asks what failure leadership is actually buying insurance against and prices it: four-region active-active is roughly 3–4× the edge cost and 2× the pipeline cost. If the real need is "redirects keep working", only the edge must be multi-region — the pipeline can stay in two regions with async shipping. Writes down the regional failure posture: which region's billing close can be delayed, who decides, and how advertisers are told.
Staff Approach — Full Reasoning
| Phase | What to Do |
|---|---|
| Design | Region-prefixed click_id; regional dedup; global deny-set only for cross-region token reuse |
| Data path | Regional Kafka + stream + OLAP; aggregates replicated to the billing region |
| Failover | DNS/anycast shifts redirects; deny-set prevents double mints; billing close for failed region delayed until spool drains |
| Testing | Region-evacuation game day: drain one region, verify counts against edge logs |
| Rollout | Second region in shadow (mint and log, but serve from primary) for 2 weeks |
Metrics to Watch:
deny_set.hits_total (cross-region retries), region.spool_unshipped_bytes, ledger.region_close_delay_hours, edge_to_ledger.gap_pct per region.
Organizational Follow-up: Each region gets an owning on-call; the global ledger has a single owner. Region evacuation added to the quarterly game-day calendar.
Ownership Question: "If a region is lost and 20 seconds of clicks are gone, who decides whether to bill estimated clicks?" Staff answer: Finance — and the default should be no. We don't bill clicks we can't prove. Engineering provides the estimate for revenue reporting only.
Key Takeaway: "Mint in one region, dedup in that region. Global coordination only for the rare cross-region retry."
What clears the Staff bar:
- Avoids global dedup by construction
- Quantifies RPO and ties it to the spool design
- Distinguishes edge availability from pipeline availability
9. Level Expectations Summary#
After studying this case study, you should be able to:
- Split "the click count" into billing-grade, dashboard and pacing products, each with a numeric correctness and freshness envelope
- Walk exactly-once end to end — edge spool, idempotent producer, keyed dedup, checkpointed windows, idempotent upsert sink — and name where it breaks
- Set watermark and allowed-lateness per consumer, and explain why billing closes in batch
- Design two-stage salted aggregation that keeps hot-ad counts exact
- Describe kappa-plus-reconciliation and why shared logic needs an independent invariant
- Design fraud reclassification as versioned corrections with a credit policy
- Choose an OLAP store from the query shape and price rollup dimensions
- Run a backfill or restatement with approval gates
The Bar for This Question#
Mid-level (L4): Builds a working pipeline: click service → queue → aggregator → database → dashboard. Knows to use a message queue and windowed aggregation. Correctness under restart, duplicates and late data is not addressed until prompted.
Senior (L5): Builds a correct-looking streaming pipeline with Kafka, Flink, event-time windows and an OLAP store. Mentions dedup and exactly-once. Tends to treat the stream as the billing number, uses one lateness setting, and has no correction or reconciliation story. Fraud is a filter.
Staff+ (L6): Starts from the consumers and their correctness bars. Treats exactly-once as a path property with an honest edge boundary. Separates the provisional stream from the batch-closed ledger, designs reconciliation with drift metrics, versions corrections and fraud reclassifications, and names owners and approval gates for anything that moves revenue. Knows the hot-ad cascade and the stalled-watermark failure. The interviewer should learn something from the answer.
10. Staff Insiders: Controversial Opinions#
10.1 "The Stream Should Never Be the Billing Number"#
| Evidence | Detail |
|---|---|
| Restarts, replays, late data | Every one of them moves a streaming number after display |
| Fraud classification | Sophisticated IVT detection needs 24–72h of context |
| Audit | Finance needs a number reproducible from raw logs with a frozen logic version |
The Staff position: Streaming previews; batch closes. Even with flawless exactly-once, billing needs a close event, a frozen logic version and an independent recount. Uber's Flink pipeline has end-to-end exactly-once and still compares against offline data.
Why this matters in interviews: Candidates who bill from the stream will be walked through restart, late data and fraud until the design collapses. Saying this early ends that line of questioning.
10.2 "Lambda vs Kappa Is the Wrong Debate"#
| Evidence | Detail |
|---|---|
| Lambda's real cost | Two implementations of business logic that drift |
| Kappa's real cost | No independent check; replay bounded by log retention |
| What actually matters | Whether there is an independent recount and a correction protocol |
The Staff position: One logic library, run in streaming and batch modes, plus an invariant that doesn't depend on that library. The architecture label doesn't matter; the audit does.
Why this matters in interviews: Interviewers who ask "lambda or kappa?" are checking whether you know why lambda existed. Answering "I want an independent recount — here's how I get it without two codebases" is the strongest answer.
10.3 "Exactly-Once Is Cheaper at the Sink Than in the Engine"#
| Evidence | Detail |
|---|---|
| Two-phase commit sinks | Add checkpoint-interval latency (30–60s) to visibility |
| Idempotent upserts | Zero added latency; replay just overwrites |
| Real failures | Most double-billing comes from append sinks, not engine bugs |
The Staff position: Design the output as an idempotent function of (key, window) and you rarely need transactional sinks. Spend the engineering effort on the key design.
Why this matters in interviews: It shows you understand why exactly-once works, not just that a flag enables it.
10.4 "Most Click-Count Disputes Are Definition Disputes, Not Bugs"#
| Evidence | Detail |
|---|---|
| Clicks vs landing-page arrivals | Users bounce before the landing page loads; 5–10% gaps are normal |
| IVT methodology | Our filter removes clicks the advertiser's tracker counts |
| Time zones | Advertiser reports in local time; ledger in UTC |
The Staff position: Build a discrepancy report and a published measurement definition before investing in more accuracy. Most tickets close with an explanation, not a fix.
Why this matters in interviews: It shows you know where the business pain actually comes from.
11. The Principal Lens (L7)#
Why L7 Sees This Problem Differently#
A Staff engineer designs a click aggregator. A Principal engineer notices that search ads, display ads, video ads and the new retail-media product each count clicks in their own pipeline, each with its own dedup, its own IVT rules and its own definition of a billable click — and that finance reconciles four ledgers by spreadsheet every month. The L7 question is not "Flink or Spark?" but "how many ways does this company count money, who certifies each, and which do we consolidate?" At org scale, click aggregation is a measurement platform with a governance layer: one event schema, one dedup service, one IVT methodology with versioning, one ledger contract with finance, and many surfaces plugged into it.
🧭 Principal Move: "Before designing another pipeline, I'd like to know how many systems compute billable events today and how finance reconciles them. If the answer is 'four, by spreadsheet', the highest-leverage design is a shared measurement platform with one ledger contract — and the new surface becomes its first customer."
The Org-Level Fault Line#
One measurement platform vs per-surface counting pipelines.
| Option | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Per-surface pipelines | Each team ships fast; fits its surface | 3–5 definitions of a billable click; IVT rules diverge; finance reconciles manually; cross-surface fraud invisible | Finance (manual recon), advertisers (inconsistent numbers), ads quality (N integrations) |
| One central pipeline, central team | One definition, one ledger | Central team becomes a bottleneck; new surfaces wait quarters | Product velocity; platform burnout |
| Shared measurement platform + surface-owned logic plugins (L7 default) | Common edge, log, dedup, IVT, ledger contract; surfaces own attribution rules within a schema | 3–6 engineer-quarters upfront; requires a versioned plugin contract | Platform budget up front; pays back at the third surface |
The deciding question: how many billable surfaces will exist in 2 years? One or two: shared library and conventions. Three or more: the platform pays for itself in finance hours and avoided restatements.
Cost Model#
Assumptions: cloud list prices, ~300 bytes per raw click, Parquet compression ~5×, managed Kafka, a self-managed or managed OLAP cluster, loaded engineer cost ~$25K/month.
| Scale | Architecture | Infra $/month | Headcount | On-Call Load |
|---|---|---|---|---|
| ~10M clicks/day, 1 region | Edge + Kafka (3 brokers) + small stream job + ClickHouse or Pinot (3 nodes) + nightly batch | ~$5K–$12K | ~1–1.5 FTE | Shared with data platform; ~1 page/month |
| ~1B clicks/day, 1–2 regions | Spooling edge, Kafka (12–20 brokers), stream job with two-stage aggregation, Pinot/Druid (20–40 nodes), lake + nightly recount, reconciler | ~$80K–$200K | 5–8 FTE (edge, pipeline, OLAP, recon/billing) | Dedicated rotation; ~3–6 pages/month |
| ~10B events/day incl. impressions, 3+ regions | Regional cells, shared dedup and IVT services, impressions sampled except CPM, federated OLAP, global ledger | ~$600K–$1.5M | 15–25 FTE across platform + per-surface | Per-region rotations; quarterly region-evacuation and replay game days |
The line that matters to leadership: at 1B clicks/day and a ~$0.50–$1 average CPC, the platform counts ~$15–30B a year. A $150K/month platform is ~0.01% of what it measures. The expensive failures are not infrastructure — a single 1% overbilling event costs more than a year of the platform.
The 3-Year Evolution Path#
Each step is triggered by a business event. Building regional cells before there's a second region or a residency requirement is how a platform team spends a year on something nobody needs yet.
One-Way Doors vs Two-Way Doors#
| Decision | Door | Reversal Cost | Why |
|---|---|---|---|
| Stream engine (Flink vs Kafka Streams vs Spark) | Two-way | 1–2 quarters | Behind a logic library and a sink contract |
| OLAP store | Two-way if the API abstracts it | 1–2 quarters | Becomes one-way if advertisers query it directly |
click_id format and minting location | One-way | Years | Dedup, fraud, attribution and partner exports all key on it |
| Definition of a billable click | One-way | Contract changes | Advertiser terms of service and revenue recognition depend on it |
| Billing close policy (T+3d) | One-way-ish | Finance process + customer notice | Revenue recognition and invoicing calendars build on it |
| Raw log retention period | One-way in the short direction | Lost data is gone | Shortening it destroys the ability to restate or defend disputes |
| Rollup dimensions | Two-way | Backfill from raw | Cheap to add if raw retention allows |
🧭 Principal Insight: Spend review time proportional to reversal cost. The OLAP debate gets an afternoon. The
click_idformat, the billable-click definition and raw retention get a design review with finance, legal and ads quality in the room.
The Standard I'd Write#
RFC: Billable Event Measurement Standard (v1)
Scope: Every system that counts events used for invoicing, revenue recognition, or advertiser-facing reporting.
Requirements:
- Billable events MUST carry a server-minted, signed, globally unique event ID assigned at the first server-side touchpoint.
- Raw events MUST be written to the immutable measurement log and retained for ≥13 months.
- Every ledger number MUST be reproducible from the raw log plus a recorded logic version.
- Sinks feeding reports or ledgers MUST be idempotent; append-only increments are not permitted.
- Pipelines MUST publish
recon.drift_pctagainst an independent count (edge or batch) and page at >0.1% on billed days.- Changes to IVT rules or attribution logic projected to move revenue >0.5% MUST run in shadow ≥7 days and be approved by finance.
- Advertiser-facing reports SHOULD expose
is_finalanddata_freshness_ts.Exceptions: Filed with the measurement architecture group; decision within 5 business days; maximum duration 2 quarters.
Success metrics: All billable surfaces on the platform within 4 quarters; zero restatements caused by double-counting; finance monthly reconciliation time reduced from days to hours; advertiser discrepancy tickets down 50%.
What I'd Tell the VP#
"We count roughly $20 billion a year of clicks across four separate systems, and finance reconciles them by hand every month. Each system has its own idea of what a valid click is, which is why advertisers sometimes see different numbers for the same campaign. I'm proposing we move all of them onto one measurement platform over the next year, with a fast preview number for dashboards and an audited number for invoices. It costs about six engineers for a year. The return is fewer refunds, a month-end close that takes hours instead of days, and a fraud system that can see across all our ad products. The biggest risk is migration — we'll run old and new side by side and only switch billing when they agree for 30 days."
Principal Interview Signals#
| Signal | What It Sounds Like |
|---|---|
| Inventories before designing | "How many systems count billable events today, and who certifies each?" |
| Prices the guarantee | "Exactly-once on every derived metric costs ~30% more compute. I'd buy it for billing and pacing only." |
| Treats definitions as contracts | "The billable-click definition is a one-way door — it's in the terms of service." |
| Designs correction governance | "Engineering computes restatements; finance approves them. The threshold for automatic credits is a finance decision." |
| Designs for correlated failure | "A shared library means stream and batch agree on bugs. I want an invariant that doesn't use the library." |
Staff answers that L7 interviewers find insufficient:
- "We'll reconcile stream against batch nightly." — correct, but silent on who acts on the diff, what threshold triggers a restatement, and who approves it.
- "We'll use two-stage aggregation for hot ads." — right mechanism, but no capacity process for tentpole events that the whole org can plan around.
- "Pinot for serving." — sound, but never asks whether the org already runs an OLAP engine, or what a second one costs to operate for five years.
Appendices
Appendix A: Windowing, Watermarks and Late Data
A.1 Window Lifecycle#
A.2 Choosing the Numbers#
| Parameter | Dashboard Path | Billing Path | Reasoning |
|---|---|---|---|
| Window | 1 min tumbling | 1 day (UTC) | Dashboards need minute grain; invoices need days |
| Watermark delay | 30s | n/a (batch) | Captures ~99% of clicks at first emission |
| Allowed lateness | 5 min | 72h | Beyond 5 min, corrections are cheaper than keeping state |
| Late beyond lateness | Side output → correction job | Included if before close | Nothing is silently dropped |
| After close | Adjustment on next invoice | Adjustment | Policy, not mechanism |
A.3 Pseudocode — Stream Side#
clicks
.assignTimestamps(e -> e.ts_event, maxOutOfOrder = 30s)
.withIdleness(60s) // idle partitions don't stall watermark
.keyBy(e -> e.click_id)
.filter(dedupWithTtl(24h)) // RocksDB seen-set; dups to side output
.map(applyIvtRules(rulesAt(e.ts_event))) // effective-dated rule version
.keyBy(e -> (e.ad_id, hash(e.click_id) % 16))
.window(Tumbling(1 min)).allowedLateness(5 min)
.aggregate(sumClicksAndSpend) // stage 1 partials
.keyBy(p -> p.ad_id)
.window(Tumbling(1 min))
.reduce(sumPartials) // stage 2 exact totals
.sinkTo(upsert(key = (ad_id, window_start, dims), value = full_totals))
lateSideOutput -> correctionJob -> upsert(version = n+1)
Appendix B: Click ID, Edge Path and Dedup
B.1 The Redirect Path#
B.2 click_id Construction#
| Field | Bits | Purpose |
|---|---|---|
| Region | 4 | Regional dedup; residency |
| Timestamp (ms) | 41 | Ordering; TTL eviction; event-time sanity check |
| Edge host ID | 12 | Debugging; per-host gap detection |
| Random | 31 | Uniqueness within host-millisecond |
| HMAC (separate field) | 64 | Forged-ID rejection |
Never trust a client-supplied click ID for dedup: a bot can generate unique IDs per click to evade dedup, or reuse a valid one to inflate a competitor's spend. See ID Generation.
B.3 Dedup State Sizing#
| Clicks/day | TTL | Raw key bytes | With RocksDB overhead | Per task (64 tasks) |
|---|---|---|---|---|
| 100M | 24h | 1.6 GB | ~3–5 GB | ~50–80 MB |
| 1B | 24h | 16 GB | ~30–50 GB | ~0.5–0.8 GB |
| 10B | 24h | 160 GB | ~300–500 GB | ~5–8 GB (use 256+ tasks) |
Why 24h: client retries cluster within seconds; spool replays within hours. Beyond 24h, the remaining duplicates are caught by the batch recount, which dedups over the full day. See Idempotency for the general pattern.
Appendix C: Reconciliation, Corrections and Backfill
C.1 The Correction Flow#
C.2 Backfill Protocol#
| Step | Detail | Gate |
|---|---|---|
| 1. Scope | Advertisers, days, logic version; estimated delta | Reporting lead |
| 2. Dry run | Run in a scratch namespace; produce per-advertiser delta report | Automatic |
| 3. Approve | Delta >0.5% of any advertiser's period spend or any closed day | Finance |
| 4. Publish | Write new versions atomically per day; OLAP segment replace / upsert | Automatic |
| 5. Notify | Account managers for advertisers above threshold | Billing ops |
Backfills run on a separate compute pool with a rate cap so they never compete with the live stream. A 30-day backfill at 1B clicks/day reads ~9 TB of compressed Parquet; at a few hundred cores that's 1–3 hours.
Appendix D: Query Serving and the OLAP Store
D.1 Rollup Design#
| Table | Grain | Dimensions | Retention | Est. Rows/Day (1B clicks) |
|---|---|---|---|---|
clicks_minute | 1 min | ad, campaign, advertiser, country, device | 7 days | ~200–500M |
clicks_hour | 1 hour | same + placement | 13 months | ~20–60M |
clicks_day_ledger | 1 day | advertiser, campaign, billable flag | 7 years | ~1–5M |
| Raw (lake) | event | all, including high-cardinality | 13 months | 1B |
D.2 OLAP Options#
| Store | Strength Here | Watch Out For |
|---|---|---|
| Apache Pinot | Kafka ingestion, upsert tables, low-latency user-facing queries | Upsert tables require partitioning by primary key |
| Apache Druid | Ingest-time rollup, interval-based segment replacement for backfill | Real-time segments can't be updated in place; corrections via reindex |
| ClickHouse | Very fast scans, SQL, ReplacingMergeTree | Dedup is eventual until merges; use FINAL or argMax |
| Time-series DB | Simple metrics dashboards | Weak multi-dimensional group-by; see Time-Series Databases |
| Cloud warehouse | Ad-hoc over raw | Per-query cost and seconds of latency for user-facing dashboards |
D.3 API Contract#
GET /v1/reports?advertiser=A&from=2026-09-01&to=2026-09-30&granularity=day&dims=campaign
200 OK
{
"rows": [...],
"data_freshness_ts": "2026-10-01T09:14:00Z",
"final_through": "2026-09-27", // days closed for billing
"is_final": false,
"logic_version": "attrib-v42/ivt-v17"
}
Advertisers query a cache in front of OLAP for repeated dashboard loads (60s TTL); export jobs use a separate rate-limited path. See Rate Limiting.
Appendix E: Observability
E.1 Core Metrics#
# Completeness
redirect.302_total # independent edge count
redirect.events_dropped_total # must be zero; page on any
edge_to_log.gap_pct # 302s vs produced, per minute
# Correctness
dedup.duplicates_total # expected ~0.5-2%
recon.drift_pct{advertiser,day} # stream vs batch
invariant.billable_gt_redirects_total # must be zero
# Freshness
watermark.lag_seconds
olap.ingest_lag_seconds
events.late_total{bucket}
# Fraud
ivt.rate{rule_version,surface}
ivt.reclassified_after_close_total
E.2 Critical Alerts#
| Alert | Condition | Severity |
|---|---|---|
| Edge drop | redirect.events_dropped_total > 0 for 1 min | Page |
| Edge-to-log gap | > 0.1% for 5 min | Page |
| Watermark stalled | watermark.lag_seconds > 120 | Page |
| Recon drift on billed day | > 0.1% | Page billing on-call |
| Invariant violation | billable > redirects for any ad-hour | Page + freeze billing close |
| IVT rate shift | ±30% vs 7-day baseline | Ticket to ads quality |
E.3 Debugging the Silent Undercount#
- Compare edge count to ledger for the period. If they differ and stream ≈ batch, the loss is before the log.
- Break down
edge_to_log.gap_pctby pod and release version. - Check dedup: sample dropped duplicates — are any distinct clicks with colliding IDs?
- Check IVT: did a rule version change on the day the gap started?
- Check late data: is a growing share arriving after billing close?
Appendix F: Scale Evolution
F.1 What Works at Each Scale#
| Scale | Design | What Breaks Next |
|---|---|---|
| < 10M clicks/day | Redirect logs to a warehouse; hourly SQL; dashboards on the warehouse | Freshness and per-query cost |
| 10M–500M/day | Kafka + stream job + one OLAP store + nightly batch | Hot ads; correction tooling |
| 500M–5B/day | Two-stage aggregation, spool edge, reconciler, versioned corrections | Multiple surfaces; finance recon |
| > 5B/day incl. impressions | Measurement platform, regional cells, sampled impressions | Governance, residency |
F.2 What You Don't Build on Day One#
- Multi-region active-active pipelines (regional edge is enough)
- ML-based post-hoc fraud (start with rules; add ML when IVT rate is measured)
- Exact impression counts for non-CPM placements (sample)
- Self-serve backfill UI (a runbook and an approval ticket suffice until backfills are weekly)
- Per-advertiser custom dimensions (price them first)
Appendix G: Multi-Tenancy, Fairness and Cost
G.1 Noisy Advertisers on the Query Side#
A single agency running hourly exports over 13 months for 10K campaigns can consume a large share of OLAP capacity. Per-advertiser query quotas (e.g., 50 queries/min interactive, exports through an async job queue) protect everyone else's dashboards.
G.2 Chargeback for Dimensions#
| Dimension Request | Cardinality | Rollup Multiplier | Est. Monthly Cost |
|---|---|---|---|
| Country | ~200 (≈20 active per ad) | ~5–10× active rows | Baseline |
| Device type | 3–5 | ~2–3× | +$3–6K |
| Placement | ~1K | ~5× | +$10–20K |
| Zip code | ~40K | Approaches raw | Lake only |