System-design interview · Extended interviews
Design event-time click analytics
Design click reports that count each accepted event in its occurrence-time window, recover without duplicate contributions and publish complete totals with traceable historical corrections.
You will learn to
- Assign events to windows independently of arrival time.
- Explain deduplication, watermarks, revisions, and checkpoint recovery.
- Distinguish complete-key top-k from unsafe merging of partial winners.
Practice in this chapter
8 interview questions with model answers and follow-ups.
Go to interview practiceUseful foundations: Message queues, event logs, delivery guarantees, and backpressure · Data partitioning and sharding · Database indexes: B-trees, composite keys and query access · Replication and durability
Workload and timing examples are interview assumptions.
Dotted concept links open the relevant explanation in a new tab.
01Problem and scope
A click-analytics platform validates incoming events, identifies retries, and counts accepted clicks in time windows to produce reports and rankings. Received events, valid clicks, distinct users and billable interactions require different identities and rules. This design counts validated clicks by occurrence time. For an example, C901 occurs at 09:00:58 but arrives at 09:02:05, revising A7 in the 09:00 minute from 99 to 100 under the declared lateness policy.
Start with one process reading a durable list of events. For an occurrence-time report, compute a window from each validated event timestamp and increment that ad/window counter. Record event IDs so a transport retry cannot increment twice. This small model exposes the two essential problems: time and repeated delivery.
I ask whether the dashboard measures received events, valid clicks, unique people, or billable interactions. We choose validated click events grouped by occurrence time. Distinct users and billing have different identities and rules and remain separate. An event ID prevents transport repetition; it does not prove the user intended a legitimate click.
Occurrence time, also called event time, determines which reporting window owns a click; arrival time determines when the service learns about it. A watermark is the processor’s declared progress through event time. It lets the system apply a defined window-closing policy despite delayed delivery, while still allowing an explicit correction path for older arrivals.
The interviewer asks for “real time.” I translate that into a normal five-second accepted-to-dashboard target, two minutes of ordinary event-time lateness, and visible preliminary/final labels. The dashboard client can understand why the count changed from 99 to 100 because the result carries a revision, watermark, and validation-policy version.
Raw accepted events remain available for a declared correction window. This preserves the evidence needed to explain a late update or a policy change rather than making the dashboard's current integer the only surviving truth.
02Functional requirements
- Collect a click. Validate envelope and durably accept a stable event ID.
- Show minute counts. Return count, policy, revision, data boundary, and finality.
- Update for ordinary lateness. Correct the original event-time window.
- Read hourly top 100. Rank complete per-ad totals under one published boundary.
- Inspect a correction. Explain source interval, policy, and superseded version.
- Rebuild an interval. Publish a new authoritative result version without mixing live output.
Scope and acceptance boundaries
Report one-minute valid-click counts, campaign rollups, and hourly top 100 ads. Show preliminary updates within seconds, allow two minutes of ordinary lateness, and retain raw events for a declared correction/audit period. Fraud-model training, attribution, distinct-user estimation, and financial settlement are separate scopes.
A repeated delivery of C901 is not a new click. Two genuine clicks may have different IDs; whether both are billable is a business validation rule tied to impression/ad provenance. Label output as preliminary or finalized under a policy version. A live dashboard should not silently become the sole billing authority.
Support campaign totals and selected country/device breakdowns, with limits on allowed values. Arbitrary user-defined dimensions need a separate cost and privacy review: every additional combination may need its own counter.
Final means complete under the declared watermark and lateness policy, not proof that no older event can ever arrive. Beyond-policy events enter a correction stream. The dashboard client can see a later revised historical result, but the earlier published report remains reproducible by its version. A repeated query pinned to a version does not silently change beneath pagination.
03Non-functional requirements
- Workload. Assume one billion events/day and a 200,000/s peak.
- Collection latency. Target durable collection p95 below 100 ms.
- Visibility and query latency. Target normal accepted-to-dashboard visibility within five seconds and narrow campaign-query p95 below 200 ms.
- Availability and durability. Target 99.9% ingestion availability under admitted load. Accepted input survives one zone failure through replicated logs.
- Retention. Keep raw validated envelopes for 30 days, online event-ID deduplication for 24 hours, and compact aggregates longer under a declared policy.
- Ordinary lateness. Allow two minutes after a window ends, measured against watermark progress. Bound timestamps through provenance and clock-skew policy so forged year-old/far-future times cannot control event-time progress.
These figures are exercise assumptions.
Processing invariants
- One contribution. Count an accepted event at most once while its identity is covered by the supported online retry policy.
- Recoverable visibility. Save the processing state and input positions behind each published result so recovery can reproduce it.
- Correction authority. An older live writer cannot overwrite a newer authoritative correction.
- Honest overload. Preserve acknowledged input and expose lag; neither silent drops nor falsely final partial windows are acceptable.
Old data needs a rebuild contract
A 24-hour deduplication window cannot make an arbitrary year-old resend safe. Backfills reconstruct an interval from immutable source under a separate job identity, instead of blindly resending history into the live counter.
04Capacity estimates
| Quantity | Calculation | Consequence |
|---|---|---|
| Average events | 1B/day / 86,400 ≈ 11,574/s | Separate peak provisioning |
| Raw payload | 1B × 200 B = 200 GB/day | Archive/replay cost |
| Peak payload | 200K/s × 200 B = 40 MB/s | Partition collector/processor load |
| Hour of minute counts | 1M active ads × 60 × 24 B = 1.44 GB | Before state/index/checkpoint overhead |
| One-day dedupe payload | 1B IDs × 32 B = 32 GB | Dedupe can exceed aggregate state |
Each distinct country/device/campaign combination may create another counter key. Choose a supported retry/dedupe horizon deliberately; a 24-hour membership cache cannot make an arbitrary year-old replay duplicate-free. Offline corrections can instead rebuild a versioned interval from immutable input.
Thirty days of 200 GB/day raw payload is 6 TB before replicas, envelopes, and indexes. Three retained copies would make 18 TB of payload. At peak 40 MB/s, a ten-minute processing outage accumulates 24 GB before overhead. A processor fleet that resumes at 300,000 events/s while 200,000/s continue arriving drains that 120-million-event backlog in 1,200 seconds, or 20 minutes.
A checkpoint is a durable recovery snapshot that pairs processing state with the input positions that produced it. For these counters, that includes both accumulated counts and remembered event identities. Its frequency controls how much work a restart must replay and, in this design, how often a complete result can become visible.
With checkpoints every two seconds, a peak interval contains 400,000 input events. That does not imply copying the full 32 GB deduplication set every two seconds: incremental state snapshots and immutable shared files can reduce write volume, at the cost of managing their lifecycle correctly. Checkpoint duration must remain below the useful publication cadence or five-second freshness becomes impossible.
A hot ad receiving 20% of peak traffic sees 40,000 events/s. Hashing only ad ID pins that work to one owner. Salting—adding a shard suffix to split one ad's counter across 16 partial keys—lowers the average hot-ad write rate per partial to 2,500/s, but requires another stage to combine the partials and track their versions. The count can only be ranked after those partial contributions are combined.
The 24-byte aggregate and 32-byte identity estimates are payload assumptions, not actual engine memory measurements. Hash-table overhead, indexes, checkpoint metadata, and retained generations can dominate small counters.
05APIs and contracts
Collect a click
POST /clicks accepts {eventId:"C901",adId:"A7",impressionId:"I88",occurredAt:"2026-09-22T09:00:58Z",collectorVersion:2} with an authorized collector credential and provenance token. It returns an accepted identity after durable log append, not a guarantee that the click is billable or already visible. Invalid provenance, impossible timestamps, excessive payloads, and exhausted admission quotas receive explicit errors.
Times such as 09:00 in the worked trace abbreviate this UTC date. Actual requests use full timestamps with an offset so different dates or time zones cannot collapse into the same reporting window.
Retry identity and conflicts
A timeout leaves acceptance uncertain; the sender retries C901 rather than inventing a new ID. The online identity is (tenant, collector source, eventId). Same-identity conflicting immutable fields are rejected or quarantined, not treated as a new valid click. The authenticated source identity comes from the credential, not a caller-supplied tenant field. Stable IDs must be generated or validated by the trusted collector or impression service; letting a malicious client choose unlimited fresh identities defeats transport deduplication as an abuse control.
Read counts and top ads
GET /campaigns/C7/counts?from=2026-09-22T09:00:00Z&to=2026-09-22T10:00:00Z&resolution=minute returns each count with policyVersion, resultGeneration, revision, watermark, and preliminary/final status. Top 100 queries use the same published generation across contributing partitions. Cursors include that generation and a deterministic count/ad-ID order.
Publish a privileged correction
A correction request is privileged and specifies the input interval, source manifest, validation policy, reason, and desired output namespace. It never masquerades as an ordinary live event batch. Reports can pin either the latest authoritative interval version or a historical version for reproducibility.
06Data model and access patterns
| Record | Identity | Role |
|---|---|---|
| Raw accepted envelope | eventId, source partition/offset | Replayable evidence with occurrence and receive times |
| Validation result | eventId, policy version | Valid, invalid, or review decision |
| Dedupe state | eventId and retained identity horizon | Prevent repeated online contribution |
| Window state | ad, window, metric, policy, salt | Partial or complete cumulative count |
| Partial contribution | ad/window, salt, partial version | Replace a salted cumulative total safely |
| Checkpoint manifest | job epoch, generation, source boundary vector | Processing state and output recovered together |
| Result row | interval authority, generation, key, revision | Queryable versioned count |
| Correction manifest | interval, policy, authority version | Select which result supersedes live history |
Raw inputs and checkpoint artifacts are durable; worker memory is not. The checkpoint saves input positions, remembered event IDs, window counts, watermark/control state and output references from the same processing point. Omitting deduplication state would count replayed clicks again after restore even if window counters were restored correctly.
Campaign membership and dimension dictionaries are versioned. If A7 moves campaigns later, the event uses the defined attribution rule and metadata version; historical totals cannot change merely because a current lookup now maps the ad differently.
The result store is a derived view. A correction can be rebuilt from source and policy, but only while the declared raw retention window remains available. Long-term aggregates without retained raw evidence cannot support arbitrary future policy reprocessing.
Deduplicate first by (tenant, source, eventId), retaining the immutable payload fingerprint, then repartition accepted contributions by ad/window. If workers deduplicate only inside an ad key, a conflicting resend that changes adId could reach another worker and be counted again. The checkpoint includes both the identity check and the contribution sent to the counter. Restoring and replaying that checkpoint must not create another contribution. Expire identities only under the documented retry horizon, and enforce an admission rule for older resends rather than silently treating them as new events.
07Basic working design
The smallest service appends validated envelopes to one durable log. One processor reads them in order and uses a local transactional store for event-ID membership and ad/window counters. For C901 it checks the ID, assigns the 09:00 minute from occurredAt, records C901, and changes A7 from 99 to 100 in the same local transaction. It advances the saved input position only in a way that can be recovered with those IDs and counters.
A simple query endpoint reads the resulting counters and labels them preliminary until the declared progress policy closes the window. Top 100 can initially sort complete hourly ad totals because all contributions are on one machine. This baseline shows what is counted, when and how retries behave before adding partitions.
The counter cannot simply use the time at which the request reached the server: that would place C901 into 09:02 and distort the campaign's 09:00 report. Nor can it use only timestamp equality for deduplication: two genuine clicks may share a millisecond.
For a small workload, a database transaction that stores processed IDs, counters, and source bookmark can implement a coherent checkpoint. A separate stream framework is not required merely to explain the problem. Distribution becomes useful when input rate, state, or recovery time exceeds that single transaction domain.
C901 belongs to 09:00; dedupe and counter change share a recoverable boundary.
Read each connection in order
- syncAppend stable C901 / occurredAtClick collectors → Durable accepted event log
- asyncRead ordered eventsDurable accepted event log → Validation / counting process
- syncAtomic identity + count + progressValidation / counting process → Event IDs + windows + bookmark
- syncRead count and finalityCampaign query API → Event IDs + windows + bookmark
08Find the baseline flaws
The first wrong answer is a stateless increment endpoint. A collector times out after C901 was applied, retries it, and A7 becomes 101. Atomic increment prevents lost updates but does not prevent duplicate logical events. Event identity and the increment must share an atomic boundary.
The second is processing-time aggregation. C901 arrives at 09:02:05 and increments the 09:02 bucket despite occurring at 09:00:58. The system is fast but answers a different question. Similarly, closing a window solely because a local wall clock crossed 09:01 ignores delayed input partitions.
The third is merging worker-local top lists before complete aggregation. One worker sees Red 6/Blue 7 and another sees Red 6/Green 7. Their top 1 candidates omit Red, although its complete total 12 wins. A top-k data structure cannot repair missing contributions.
At 200,000 events/s, a single transaction per click can overwhelm one database and a hot ad can monopolize a keyed worker. A checkpoint that copies a large state set too frequently may make recovery protection itself the bottleneck. Use partitions, batches and measured snapshot intervals while preserving duplicate detection, event-time windows and complete totals.
09Improve the design, step by step
First, use replicated partitioned ingestion. The trigger is collector bursts and processor downtime. Authenticated collectors append accepted events to a durable log; processors catch up independently. This protects admitted input and scales intake, at the cost of log storage and visible processing lag. Direct transactional ingestion remains simpler for low traffic. Reduce new admissions before retention would delete accepted records that still need processing.
Second, partition keyed state with recoverable checkpoints. The trigger is per-process counter and dedupe capacity. Route ad/window work to keyed processors and checkpoint source positions, identity state, counters, and progress together. This distributes memory and CPU but introduces cross-partition watermarks, ownership changes, and checkpoint coordination. One database remains preferable while its capacity and restore time fit the target.
Third, split only measured hot keys. The trigger is A7 receiving 40,000 events/s on one owner. Hash event IDs across 16 partial counter keys. Send each partial’s cumulative total and version to a worker that combines the complete ad/window total. This lowers the hot ingress rate per worker but adds a second stage and freshness delay. Uniform hashing by ad is simpler for ordinary skew; unnecessary salting increases state and network work for every ad.
Fourth, separate published query generations and corrections. The trigger is queries seeing partial checkpoint output or a backfill racing live writers. Workers prepare output without exposing it. A metadata transaction checks the current job epoch before publishing a completed checkpoint. A historical correction then selects a new authoritative version for its interval. This gives reproducible results and crash-safe visibility, but costs retained generations, manifest management, and a checkpoint-sized freshness delay. Per-record transactional output is an alternative when an appropriate sink can bear the write rate.
The design does not advertise arbitrary exactly-once effects. It explains which identities and checkpoint/output boundaries make these particular counters reproducible, and which retention and connector assumptions bound that claim.
10Detailed architecture
Identified ingestion and counting
Collectors authenticate through an ingestion gateway that validates envelopes and provenance before durable append. A partitioned replicated log retains accepted input. Raw archival workers preserve source manifests for correction and audit. First group events by their trusted identity, check immutable fingerprints and apply the validation policy. Then route accepted contributions to ad/window counters. Its deduplication state and the count workers’ window state share the checkpoint boundary.
Hot keys optionally pass through salted partial counters and a complete-total reducer. The combining worker tracks each partial’s identity and version. A newer cumulative value replaces the previous value; a retry must not add that total again. The checkpoint coordinator captures a consistent processing boundary and the metadata authority publishes its output generation only after state and result artifacts are durable.
Queries and correction authority
Query servers pin a complete result manifest, read counters and top lists from that generation, and expose watermark/finality. A correction pipeline reads a frozen raw source interval and produces a separate candidate authority version. Publication switches that interval to the corrected version in one transaction and prevents old live workers from overwriting it.
Publication and artifact lifetime
Ingestion acceptance is synchronous. Validation, counting, archival, checkpointing, and corrections are asynchronous. Dashboard visibility waits for the published generation, assumed every two seconds in healthy operation. The metadata service records staging grants that protect objects during upload, manifest references that retain published objects, and reader pins that protect active reads. Cleanup checks these records before marking an object for deletion; an object marked deleting cannot later be published. Waiting before deletion may reduce races, but the atomic metadata checks prevent deletion of an object being published or read.
Implementation option and limits
A practical implementation can use Kafka for replayable input and Flink for identity-keyed state, repartitioning, event-time windows and checkpoints, with durable object storage for snapshots and raw evidence. Use a sink with verified checkpoint integration or implement the staged-generation publication described here. Flink’s operator-state guarantee alone does not make an arbitrary OLAP database transaction part of its checkpoint. Choose the query store for indexed campaign/time reads and version retention; benchmark the complete sink and publication path against the five-second target.
One committed manifest selects both the visible results and the state needed to restore them. A backfill publishes a separately versioned replacement for its historical interval.
Read each connection in order
- sync1. C901 with provenanceAuthorized collectors → Validation / admission gateway
- sync2. Durable accepted appendValidation / admission gateway → Replicated input partitions
- asyncArchive source manifestsReplicated input partitions → Raw source archive
- async3. Replayable inputReplicated input partitions → Keyed validation / partial counters
- asyncVersioned partial totalsKeyed validation / partial counters → Complete ad/window reducers
- sync4. Capture coherent state boundaryCheckpoint coordinator → Keyed validation / partial counters
- syncCapture complete totalsCheckpoint coordinator → Complete ad/window reducers
- syncStage state / dedupe snapshotKeyed validation / partial counters → Durable state / result artifacts
- syncStage result generationComplete ad/window reducers → Durable state / result artifacts
- sync5. Publish durable checkpoint + outputCheckpoint coordinator → Epoch / checkpoint / interval authority
- sync6. Pin selected result authorityCampaign and top-k query API → Epoch / checkpoint / interval authority
- syncRead complete count / top-kCampaign and top-k query API → Durable state / result artifacts
- syncQuery with policy and finalityCampaign dashboards → Campaign and top-k query API
- syncRead frozen source intervalVersioned correction jobs → Raw source archive
- syncBuild correction candidateVersioned correction jobs → Durable state / result artifacts
- syncAtomically supersede intervalVersioned correction jobs → Epoch / checkpoint / interval authority
11Write path and acknowledgement
Recovery must restore both remembered event identities and their count updates. Replicated input allows replay, but the replay must not count a click twice.
| Operation/data | Example |
|---|---|
| Collect | POST /clicks {eventId:C901,adId:A7,impressionId:I88,occurredAt:"2026-09-22T09:00:58Z",collectorVersion:2} |
| Deduplication | (tenant,source,C901) plus immutable fingerprint and policy v3 |
| Aggregate key | (A7,09:00–09:01,validClicks,v3) |
| Published result | {count:100,revision:5,final:false,watermark:09:01:30} |
- Validate the envelope/provenance and append C901 before acknowledging durable acceptance.
- The identity-keyed processor verifies the immutable fingerprint and that C901 is new within the supported retry horizon, then validates its event time and policy result.
- Its accepted contribution reaches the ad/window owner and changes the 09:00 count from 99 to 100. Deduplication, in-flight contributions, source positions and window state belong to a coherent checkpoint.
- The worker stages cumulative count 100/revision 5 under an attempt-specific epoch/generation namespace. Repeated row delivery must match that value; older revisions cannot overwrite it. Staging does not make it visible.
- The coordinator completes durable state and result artifacts, then atomically publishes their checkpoint manifest under the current job epoch. Readers remain on count 99 until this boundary commits.
- The dashboard query pins the published generation and sees 100 in 09:00, with watermark 09:01:30 and a preliminary label. The archive retains the event and policy evidence explaining the change.
- A crash before publication restores the earlier checkpoint and replays C901. A crash after publication restores identity and window state that already include it. Neither case adds a second logical contribution.
- At the closing progress boundary, the window becomes final under policy v3. Beyond-policy events enter a correction path; a privileged job can publish a new authoritative interval version with provenance.
The valid-click status comes from provenance and policy. If a fraud decision changes later, record a new policy result or versioned recomputation. Transport deduplication alone does not prove that distinct human clicks should both be billed.
12Read and delivery path
Read one published generation with its revision, watermark and validation policy. Rank each ad’s complete total and show whether results are preliminary or corrected.
Suppose the 09:00–09:01 window emitted revision 4 when progress crossed 09:01. C901 arrives while the watermark is 09:01:30, still inside its two-minute allowed-lateness period ending at watermark 09:03. Update and emit revision 5. Data arriving beyond the online policy goes to a correction path rather than disappearing invisibly.
The read path first authenticates the dashboard client's campaign scope and pins the current interval/result manifest. It then selects minute rows under one generation, sums only compatible metric and policy versions for rollups, and returns preliminary/final status alongside the count. A top 100 request uses complete hourly ad totals at that same boundary, not whichever partial worker responded most recently.
The watermark is included because count 100 has different meaning while progress is 09:01:30 than after the allowed-lateness boundary. Missing or idle inputs follow a documented rule. Declaring an input idle permits progress, but returning data is still subject to the late-event policy; idleness is not proof that the source will never emit an older record.
Pagination pins the result generation and deterministic order. If the generation expires, the API asks the client to restart rather than mixing pages from different ranking states. Cached results use interval authority, generation, policy, and query dimensions as their identity, so a corrected historical count does not remain hidden behind a stale unversioned cache.
13Correctness deep dive
For a hot ad, salt its key across partial counters, then reduce those partials. Give each partial snapshot an identity and version so a retry replaces its previous total rather than adding the whole count again. A heap tracks the largest current complete totals; count corrections can decrease a winner and require reconsidering candidates.
The crash boundary is as important as the top-k proof. Let checkpoint 4 contain source offset 117, dedupe without C901, and count 99. Worker A processes offset 118, stages count 100/revision 5, and begins checkpoint 5. Output is not yet queryable. Each attempt writes immutable artifacts under its own job epoch and artifact identities, so an obsolete epoch cannot overwrite a replacement’s files even when both use logical generation 5. The following publication call is for the replacement job in epoch 8.
publishCheckpoint(epoch=8, generation=5):
require durable state snapshot and staged result artifacts
require source vector and watermark match that snapshot
begin metadata transaction
require current job epoch == 8 and predecessor == 4
require artifacts are ready, protected, and not deleting
install checkpoint5 and active result generation5 together
transfer staged references to retained manifest references
commit
Per-window revisions still reject repeated or stale row deliveries within a generation, but revision numbers alone are insufficient if two recovered writers can invent conflicting revision 5 values. Checking the current job epoch during publication decides which writer may publish, removing that ambiguity. All state influencing replay, including event-time control progress and validation-policy version, belongs to the checkpoint.
For salted totals, retain each salt's latest cumulative value and version. Updating salt 3 from 6 to 8 contributes a difference of 2 to the complete total, not another 8. Rank only after all required partials at the published boundary are included. Approximate heavy-hitter sketches are an alternative only with their error semantics stated.
A heavy-hitter sketch is a compact approximate summary used to find frequently occurring keys without retaining an exact counter for every key. It is an alternative when memory limits justify a declared approximation. The design here keeps exact complete totals; the comparison later distinguishes that contract from sketch-based candidate selection.
The manifest publishes only output with a completed recovery checkpoint. When a replacement job takes a newer epoch, the metadata service rejects the old job's publication attempts.
Read each connection in order
- syncStage count 100 and checkpoint 5Worker A / epoch 7 → State and output artifacts
- syncCrash before manifest publicationWorker A / epoch 7 → Worker A / epoch 7
- syncRead active result generationQuery server → Manifest authority
- returnGeneration 4: count 99Manifest authority → Query server
- syncAcquire epoch 8; restore checkpoint 4Recovery job / epoch 8 → Manifest authority
- syncReplay C901; stage coherent generation 5Recovery job / epoch 8 → State and output artifacts
- syncPublish checkpoint 5 + result5 atomicallyRecovery job / epoch 8 → Manifest authority
- syncResume late publication from epoch 7Worker A / epoch 7 → Manifest authority
- blockedReject obsolete epochManifest authority → Worker A / epoch 7
- syncPin generation 5: count 100Query server → Manifest authority
14Failure and recovery
Processing approach comparison
| Approach | Benefit | Cost/error |
|---|---|---|
| Exact per-key counters | Auditable counts | More active keys require more counter state |
| Salt then reduce | Distributes hot-key writes | Extra stage and version tracking |
| Heavy-hitter sketches | Bounded approximate candidate state | Algorithm-specific error bounds |
| Batch recomputation | Reproducible corrections | More latency and archived-input work |
Recovery contract
Save input positions together with processing state, remembered event IDs and window counts; restore them together before replay. A checkpoint’s guarantees depend on connectors and sink integration; it cannot magically transact with any external database. Flink checkpointing. Mark idle inputs deliberately so they do not freeze progress indefinitely, and route their returning late data through the same lateness policy.
| Failure or race | Required response and boundary |
|---|---|
| Lost reply or processor restart | A collector loses its acceptance reply and retries C901; retained event identity prevents another contribution. A processing worker dies after staging output but before checkpoint publication; the staged output stays invisible and recovery replays from the last manifest. If it dies after publication, the new worker restores that published checkpoint. A partitioned old worker cannot publish because the metadata authority has advanced its job epoch. |
| Backfill races live output | A backfill builds interval 09:00–10:00 under correction authority c2 from a fixed source manifest. Publication compares the interval's current authority, installs c2, and fences further live writes to that historical namespace. Live processing continues for other intervals. Queries cannot accidentally add live count 100 and backfilled count 100 together: the manifest selects one authority for that interval. |
| Overload threatens retention | During overload, preserve accepted input, show watermark lag, and increase net processing capacity or tighten new admission. Do not advance watermarks merely to make finality look healthy. When input retention is in danger, alert on the oldest required offset and make a controlled recovery decision before irreversible deletion. |
15Operations, security, and cost
If a worker crashes after staging revision 5 but before completed-checkpoint publication, that output remains invisible; recovery replays from the last published checkpoint. The publication transaction checks the job epoch before making output visible; increasing a row revision alone is insufficient. A new backfill publishes a distinct higher-authority interval version so it does not race silently with live output. Retain correction provenance and validation-policy versions.
Authenticate collector tokens, limit forged/future timestamps, minimize user identifiers, and isolate tenant queries. Monitor accepted versus persisted events, duplicate fraction, watermark lag, late-event rates, hot-key skew, checkpoint duration, sink conflicts, and streaming/batch differences. Backpressure before dropping acknowledged input. A change from 99 to 100 is explainable only if the system preserves both the event and its processing rules.
Measure accepted-to-visible lag, checkpoint completion time, oldest unprocessed event age, late-event fraction by source, and correction discrepancies. A rising valid-click count may reflect traffic, duplicate identities, or a validation-policy change; dashboards should expose the policy version so investigators can distinguish those explanations.
A rollout replays a fixed raw interval into a separate candidate namespace and compares event classifications, minute totals, and top-k results before publication. Recovery tests crash before and after manifest commit, resume a stale job epoch, replay duplicate collector batches, and return an idle input with late data. Delete tests verify that staged/current artifacts and query pins prevent premature cleanup.
At one billion IDs/day, 32 GB is only raw identity payload. Keeping seven days instead of one multiplies that raw dedupe set to 224 GB before index overhead. A longer online retry horizon therefore has a real cost. Batch correction from retained raw input may be a better contract for old data than retaining every ID in the low-latency live state indefinitely.
16Decision ledger and limitations
| Choice | Benefit | Cost or limitation |
|---|---|---|
| Event-time windows | Count in the business occurrence interval | Lateness and progress policy |
| Exact retained event IDs | Suppress supported transport retries | Longer retry coverage retains more event IDs |
| Salted partial counters | Relieve a hot ad | Extra reduction stage and partial versions |
| Checkpoint-published generations | Recoverable coherent query state | A checkpoint-sized visibility delay |
| Versioned interval corrections | Reproducible historical repair | Separate authority and audit workflow |
| Approximate heavy hitters | Bounded candidate memory | Error bounds and possible candidate misses |
We choose exact counts for supported validated events, while stating that provenance and fraud rules define which events qualify. “Exactly once” without that identity and policy boundary would be an empty claim. Similarly, finalized under two-minute lateness is a publishing rule, not omniscient knowledge of all future arrivals.
The next limit may be the worker combining one hot ad’s total, memory for remembered event IDs, or time spent saving checkpoints. Measure which one dominates before adding another aggregation layer. If the requirement changes to approximate trends at very large scale, sketches and sampling can be appropriate, but the API and memory cards must stop describing those numbers as exact auditable counts.
17Interview closing
“Clicks belong to event-time windows based on when they occurred, even when delivery is delayed. I durably accept identified events, validate them under a named policy, and maintain event-time counters plus a bounded deduplication horizon. Watermarks control preliminary and final publication, with an explicit correction path for later data.
“At scale, partitioned input and keyed state handle throughput; measured hot ads can be salted, then reduced to complete totals before top-k ranking. I aggregate every key completely before selecting top-k: a globally winning key may be absent from every partial worker top list. Queries pin a complete result generation.
“The hardest crash case is output written before recovery state is saved. I stage output and atomically publish its completed checkpoint manifest under a current job epoch. A crash exposes either the previous recoverable boundary or the new one, not an unrepeatable count. Historical backfills publish a separate interval authority so they cannot race live output.
“The costs are raw retention, dedupe state, checkpointing, and a few seconds of visibility delay. My next test replays a late duplicate through a crash and a backfill, then proves both the count and its provenance remain explainable.”
If the interviewer changes the output into billing, I would add the financial validation, dispute, and settlement contract explicitly. A live operational dashboard would not automatically become the billing ledger.
Practise the interview questions
Say your answer aloud before opening the model answer. Then answer the follow-up and compare the reasoning.
For a metric counting clicks by occurrence time, an event occurred at 09:00:58 and arrived at 09:02:05. Which minute receives it?
Reveal a model answer
For the occurrence-time metric we chose, it belongs to 09:00–09:01 after validating the timestamp. Arrival time tells us when we can process it, not which reporting interval it describes. A processing-time metric would be different and must be labeled accordingly.
Interviewer follow-up
Can you trust every browser timestamp?
Reveal the follow-up answer
No. Validate reasonable bounds and signed/impression provenance, and define how suspicious timestamps are handled. Event-time design is not permission for clients to place events arbitrarily in history or the future.
What the answer must demonstrate: Specify metric semantics before writing a window operator.
Does watermark 09:01 mean all earlier clicks are definitely present?
Reveal a model answer
It is a progress declaration, not an infallible fact about mobile networks. I use it to emit results, then apply the allowed-lateness policy to events that arrive behind that progress. The output carries finality and revision so downstream users understand when counts can still change.
Interviewer follow-up
Why does an idle input matter?
Reveal the follow-up answer
Combined progress often follows the slowest relevant input. An idle partition can stall windows unless marked idle deliberately; when it returns, delayed records still need defined handling.
What the answer must demonstrate: A watermark requires an operational lateness policy.
Revision 5 says count 100. What happens if the sink receives it twice?
Reveal a model answer
Within one authoritative generation, this is a cumulative total: install revision 5 once and require a repeated revision 5 to carry the same value. Do not add 100 twice. Older revisions cannot overwrite newer totals. Publication still waits for the completed checkpoint; row versions alone do not establish a recoverable result.
Interviewer follow-up
What about replay after restoring an old checkpoint?
Reveal the follow-up answer
The recovered writer needs a consistent revision/epoch scheme or idempotent transactional sink. Merely choosing a counter in memory is insufficient if restart can reuse incompatible output versions.
What the answer must demonstrate: Explain how restart preserves output identities and rejects stale writers.
Worker L sees Red=6/Blue=7; worker R sees Red=6/Green=7. Why do their local top 1 lists miss the global winner?
Reveal a model answer
Each worker has only part of Red’s traffic: six on each, for twelve total. Their local sevens win only against partial counts. I must combine complete per-ad window totals before ranking, or use an approximate algorithm with a clearly stated candidate/error guarantee.
Interviewer follow-up
How do you still split a very hot ad?
Reveal the follow-up answer
Spread its incoming updates across partial counters. A second worker tracks each partial’s identity and version and combines them into the complete count needed for ranking.
What the answer must demonstrate: Do not apply complete-owner top-k proofs to partial counters.
A fraud correction changes last week’s count. How do you publish it?
Reveal a model answer
Recompute the relevant interval from retained events under a recorded validation policy and publish a new result version with correction provenance. I would not add a fresh batch total on top of the existing aggregate or erase the reason for the change.
Interviewer follow-up
Is a one-day dedupe cache enough for that replay?
Reveal the follow-up answer
No. It only covers its stated horizon. A batch rebuild can deduplicate the complete input interval and replace the result version instead of depending on expired online membership.
What the answer must demonstrate: Retry dedupe and historical correction have different boundaries.
Does enabling checkpoints make every dashboard write exactly once?
Reveal a model answer
Only if source positions, processing state, and sink behavior cooperate under the recovery protocol. A worker may write output and fail before checkpoint completion. The sink must transactionally coordinate or recognize repeated/older output versions; otherwise replay duplicates effects.
Interviewer follow-up
What should you compare to detect mistakes?
Reveal the follow-up answer
Reconcile sampled or full window totals against a reproducible batch computation from raw events, and monitor duplicate/revision conflicts and watermark lag rather than relying only on uptime.
What the answer must demonstrate: Trace the crash between output and checkpoint.
A published checkpoint contains count 99. A worker stages count 100 for the next checkpoint and crashes before publishing it. What can a reader see?
Reveal a model answer
In this design the write is staged, so readers still see the previous published generation. Only a metadata transaction that installs a durable checkpoint and its output references together makes 100 visible. A crash before that transaction replays from 99; a crash after it restores the state that already includes the click.
Interviewer follow-up
Why are increasing row revisions alone not sufficient?
Reveal the follow-up answer
Two unfenced recovery attempts could produce the same revision with different state or overwrite one another. The current job epoch and predecessor generation decide which complete output is authoritative.
What the answer must demonstrate: Identify publication and recovery as one coherent boundary.
A backfill and live processor both produce 09:00 counts. How do they avoid overwriting each other?
Reveal a model answer
The backfill writes a separate candidate authority version from a fixed source interval. A manifest transaction selects that version and fences live writes for the historical interval. Queries select one authority; they do not sum both copies.
Interviewer follow-up
Can you recompute a three-year-old interval after retaining raw events for only 30 days?
Reveal the follow-up answer
Not under this stated contract unless another archive preserved the necessary evidence and policy inputs. Aggregate counters alone do not support arbitrary future reclassification.
What the answer must demonstrate: Tie correction capability to actual retained evidence.
Blank-page exercise · 45 minutes
Build the answer yourself
Build the dashboard client’s minute counts and hourly top ads. Deliver C901 late and twice, crash after revision 5, and prove the Red/Blue/Green partial-winner example.
- Define event identity, valid-click rules, and occurrence time.
- Calculate payload, aggregate state, and dedupe retention.
- Trace a late event into a versioned replacement result.
- Demonstrate why partial local top-k is unsafe.
- Recover checkpoints and publish an audited historical correction.
Check that each component and design decision follows from your requirements and workload.
Recall the key ideas
Answer from memory before opening each card. Explain why the choice works and what it costs. Revisit missed cards tomorrow.
Design event-time click analyticsWhich clock chooses the reporting window?Recall first, then reveal
The declared event time, after validation, when the metric is defined by when the click occurred.
Occurrence time chooses the window; arrival time determines when processing can begin.
Return to lessonDesign event-time click analyticsWhat is a watermark?Recall first, then reveal
A declared event-time position indicating how far an input has progressed. The processor uses it to emit or close windows; older events can still arrive and need a lateness policy.
Progress with a lateness contract.
Return to lessonDesign event-time click analyticsWhy replace aggregate revision 4 with 5?Recall first, then reveal
A retried cumulative total must not be added again; a newer version replaces the prior materialized result.
Version the total, do not double-add it.
Return to lessonFinal revision
Summary and interview notes
Count validated clicks in the windows where they occurred, with explicit retry and lateness rules. Restore event identities and counters together, rank complete ad totals, and expose only published generations. A historical correction selects a new authoritative version for its interval.
Remember these points
- Deduplicate the tenant/source/event identity before partitioning by mutable event dimensions such as ad ID.
- Watermarks express event-time progress; allowed lateness and correction rules bound finality.
- Rank complete ad totals: merging winners from partial counts can omit the true winner.
- Publish durable checkpoint state and result references together, fenced by job epoch.
- Online event IDs and raw evidence expire. An old-interval rebuild must deduplicate its complete retained input rather than rely on expired online IDs.
Interview tips
- Work through one late duplicate, then crash both before and after publication.
- Use the Red=6+6 versus Blue=7 and Green=7 example to demonstrate the top-k boundary.
- Translate “real time” into collection latency, processing lag and result-publication cadence.
Important qualifications
- A stream framework’s state guarantee does not automatically include an arbitrary external sink.
- A valid-click dashboard is not a billing ledger; provenance and financial validation are separate contracts.
Technical references
- Flink event time and watermarksPrimary explanation of event time, parallel watermarks, late records, and windows.
- Flink checkpointingExplains recoverable stream state and the conditions behind checkpoint-based processing guarantees.
Practice marks stay in this browser.