Learning pathsA
GUIDED PRACTICE

Review: Ad Click Aggregator

Count ad clicks per campaign and time window while accepting retries, late arrivals, and reprocessing.

Interview scope and guarantees

Ingest click events, deduplicate transport retries, aggregate by campaign and time, and expose provisional and finalized totals. Fraud classification is a separate versioned process. Define whether billing uses these totals or a separately reconciled ledger.

Capacity worksheet

Assume 2 billion events/day: 23,148/s average and 115,741/s at 5× peak. At 200 bytes/event that is 400 GB/day raw. A 30-day replay log is 12 TB before replicas. Partitioning only by campaign can overload the shard for one major advertiser.

Concrete API contract

Contract / pseudocode
POST /clicks {eventId,campaignId,adId,eventTime,impressionId}
GET /campaigns/{id}/counts?from=...&to=...
Return count, asOf, watermark, status and aggregationVersion.

Data model and access paths

Contract / pseudocode
raw_clicks(event_id,campaign_id,event_time,ingest_time,payload_version)
partial_counts(campaign_id,bucket,salt,count,checkpoint)
final_counts(campaign_id,bucket,version,count,status)
dedup(event_id,expires_at)

Evolve a solution and explain each change

Three architecture decisions for Ad Click Aggregator, including the pressure each introduces.
Scroll to inspect the diagram, or open it at full size.

Figure — Three architecture decisions for Ad Click Aggregator, including the pressure each introduces.

Step 1: Record raw clicks

Accept validated event identities into a durable log. Synchronous counter writes are vulnerable to retries and hot campaigns.

Step 2: Aggregate event windows

Key state by campaign/window and checkpoint progress with output semantics. Late events and reprocessing can double count materialized totals.

Step 3: Separate estimates from billing

Expose provisional dashboards; finalize billable records with corrections and auditability. A displayed count must not silently become an invoice.

Responsibility overview

Connected responsibilities for Ad Click Aggregator. Trace the authoritative and derived paths separately.
Scroll to inspect the diagram, or open it at full size.

Figure — Connected responsibilities for Ad Click Aggregator. Trace the authoritative and derived paths separately.

Worked end-to-end scenario

Click e17 arrives twice after a client retry. Ingestion validates the ad/campaign association and retains enough identity to deduplicate within policy. An aggregator updates the ten-minute campaign window. It crashes after writing a total but before checkpointing its input offset. On replay, a naive increment counts e17 again. Either coordinate checkpoint and output through the selected processing model or write deterministic versioned aggregates that replace the same window result idempotently. A click arriving after the watermark revises the provisional dashboard or becomes a correction under the billing contract; it does not disappear without explanation.

Why these access paths matter

Keep immutable raw events partitioned for replay and aggregate rows by campaign, window, and revision. Event time and ingestion time answer different questions. Dedup retention must cover the retry/replay horizon. Hot campaigns may need partial aggregates merged into an exact total; the merge key must include the input range or revision to avoid adding the same partial twice.

Build the baseline first

Ad Click Aggregator: baseline request paths.
Scroll to inspect the diagram, or open it at full size.
  1. Validated event → durable log → partition consumer.
  2. Event-ID dedup → time bucket → counter update.
  3. Bucket totals → query store → dashboard.

Evolve the design under load

Ad Click Aggregator: additional scaling and recovery paths.
Scroll to inspect the diagram, or open it at full size.
  1. Hot campaign → salted partial aggregates.
  2. Partial merge → campaign totals → versioned publication.
  3. Raw archive → bounded replay → replacement generation.

Defend the hardest decision

Use event time for business windows and ingestion time to diagnose delay. A watermark expresses progress, not proof that no late event can ever arrive. Define an allowed lateness period, then publish corrections or quarantine later events. If an event ID is retained for only one day but the log can be replayed for thirty, a replay can double count unless the sink uses a generation or source-position transaction.

Failure and recovery analysis

A processor commits an aggregate and crashes before advancing its offset. Replay must not increment the same contribution again. Options include transactional state/checkpoints or storing idempotent per-batch deltas keyed by source ranges. Reprocessing fraud rules should create a new version rather than silently rewriting billing history.

Security and privacy boundary

Validate campaign ownership for queries, rate-limit ingestion abuse, and minimize personal tracking fields.

Interview follow-ups with reasoning

Question: How do you handle campaign skew?

Show answer and explanation

Answer: Salt partial keys, then merge bounded subaggregates.

Question: How do you validate correctness?

Show answer and explanation

Answer: Compare exact raw-event totals for sampled buckets with the published generation.

Question: How do you recover from a bad deployment?

Show answer and explanation

Answer: Replay into a parallel generation and switch queries after reconciliation.

Operate and verify the design

Source lag, dedup hits, late-event rate and aggregate-to-raw discrepancies.

Crash after aggregate commit but before offset advancement and confirm no double increment.

A second scenario to test transfer

Three events arrive at 10:04: a new click that occurred at 10:03, a retry of an already counted 10:03 click, and a delayed click from 10:02. The first increments the 10:03 window. The second changes nothing. The third updates 10:02 according to the late-arrival policy. Simply incrementing the current minute for all three would get both identity and time wrong.

Walk through a duplicate arriving after the normal dedupe retention window. What can the service honestly guarantee?

Show answer and explanation

Answer: The short-lived dedupe state alone cannot recognize it. Either retain durable identity evidence for the billing/replay horizon, restrict the accepted replay contract, or reconcile against an authoritative ledger. State the limitation rather than promising unbounded exactly-once counting.

Compare alternatives

ChoiceAdvantageConsequence
Event-time windowsReflect actual event timeLate-data policy required
Ingestion-time windowsSimple operational timingDelayed events shift periods
Provisional totalsFast feedbackVisible corrections
Billing finalizationStronger audit boundaryMore latency and reconciliation

A design-changing exercise

A processor advertises exactly-once state updates. Are external invoice writes automatically exactly once?

Show answer and explanation

Answer: No. The guarantee must include the sink boundary. Use sink transactions, idempotent records, or reconciliation with a durable billing ledger.

Design workshop: checkpoint one aggregate revision

Core scope is eligible click ingestion, minute/hour reporting and reconciled billable totals. Fraud scoring is upstream. Choose provisional dashboard refresh every ten seconds, two-minute event-time lateness and durable raw history for a 30-day reconciliation horizon. A billing invoice is a separately finalized object, not the latest dashboard number.

Choose a transactional batch sink: each partition processes a contiguous source range, writes event dedup identities and aggregate deltas, and advances its durable offset in the same database transaction. On restart, seek from that database offset. Do not both commit an unrelated broker offset early and claim the database output cannot be lost. If a stream engine owns the checkpoint, its external-sink integration needs an equivalent recoverable contract.

sql
BEGIN;
SELECT next_offset FROM source_progress WHERE partition=:p FOR UPDATE;
-- Require the batch to start at the recorded offset.
-- Insert unique event IDs; apply only newly inserted eligible deltas.
-- Update campaign/window counters and their revision IDs.
UPDATE source_progress SET next_offset=:end_exclusive WHERE partition=:p;
COMMIT;

This transaction belongs to one owned input/sink shard, not one global database. Partition source ranges so their state keys and progress record co-reside, and scale sink shards independently. A range spanning independent sinks needs a different checkpoint protocol; it cannot reuse this single-transaction proof.

Dedup rows retain identities for the accepted replay horizon. If operators replay older history after short dedup retention, rebuild a new projection from authoritative events rather than append again to the live counters. Stable batch-range uniqueness prevents duplicate batch commits but does not by itself remove producer-created duplicates carrying different offsets with the same event ID.

Output and offset survive the same crash boundary. Log partition to Aggregator: Batch offsets 100..119; Aggregator to Transactional sink: Dedup events, apply deltas, set next_offset=120; Transactional sink to Aggregator: Commit succeeds but acknowledgement lost; Aggregator to Transactional sink: Read next_offset after restart: 120; Aggregator to Log partition: Resume at 120, not increment 100..119 again; Dashboard to Transactional sink: Read campaign/window revision and provisional flag
Scroll to inspect the diagram, or open it at full size.

Figure — Output and offset survive the same crash boundary.

Salt a hot campaign into S subkeys by event identity. If a minute has salt totals [400,300,200,100], merge to 1,000 before campaign reporting. Checkpoints must identify complete salt revisions; adding an old partial total to a new one produces an invalid window. Negative fraud corrections reference the original eligible click and are applied once, yielding a new aggregate revision.

Event time, lateness and correction arithmetic. 10:03 click arrives 10:04 / 10:03 +1 / Provisional revision increases; Same event ID replayed / No change / Duplicate counted as processing metric only; 10:02 click arrives before watermark closure / 10:02 +1 / Late but accepted correction; Fraud reversal of one counted click / Original window -1 once / Revised total with audit identity; Invoice finalized at 1,000 then valid correction -1 / Adjustment of -1 / Do not silently rewrite an issued invoice
Scroll to inspect the diagram, or open it at full size.

Figure — Event time, lateness and correction arithmetic.

Use the minimum of active input watermarks to close a window; an idle partition needs explicit idle detection and re-entry late-data policy. A stalled partition cannot silently be ignored without changing completeness. Store rollups by campaign and window for advertiser reads, with currency/rate version only at billing boundaries. Reconcile: accepted eligible IDs − reversed IDs = finalized billable count; compare independently derived raw-history totals against projection totals.

For one million events/s at 100 bytes, raw ingestion is 100 MB/s and 8.64 TB/day before compression/replication. State size depends on active campaigns, dedup horizon and buckets, not only event throughput. Alert on offset lag, watermark lag, duplicate fraction, correction age and reconciliation delta; keep dashboard freshness separate from billing finality.

Exercise: Batch commit succeeds but the worker dies before a broker offset commit. Will replay double bill?

Show answer and explanation

Answer: This selected design resumes from the sink's transactional next_offset and deduplicates stable event IDs. If another consumer incorrectly follows only a stale broker offset, the source-progress guard rejects the already committed range. Invoice creation has its own unique finalized billing identity.

Technical references

Kafka processing and external-sink boundaries.

PostgreSQL transaction isolation and concurrent updates.

11:00Self-guided practice timer
The timer resets when you leave this page. Save your design separately.
Your challenge

Walk through a duplicate arriving after the normal dedupe retention window. What can the service honestly guarantee?

Your design draft

Clarify assumptions, explain your approach, and test the difficult cases. Save your draft, then compare it with the study notes.

Read study notes

Self-review checklist

Self-guided practice. Automated AI feedback and code execution are not connected.