Hiring BarSupport

Design a Stream Processing System — Staff-Level Case Study

Case study57 min read6 diagrams

Technologies referenced in this case study: Apache Kafka · Flink & Stream Processing · Redis · Cassandra · Time Series DBs

Related: Data Pipeline Patterns · Message Queues · Metrics & Monitoring · Scaling Writes · Degraded Mode Framework

How to Use This Case Study#

Organized for interview use first, reference second. Read front-to-back once. Return to individual sections for targeted review.

ModeTimeWhat to Read
Quick Review15 minExecutive Summary → Interview Walkthrough → Fault Lines table → Active Drills 1–3
Targeted Study1–2 hrsExecutive Summary → Walkthrough → Section 3 (Fault Lines) → Section 4 (Failures) → weak-spot Deep Dives
Deep Dive3+ hrsEverything, including the Principal Lens and appendices
What is Stream Processing? — Why interviewers pick this topic

Stream processing computes continuously over unbounded event data: aggregating ad clicks per campaign per minute, detecting fraud within seconds of a card swipe, maintaining a live view of inventory. Instead of "run a query over a table," it is "keep a query's answer up to date as rows keep arriving — including rows that arrive late, out of order, twice, or not at all."

Before vs After — Ad-click billing scenario:

With a nightly batch job:
t=0:        Advertiser's campaign goes viral at 10:00; budget $5,000/day
t=+2h:      Spend actually hits $5,000 — nothing stops it
t=+14h:     Nightly batch computes spend: $41,200
t=+1d:      Company eats $36,200 of over-delivery; advertiser disputes the invoice

With stream processing (event-time windows, 1-min tumbling, 2-min allowed lateness):
t=0:        Same viral spike
t=+1min:    Per-campaign spend updated every minute from click events
t=+2h:      Spend crosses 95% of budget → pacing service throttles serving within ~90 s
t=+2h02m:   Campaign paused at $5,040 — 0.8% over-delivery, within contract tolerance
t=+1d:      Batch reconciliation over the same events matches streaming totals within 0.05%

Why interviewers reach for this question: Every data-heavy company has a streaming pipeline, and almost every one has had an incident where the numbers were silently wrong. The question tests whether you understand time (event vs processing), completeness (watermarks, late data), correctness under failure (exactly-once, idempotent sinks), and the operational reality of state and replay — not whether you can name Flink operators.

Mechanics Refresher: Windows and Processing Guarantees
ConceptHow It WorksProsCons
Tumbling windowFixed, non-overlapping (e.g., each minute)Simple, one result per key per windowBoundary effects; bursty emission
Sliding / hopping windowFixed size, advancing by a smaller step (5 min every 1 min)Smooth trendsEach event in size/step windows → 5× state/work
Session windowGap-based; closes after N min of inactivityModels user behaviorUnbounded length; merges on late data
Global window + triggersOne window per key; custom emissionFlexible (running totals)You own all correctness logic
At-most-onceCommit offsets before processingLowest latencyLoses data on failure
At-least-onceCommit after processingNo lossDuplicates on replay
Exactly-once (effectively-once)Checkpointed state + offsets, transactional/idempotent sinksCorrect countsLatency ≈ checkpoint interval for transactional sinks; complexity

For most production systems: Kafka as the log, Flink (or Kafka Streams for simple per-key work) as the engine, event-time tumbling windows with watermarks, checkpointed state, and idempotent upserts into the serving store keyed by (key, window_start). The engine is almost never the interview question — time, completeness and correctness under replay are.


Executive Summary

If you only read one section, read this. Everything in the case study flows from the contrast below.

What This Interview Actually Tests#

Stream processing is not a framework question. Everyone can say "Kafka plus Flink."

It is a correctness-under-time question that tests:

  • Whether you separate event time from processing time — and decide what "the answer for 10:05" means
  • Whether you state a completeness policy: how late is too late, and what happens to late events
  • Whether "exactly-once" in your mouth means an end-to-end property you can defend, including the sink
  • Whether you can replay, backfill and upgrade a stateful job without corrupting outputs
  • Whether you name who owns lag, late data, schema breaks and reconciliation

The key insight: Every streaming answer is provisional. Staff engineers decide, per use case, how long to wait for completeness, how to correct results after they're emitted, and how to prove the numbers against a batch source of truth.

The L5 vs L6 Contrast — Start Here#

BehaviorSenior (L5)Staff (L6)Principal (L7)
First move"Kafka → Flink → Cassandra""What's the consumer of the output, and what's the cost of a wrong vs late number?""Is this one more bespoke pipeline, or should it be a query on the org's streaming platform with a shared event contract?"
TimeProcessing-time windowsEvent-time windows with watermarks; states allowed lateness with a numberStandardizes event-time semantics and timestamp ownership across producers org-wide
Exactly-once"Flink has exactly-once"End-to-end: replayable source + checkpointed state + idempotent/transactional sink; names the boundary where it breaksDefines which metrics are "billing-grade" (reconciled) vs "operational-grade" (approximate), with different platforms and SLOs
Late dataDrops itAllowed lateness + retractions/upserts + side output + batch reconciliationSets org policy on restatement: how far back numbers can change, who's notified
State"State in RocksDB"Sizes state (keys × window × bytes), TTL, checkpoint duration, rescale strategyPrices state and replay capacity; retention on the log as a recovery budget
OwnershipPipeline team owns everythingProducers own schemas + timestamps; platform owns engine; consumers own lateness policyRedraws boundaries: data contracts with producer SLAs; streaming platform as a product
Why "time" separates levels

L5: Windows by arrival time at the processor. Works in a demo. In production, a mobile client uploads 40 minutes of buffered events after regaining signal, a Kafka partition stalls for 3 minutes, or a job restarts and replays an hour — and each of these moves events into the wrong windows. The numbers are wrong and nothing alerts.

L6: "Windows are by event time. The watermark is my estimate of 'I've seen everything up to T.' I'll set it at max observed event time minus 30 seconds, allow 5 minutes of lateness with updates, and route anything later to a side output that feeds the batch correction. Dashboards show a 'provisional' badge until the window is final."

Why "exactly-once" separates levels

L5: Enables the engine's exactly-once mode and considers it solved. But exactly-once in Flink covers state and offsets inside the job; if the sink is an HTTP call to a payments service or a non-transactional database insert, replays after a failure duplicate side effects.

L6: Treats exactly-once as an end-to-end chain: "Source is replayable Kafka; state and offsets are checkpointed together every 30 seconds; the sink is an upsert keyed by (campaign_id, window_start), so replays overwrite rather than add. For side effects like sending an alert, I dedupe on an idempotency key at the receiver. Where I can't do that, I say it's at-least-once."

Why "ownership" separates levels

L5: The pipeline team owns everything from the producer's SDK to the dashboard. Every upstream schema change and every clock bug becomes their incident.

L6: Draws the contract lines: producers own event schemas (registry, compatibility checks) and the correctness of event_time; the streaming platform owns the engine, checkpoints and capacity; each consuming product owns its lateness policy and what "final" means for its users.

The Staff Positions#

PositionRationale
Event time, always, for anything user-facing or billedProcessing time makes results depend on outages and replays
Idempotent sinks over transactional magicUpsert by (key, window) makes replay safe regardless of engine guarantees
State allowed lateness explicitly (e.g., 5 min) and route the restWaiting forever is unbounded state; dropping silently is wrong numbers
Batch is the source of truth for moneyStream for speed, batch (or re-stream) for reconciliation; publish the diff
Keep the log long enough to replayKafka retention ≥ 7 days (or tiered storage) is your recovery budget
Pre-aggregate at the edge of the hot keyTwo-phase (local then global) aggregation for skewed keys
Version stateful jobs like databasesSavepoints, state schema evolution, blue/green for incompatible changes

The Three Intents#

IntentConstraintStrategyFailure ModeCorrectness Bar
Real-time analytics / billing aggregationCounts must reconcile with moneyEvent-time windows, idempotent upserts, batch reconciliationDouble/under-counting on replay or late data≤ 0.1% diff vs batch; finalized windows
Low-latency detection / alerting (fraud, anomaly)Decide in < 1 sPer-key stateful rules/CEP, processing-time timers where neededFalse negatives during lag; alert storms on replayLatency SLO p99 < 1 s; recall over precision
Continuous materialized views / CDC propagationDownstream view mirrors sourceOrdered per-key changelog, upsert/delete semanticsOut-of-order updates overwrite newer statePer-key ordering; convergence within seconds

🎯 Staff Move: "I'll assume billing-grade ad-click aggregation, because it forces the hardest questions — event time, late data, exactly-once to the sink, and reconciliation. Fraud detection optimizes for latency and tolerates approximate counts; CDC views optimize for per-key ordering. Those change the design."

The Five Fault Lines#

#Fault LineThe Tension
1Latency vs CompletenessEmit fast with incomplete data, or wait for stragglers? (watermarks, allowed lateness)
2Exactly-Once vs Throughput & SimplicityTransactional sinks and aligned checkpoints, or at-least-once + idempotency?
3Local State vs External StateEmbedded RocksDB state (fast, hard to rescale) or remote store (simple, slow)?
4Streaming-Only vs Stream + Batch ReconciliationOne code path (Kappa), or a batch source of truth (Lambda-ish)?
5Platform vs Team-Owned PipelinesShared engine and contracts, or each team runs its own jobs?

In the Wild: Real Production Systems#

Why this section belongs here: Citing specific production systems demonstrates you've studied operational reality, not textbook designs.

Google — MillWheel and the Dataflow Model#

Google's MillWheel paper (2013) introduced low watermarks and exactly-once per-key processing via persistent state and deduplication; the Dataflow Model paper (2015) generalized it into what / where / when / how: what results are computed, where in event time (windows), when in processing time they're emitted (triggers + watermarks), and how refinements relate (discarding, accumulating, retracting). Apache Beam is the open-source descendant.

Staff insight: The four questions are the interview. If you answer "when do results fire, and how do late updates modify them," you've covered what most candidates never reach.

LinkedIn — Kafka and Samza#

Kafka was created at LinkedIn as a durable, replayable log; Samza was built there to process it with local state backed by a changelog topic, so a failed task restores state by replaying its changelog. Jay Kreps's "Questioning the Lambda Architecture" argued for reprocessing by replaying the log with a new job version instead of maintaining separate batch code.

Staff insight: Replayable log + changelog-backed state is the foundation of every modern backfill story. Retention is not a storage setting — it's how far back you can recover.

Both have publicly described large Apache Flink deployments: Uber for real-time pricing, marketplace metrics and its AthenaX SQL-on-streams platform; Netflix for its Keystone data pipeline routing trillions of events per day to sinks and real-time processing jobs. Both evolved toward self-service streaming platforms where teams submit SQL or config rather than operating clusters.

Staff insight: At scale, the org problem (hundreds of jobs owned by dozens of teams) dominates the engine problem. The platform — templates, guardrails, lag alerts, managed savepoints — is the product.

What Interviewers Probe#

After You Say...They Will Ask...(What They're Evaluating)
"1-minute windows""Event time or processing time? An event arrives 10 minutes late — which window?"Time semantics
"Watermarks""How do you set it? One partition goes idle — what happens?"Watermark mechanics and stalls
"Exactly-once""Your sink is Postgres/an HTTP API. Still exactly-once?"End-to-end reasoning
"State in RocksDB""How big? How long to checkpoint? How do you rescale from 32 to 64?"State operations
"We'll reprocess""How, without double-writing the serving store?"Backfill design
"Kafka partitions by campaign""One campaign is 40% of clicks. Now what?"Skew handling

System Architecture Overview#

Diagram: System Architecture Overview

Reading the diagram: Kafka is both the input and the recovery mechanism — retention defines how far back you can replay. Flink holds keyed state (dedup set, open windows) checkpointed to object storage every 30 s. Results are upserts keyed by (campaign, window_start), so replays and late updates overwrite instead of add. The batch path recomputes from the lake daily and is the source of truth for invoices; the diff job is what tells you the stream is lying.

Quick-Reference: The 30-Second Cheat Sheet#

TopicThe L5 AnswerThe L6 Answer — Say This
Time"1-minute windows""Event-time tumbling windows; watermark = max event time − 30 s; 5-min allowed lateness with upserts; later goes to a side output."
Exactly-once"Flink exactly-once mode""Replayable source + checkpointed state/offsets + idempotent sink keyed by (key, window). Side effects dedupe at the receiver."
Late data"Drop it""Update the window within lateness; after that, side output → batch correction. Dashboards mark windows provisional until final."
State"RocksDB""~400 GB across 128 slots, incremental checkpoints every 30 s, state TTL 24 h on dedup; rescale via savepoint."
Backfill"Rerun the job""Run v2 from a savepoint or from an earlier offset into a shadow table, diff, then swap the read alias."
Skew"More partitions""Two-phase aggregation: salt hot keys into N sub-keys, pre-aggregate, then merge."

Key Numbers Worth Memorizing#

MetricValueWhy It Matters
Kafka partition throughput~10–50 MB/s write per partition (hardware-dependent)Sizing partitions; ~1 KB events → 10–50K events/s/partition
Flink per-slot throughput (simple keyed agg)~50K–200K events/s1M events/s ≈ 10–20 slots before headroom
Checkpoint interval10 s – 5 min (30–60 s typical)Transactional sinks commit at checkpoints → end-to-end latency floor
Checkpoint duration alert> 50% of intervalBackpressure or state growth; risk of never completing
Watermark delay5 s – 2 min typicalThe price of completeness paid in latency
Allowed lateness1 min – 24 hEvery unit of lateness is retained state
Mobile event delayp99 minutes, p99.9 hours–daysWhy a batch correction path exists
Dedup state sizekeys/day × ~50–100 B2B click IDs/day ≈ 100–200 GB of state
Log retention for replay≥ 7 days hot; 30–90 days tieredRecovery budget for bugs found late
Stream vs batch diff tolerance≤ 0.1% billing, ≤ 1–2% ops dashboardsDefines "correct enough" per intent
Restore time from checkpoint~1–5 min per 100 GB (parallel, local SSD)Drives RTO for stateful jobs

Interview Walkthrough

The most common mistake: Candidates spend 20 minutes drawing Kafka → Flink → database and naming operators, then run out of time before discussing late data, replay and correctness — the only parts that determine level. Get the pipeline on the board in ~10 minutes.


Phase 1: Requirements & Framing (2–3 min)#

State scope in one sentence:

"We're aggregating ad clicks into per-campaign spend per minute to drive budget pacing and, eventually, invoices."

Then the non-functional constraints that pick the design:

"Three things decide this: how fresh the numbers must be, how correct they must be — and against what source of truth — and how late events can arrive. I'll assume 1M clicks/s peak, pacing needs spend within ~1 minute, invoices need ≤ 0.1% error, and mobile clicks can arrive hours late."

Commit:

"Pacing and invoicing have different correctness bars. I'll build one streaming path optimized for pacing, and use a batch recompute as the invoice source of truth, with a continuous diff between them."

🎯 Staff Move: Separate the decision the output drives (throttle a campaign) from the record it creates (an invoice). The first wants speed; the second wants completeness. One pipeline rarely does both well.


Phase 2: Core Entities & API (1–2 min)#

  • ClickEvent (click_id, campaign_id, ad_id, event_time, ingest_time, cost_micros, device_id)
  • WindowAggregate (campaign_id, window_start, clicks, spend_micros, is_final, version)
  • Watermark (per source partition and per operator; the completeness frontier)
  • LateCorrection (events beyond allowed lateness, written to the lake for batch)
Produce(ClickEvent)                                    # at-least-once from ad servers, idempotent by click_id
GetSpend(campaign_id, from, to) → [WindowAggregate]    # includes is_final flag
GetSpendRunningTotal(campaign_id, day) → micros        # for pacing, ~1 min fresh

"is_final is part of the API. Consumers must know whether a number can still change."


Phase 3: High-Level Architecture (≤5 min)#

Diagram: Phase 3: High-Level Architecture (≤5 min)

Walk the flow:

  1. Ad servers stamp event_time at the edge and produce to Kafka keyed by campaign_id.
  2. Flink dedupes by click_id (at-least-once producers retry).
  3. Event-time tumbling windows fire when the watermark passes window_end; late events within 5 minutes re-fire an updated result.
  4. Results upsert into the serving store by (campaign_id, window_start); pacing reads running totals.
  5. Raw events also land in the lake; daily batch recomputes invoices and diffs against the stream.

🎯 Staff Move: "This works on a good day. What decides whether it's correct is what happens with late events, restarts, a hot campaign and a code bug we find three days later. Let me go there."


Phase 4: Transition to Depth (1 min)#

"Three areas: time and completeness — watermarks and late data; correctness under failure — exactly-once to the sink; and state operations — sizing, rescaling and backfill. I'd start with time, because it's where streaming numbers silently go wrong. Which would you prefer?"


Phase 5: Deep Dives (25–30 min)#

Deep dive 1: Watermarks and late data (8–10 min)

"The watermark is a heuristic: bounded out-of-orderness of 30 s means I assume nothing older than max_event_time − 30 s is still coming. The window [10:05, 10:06) fires when the watermark passes 10:06, i.e., roughly when I've seen an event at 10:06:30. Events for that window arriving within 5 more minutes re-fire an updated aggregate — the sink upserts, so the number converges. After that, events go to a side output and the lake, and the batch job corrects them."

Quantify: "From the lake I measured: 99.2% of clicks arrive within 30 s, 99.9% within 5 min, 99.99% within 6 h. So a 30-s watermark with 5-min lateness captures 99.9% in-stream; batch catches the rest."

Idle partitions: "If one Kafka partition gets no events, its watermark doesn't advance and holds back the whole job — windows never fire. Configure idleness detection (e.g., 60 s) so idle partitions are excluded from the min."

Deep dive 2: Exactly-once end-to-end (7–8 min) — Replayable source + checkpoints + idempotent sink; transactional Kafka sink when chaining jobs; dedup for producer retries. Name where it breaks: external side effects.

Deep dive 3: State and skew (5–7 min) — Size state, pick RocksDB with incremental checkpoints, two-phase aggregation for hot campaigns.

Deep dive 4: Backfill (3–5 min) — Fix a bug by replaying from Kafka/tiered storage into a shadow table, diff, and swap the read alias.


Phase 6: Wrap-Up (2–3 min)#

"Summary: Kafka keyed by campaign, Flink with dedup and event-time windows, 30-s watermark, 5-min lateness with upserts, side output to the lake, batch as invoice truth with a 0.1% diff alarm, 7-day hot retention for replay. Next I'd build per-producer freshness SLOs and a self-service backfill tool. I wouldn't build a custom stream engine or a real-time invoice — invoices can wait a day."

Common Timing Mistakes#

MistakeTime LostFix
Explaining Kafka internals (ISR, segments)5–8 min"Kafka, RF=3, 256 partitions, 7-day retention."
Enumerating window types3–5 minPick one and justify it
Designing the dashboard3–5 min"Serving store supports range reads by campaign."
Saying "exactly-once" and moving on—Interviewer will drill; state the sink contract first
Skipping reprocessing—It's the day-3 incident; raise it yourself

1. The Staff Lens#

1.1 Why This Problem Exists in Staff Interviews#

Streaming pipelines fail quietly. A broken web service returns 500s; a broken stream job returns plausible numbers that are 3% low. The Staff skill is designing for verifiable correctness: explicit completeness policies, idempotent outputs, reconciliation against a source of truth, and replay as a routine operation. Interviewers use streaming to see whether you think about time and failure as first-class, and whether you know who gets hurt when a number is wrong.

1.2 The L5 vs L6 Contrast — Visual#

Diagram: 1.2 The L5 vs L6 Contrast — Visual

1.3 The Staff Question That Cuts Through Everything#

"When is a number final, and what happens if it changes after someone acted on it?"

If the answer is "it's final when emitted," you need to wait long enough to be complete and pay in latency. If numbers can change, every consumer must handle updates (upserts, retractions, is_final) and someone must decide how far back restatements are allowed. Everything else — watermarks, lateness, sinks, reconciliation — is implementation of that answer.


2. Problem Framing & Intent#

2.1 The Three Intents — Explained#

Real-time analytics / billing aggregation. Correctness is measured against a ledger. Numbers must reconcile. Design: event time, dedup, idempotent sinks, finalization, batch truth. Latency target: ~1 min. Typical scale: 100K–10M events/s.

Low-latency detection (fraud, anomaly, abuse). Value decays in seconds: a fraud decision after the card transaction completes is useless. Design: per-key state (velocity counters, recent history) with sub-second processing; approximate is fine; missing an event is worse than a false positive. Replays must not re-fire actions (suppress alerts older than N minutes).

Continuous materialized views / CDC. Mirror a source database into a cache, search index or read model. Design: per-key ordering from the changelog, upsert/delete semantics, version-guarded writes so an older update never overwrites a newer one. Correctness: convergence within seconds.

🎯 Staff Move: "Fraud tolerates approximation but not latency; billing tolerates latency but not approximation. If product wants both from one job, I'll push back and split them — sharing a job couples their failure modes."

2.2 When NOT to Use Stream Processing#

SituationBetter AlternativeWhy
Consumers read results dailyHourly/daily batch (Spark, warehouse SQL)Cheaper, simpler, naturally complete
Freshness need ≥ 15 minMicro-batch or scheduled SQL on the warehouseNo watermarks, no stateful ops
Simple per-event transform, no stateStateless consumer / serverless functionEngine overhead isn't worth it
Exact answers over all historyBatch or OLAP queryUnbounded state in a stream
Low volume (< 1K events/s) with a DB alreadyDB triggers / outbox consumersOne fewer distributed system
Team has no one to own stateful jobsManaged streaming SQL (e.g., cloud-managed Flink/ksqlDB) or batchStateful jobs need an owner at 3am

"If the dashboard is looked at once a morning, a streaming pipeline is a cost with no customer."

2.3 What the Interviewer Leaves Underspecified#

GapWhy It MattersWhat to Say
Lateness distributionSets watermark and lateness"Do we know p99/p99.9 delay per producer? I'll measure it from the lake."
Consumer of the outputSets finality semantics"Does anyone act on a number that may change?"
Duplicate rate from producersSets dedup design"Are producers at-least-once with retries?"
Key cardinality and skewSets state size and hot-key handling"How many campaigns? Top campaign's share?"
Correction windowSets restatement policy"Can yesterday's numbers change? Last month's?"
Source of truthSets reconciliation"Is the batch/ledger the truth, or is the stream?"

2.4 Precise Terminology#

TermPrecise MeaningCommon Confusion
Event timeWhen the event happened (set by producer)Confused with ingest or processing time
Processing timeWall clock of the operator processing itSeems simpler; results depend on lag
WatermarkAssertion that no events with timestamp ≤ W are expectedNot a guarantee — a heuristic
Allowed latenessHow long after the watermark passes a window it still accepts updatesNot the same as watermark delay
TriggerWhen a window emits (on watermark, early every N s, on late event)Assumed to be "once at window end"
Accumulating vs retractingLate firing emits the new total vs emits (−old, +new)Downstream sums double-count with accumulating + append sinks
CheckpointConsistent snapshot of state + source offsetsNot the same as a savepoint (user-triggered, portable)
Exactly-onceEach event affects state/output exactly once as observedNot "processed once" — replays happen; effects are deduped
BackpressureSlow operator throttles upstreamSymptom, not a failure; lag grows
Kappa vs LambdaReprocess by replaying the stream vs a separate batch layerPresented as a religion; it's a cost tradeoff

3. The Five Fault Lines#

3.1 Fault Line 1: Latency vs Completeness#

The tension: Every second you wait before emitting a window captures more late events and costs freshness plus state. Never waiting means emitting partial numbers.

StrategyWhat WorksWhat BreaksWho Pays
Processing-time windowsLowest latency, trivialResults depend on lag/replays; wrong after outagesConsumers (silently wrong numbers)
Event time, tight watermark (5 s), no latenessFast, final on emitDrops stragglers (~1–5%)Advertisers / finance (under-count)
Event time, 30 s watermark + 5 min lateness + upsertsConverges, ~99.9% capturedNumbers change after first emit; state × 6 windowsConsumers (handle updates) + platform (state)
Long watermark (1 h)Nearly complete1-hour latencyPacing (useless for real-time)
Early + on-time + late triggersSpeculative early results, refinedConsumers must understand refinementProduct (UX for provisional)
Diagram: 3.1 Fault Line 1: Latency vs Completeness

Staff default: Bounded-out-of-orderness watermark set from measured delay (p99 of source-to-ingest delay, e.g., 30 s), allowed lateness at ~p99.9 (5 min), side-output beyond that, and a batch correction for the tail.

When to deviate: Fraud — emit per event with processing-time timers; completeness doesn't matter. Invoicing — don't stream; batch after a 24-hour completeness delay.

🎯 Staff Move: "I don't pick the watermark by intuition — I measure the delay distribution per producer from the lake and set watermark at p99 and lateness at p99.9. Then I alert on late_events_dropped_total so I know when the distribution shifts."

3.2 Fault Line 2: Exactly-Once vs Throughput & Simplicity#

The tension: End-to-end exactly-once requires coordinating checkpoints with sink commits — adding latency (results visible only at checkpoint) and complexity (transaction timeouts, zombie fencing). At-least-once is simpler and faster, but duplicates must be absorbed.

ApproachWhat WorksWhat BreaksWho Pays
At-most-onceFastData loss on crashEveryone downstream
At-least-once, append sinkSimple, fastDuplicates on replay (e.g., 30 s of events re-added)Consumers (inflated counts)
At-least-once + idempotent upsert sinkReplays overwrite; no coordinationNeeds deterministic keys; not for side effectsSink design (key schema)
Two-phase commit sink (Kafka transactions)True end-to-end EOS into KafkaOutput visible only after checkpoint (30 s+); read_committed consumers; transaction timeoutsLatency budget
Side effects with idempotency keysSafe external callsReceiver must support dedupReceiving team
Diagram: 3.2 Fault Line 2: Exactly-Once vs Throughput & Simplicity

Staff default: Checkpointed state + replayable source + idempotent upsert sink with deterministic keys. Use transactional Kafka sinks only between jobs where a downstream job reads with read_committed. For external side effects (alerts, payments), use idempotency keys derived from (key, window, version).

When to deviate: If the sink cannot upsert (e.g., an append-only event bus consumed by others), use transactional producers and accept checkpoint-interval latency.

3.3 Fault Line 3: Local State vs External State#

The tension: Embedded state (RocksDB on local SSD) gives microsecond access and scales with parallelism, but makes rescaling, recovery and upgrades heavy. External state (Redis, Cassandra) makes jobs stateless and easy to scale but adds a network hop per event and moves consistency problems to the store.

ApproachWhat WorksWhat BreaksWho Pays
Heap stateFastestGC pauses, bounded by memoryOn-call (OOM)
RocksDB + incremental checkpointsTB-scale state, ~µs readsRestore/rescale time (minutes per 100 GB); compaction tuningPlatform (state ops)
External KV (Redis/Cassandra)Stateless workers, easy rescale0.5–2 ms per access; 1M events/s → 1M+ store ops/s; exactly-once harderStore owners + latency
Hybrid: local cache + externalBalanceCache coherence on rebalanceJob owners

Staff default: RocksDB with incremental checkpoints to object storage and local recovery enabled; state TTLs on everything; size state explicitly.

State sizing example: 2B click IDs/day for dedup × ~80 B = ~160 GB; open windows: 500K active campaigns × 6 open windows (1 + 5 min lateness) × ~200 B = ~0.6 GB. Dedup dominates — so consider a bloom filter per hour or dedup only within 10 minutes (where 99.9% of retries fall) to cut it 100×.

3.4 Fault Line 4: Streaming-Only vs Stream + Batch Reconciliation#

The tension: A single streaming code path (Kappa) avoids maintaining two implementations; a batch recompute provides an independent truth and catches streaming bugs, at the cost of two code paths that can disagree.

ApproachWhat WorksWhat BreaksWho Pays
Stream onlyOne codebase; replay to fixNo independent check; bugs undetectedFinance (silent errors)
Lambda (separate batch + stream code)Independent truthTwo implementations drift; double maintenancePipeline team (2× code)
Same logic, two runners (Beam/Flink batch mode, SQL)One definition, two executionsEngine semantic differences at edgesPlatform
Stream + batch diff only for billing-grade outputsTruth where it mattersDiff pipeline to ownPipeline team (moderate)

Staff default: Unified logic (same SQL/Beam pipeline) executed in streaming for freshness and in batch over the lake daily for billing-grade outputs; a diff job alerts at > 0.1%. Operational dashboards are stream-only.

3.5 Fault Line 5: Platform vs Team-Owned Pipelines#

OptionWhat WorksWhat BreaksWho Pays
Each team runs its own Flink clusterAutonomy30 clusters, 30 upgrade schedules, inconsistent checkpointingEvery team's on-call
Central platform, teams write jobsShared ops, standard alertsPlatform bottleneck; noisy neighbors on shared clustersPlatform team
Self-service SQL / templates on platformFast onboarding, guardrailsLimited expressivenessPower users
Managed cloud serviceNo cluster opsCost; version lag; lock-inBudget

Staff default: Platform-owned engine (per-job clusters via Kubernetes operator for isolation), standard templates for dedup/windowing/sinks, mandatory lag and checkpoint alerts; product teams own job logic, lateness policy and output contracts.

🧭 Principal Move: "The scarce resource isn't compute — it's people who can debug a stuck checkpoint at 3am. A platform with 5 templates that cover 80% of jobs is worth more than infinite flexibility."


4. Failure Modes & Operational Reality#

4.1 The Stalled Watermark — Windows Never Fire#

Scenario: One of 256 Kafka partitions stops receiving events after a producer-side routing change.

t=0:       Producer deploy routes no traffic to partition 118
t=+0s:     Partition 118's watermark frozen at 14:02:10
t=+1min:   Job watermark = min over partitions = 14:02:10; no windows after 14:02 fire
t=+5min:   Pacing service sees spend frozen; campaigns keep serving past budget
t=+20min:  Open window state grows 20× (every campaign, every minute, unfired)
t=+35min:  Advertiser support: "spend hasn't updated"; on-call sees consumer lag = 0 and is confused
t=+50min:  Root cause: idle partition; idleness timeout not configured

Why it's nasty: consumer_lag looks perfect — the job is reading everything. Only watermark_lag_seconds (now − current watermark) exposes it. Detection: watermark_lag_seconds > 2× configured watermark delay for 2 min; window_fired_total rate drops to 0. Blast radius: Every output of the job; downstream pacing over-delivers — ~$X per minute of ad spend uncapped. Mitigation: Enable source idleness (e.g., 60 s) so idle partitions don't hold back the watermark; restart is not needed once configured. Prevention: Template default for idleness; alert on watermark lag, not just consumer lag. Owner: Streaming platform (template + alert); pipeline owner (response).

🎯 Staff Move: "Consumer lag tells you whether you're reading. Watermark lag tells you whether you're answering. Alert on both."

4.2 Replay Double-Count — The Append Sink#

Scenario: A job with an append-only sink (INSERT of per-window deltas) restarts from a checkpoint taken 4 minutes before the crash.

t=0:       TaskManager OOM; job restarts from checkpoint at t−4min
t=+1min:   Job replays 4 min of events (240M clicks) and re-inserts their deltas
t=+1min:   Serving store now contains 4 minutes of duplicated spend
t=+2min:   Pacing thinks 1,800 campaigns hit budget early → throttles them
t=+6h:     Advertisers complain of under-delivery; batch diff shows stream +3.1% for that hour

Detection: reconcile.stream_vs_batch_diff_pct > 0.1%; job.restarts_total correlated with spikes in sink.rows_written. Mitigation: Recompute affected windows from Kafka and overwrite. Prevention: Sinks upsert absolute window values keyed by (campaign_id, window_start), never deltas; or transactional sinks. Owner: Pipeline owner; platform enforces sink template.

4.3 Checkpoint Death Spiral#

Scenario: State grows past what the checkpoint interval can snapshot under backpressure.

t=0:       Traffic 1.6× normal (holiday); operator backpressured
t=+5min:   Checkpoint barriers queue behind buffered records; checkpoint takes 4 min (interval 1 min)
t=+15min:  Checkpoints start timing out (10-min timeout); last successful one ages
t=+40min:  TaskManager fails; restore from 40-min-old checkpoint → replays 40 min at 1.6× load
t=+60min:  Replay creates more backpressure; checkpoints fail again; lag grows 3 hours

Detection: checkpoint_duration_ms > 50% of interval; checkpoint_failures_total > 0; last_successful_checkpoint_age_s > 3× interval. Mitigation: Unaligned checkpoints or buffer debloating to decouple barriers from backpressure; temporarily scale out (via savepoint); shed non-critical outputs. Prevention: Capacity at 2× peak; incremental RocksDB checkpoints; state TTL; load test at holiday peak. Owner: Platform (checkpoint config) + pipeline owner (capacity).

4.4 Hot Key — One Campaign at 40% of Clicks#

A Super Bowl campaign sends 400K clicks/s to one key → one subtask at 100% CPU, the rest at 15%. Backpressure propagates; the whole job's lag grows. Detection: subtask.busy_time_pct max/median > 3; records_in_per_s skew by subtask. Mitigation: Two-phase aggregation — salt campaign_id with hash(click_id) % 32 for local pre-aggregation, then merge by campaign_id. Hot-key detection can apply salting only to the top-N keys. Owner: Pipeline owner.

4.5 Schema Break Upstream#

A producer renames cost_micros to cost without registry compatibility checks. Deserialization fails; the job either crashes in a loop or (worse) defaults cost to 0. Detection: deserialization_errors_total > 0; spend_per_click_avg drops to 0. Mitigation: Route bad records to a DLQ topic; fail fast on schema-incompatible changes at produce time (registry enforced backward compatibility). Owner: Producer team (schema contract); platform (registry enforcement).

4.6 Late-Data Surge After Mobile Outage#

A mobile SDK bug holds events for 6 hours; the fix releases 800M delayed clicks in 20 minutes. All exceed allowed lateness → side output → lake. Stream under-reports yesterday by 4%; pacing made decisions on under-counts. Detection: late_events_side_output_total spike; event_delay_seconds p99 jumps. Mitigation: Batch correction restates affected windows; finance notified per restatement policy. Owner: Mobile team (SDK); pipeline owner (restatement); finance (sign-off).

4.7 Operational Reality Matrix#

FailureDetection SignalBlast RadiusMitigationOwner
Stalled watermarkwatermark_lag_seconds > 2× delayAll windowed outputsIdleness config; alert on watermark lagPlatform + pipeline
Replay double-countstream_vs_batch_diff_pct > 0.1%Windows since last checkpointUpsert sinks; recomputePipeline owner
Checkpoint spiralcheckpoint_duration_ms > 50% intervalWhole job; hours of lagUnaligned checkpoints; scale via savepointPlatform
Hot keysubtask busy-time skew > 3×Whole job via backpressureTwo-phase aggregationPipeline owner
Schema breakdeserialization_errors_totalAffected fields/recordsDLQ; registry enforcementProducer team
Late surgelate_events_side_output_totalPast windows under-countedBatch restatementPipeline + finance
Consumer lagconsumer_lag_seconds > SLOFreshnessScale out; shedPipeline owner
State explosionstate_size_bytes growth > 20%/dayCheckpoints, restoresTTL; reduce dedup horizonPipeline owner
Kafka retention exceeded before fixoldest_needed_offset < log_start_offsetUnrecoverable from streamRebuild from lakePlatform

5. Evaluation Rubric#

5.1 Level-Based Signals#

DimensionSenior (L5)Staff (L6)Principal (L7)
Time semanticsWindows by arrivalEvent time, watermark from measured delay, allowed lateness, side outputsOrg-wide timestamp contract; producers own event_time accuracy with SLAs
Correctness"Exactly-once enabled"End-to-end chain with idempotent sink; names where it breaksBilling-grade vs operational-grade tiers with different platforms, SLOs and audit
StateUses keyed stateSizes state, TTLs, checkpoint budget, rescale strategyPrices state + retention as recovery budget; capacity model per tier
Reprocessing"Rerun the job"Shadow output + diff + alias swap; retention sized for itSelf-service backfill product; restatement policy with finance/legal
OperationsConsumer lag alertWatermark lag, checkpoint health, late-data, diff vs batch — with ownersError budgets per pipeline tier; game days for replay
OwnershipPipeline team owns allProducers own schemas/timestamps; platform owns engine; consumers own latenessData contracts as org standard; platform-as-product with templates

5.2 Strong Hire Signals#

SignalWhat It Sounds Like
Measured watermark"I set watermark at p99 delay (30 s) and lateness at p99.9 (5 min), measured per producer from the lake."
Idempotent sink"Upsert absolute values by (campaign, window_start); replays overwrite, never add."
Watermark lag alert"Consumer lag can be zero while the watermark is stuck — alert on both."
Reconciliation"Batch is truth for invoices; stream-vs-batch diff > 0.1% pages."
Finality in the API"Every aggregate carries is_final; consumers know what can change."

5.3 Lean No-Hire Signals#

SignalWhy It Misses the Bar
Processing-time windows for billingResults depend on outages
"Exactly-once" with an HTTP sinkDoesn't understand end-to-end boundary
No late-data policySilently wrong numbers
No replay planCan't fix bugs found after the fact
Unbounded stateJob dies weeks later

5.4 Common False Positives#

  • Flink API fluency ≠ streaming judgment — knowing ProcessFunction doesn't answer "when is a number final?"
  • Lambda vs Kappa debate — reciting the religion without a cost argument is a Senior move.
  • Exotic windows (session merging, CEP) when the problem needs tumbling counts.
  • Choosing Spark vs Flink vs Kafka Streams at length — the engine rarely decides level.

6. Interview Flow & Pivots#

6.1 Typical 45-Minute Shape#

PhaseTimeGoal
Framing0–3 minDecision vs record; freshness, correctness bar, lateness
Entities + API3–5 minEvents, window aggregates with is_final
High-level design5–10 minKafka, dedup, event-time windows, upsert sink, lake + batch
Transition10–11 minOffer time / exactly-once / state & backfill
Deep dives11–40 minWatermarks & late data → sink correctness → skew & state → backfill
Wrap-up40–45 minEvolution, ownership, what not to build

6.2 How Interviewers Pivot — And What They're Testing#

Interviewer SaysWhat They're TestingWhere to Go
"Events arrive hours late"Completeness policyLateness, side output, batch correction
"The job crashed"Exactly-once reasoningCheckpoint + replay + idempotent sink
"We found a bug from last week"ReprocessingRetention, shadow run, alias swap
"One key is 40% of traffic"SkewTwo-phase aggregation
"Make it sub-second"Latency tradeoffsEarly triggers, drop lateness, processing-time timers
"Why not just batch?"JudgmentWhen NOT to stream

6.3 What to Deliberately Skip#

TopicWhy L5 Goes HereWhat L6 Says Instead
Kafka broker internalsFamiliar"RF=3, min ISR 2, acks=all."
Chandy–Lamport algorithm detailSounds impressive"Aligned barriers snapshot state + offsets consistently."
Window type catalogEasy"Tumbling 1 min; sliding would cost 5× state for no user benefit."
Serving-store internalsComfortable"Any KV/OLAP supporting upsert and range by campaign."

6.4 Follow-Up Questions to Expect#

  1. "How do you choose the watermark delay, and what happens when it's wrong?"
  2. "Walk me through a crash between writing to the sink and checkpointing."
  3. "How do you change the aggregation logic without losing state?"
  4. "How do you backfill 30 days after a bug?"
  5. "How do you handle a hot key?"
  6. "How do you know the numbers are right?"
  7. "How do you rescale from 64 to 128 parallelism?"

7. Active Drills#

Drill 1: The Opening#

Prompt: "Design a system that counts ad clicks per campaign in real time."

Staff Answer

"First: who uses the counts, and what's the cost of a wrong number versus a late one? If it's budget pacing, I need ~1-minute freshness and can tolerate small, self-correcting errors. If it's invoicing, I need ≤ 0.1% error and can wait a day. I'll design for pacing with a batch reconciliation feeding invoices.

Kafka keyed by campaign; Flink dedupes by click_id; event-time 1-minute tumbling windows with a watermark from measured delay — say 30 s — and 5 minutes of allowed lateness; results upsert into a serving store by (campaign, window_start) with an is_final flag. Events later than 5 minutes go to the lake. Daily batch over the lake is the invoice source of truth, with a diff alarm at 0.1%."

Why this is L6:

  • Separates decision (pacing) from record (invoice) and gives each a correctness bar
  • Event time and lateness with numbers, not adjectives
  • Reconciliation designed in, not bolted on

What L7 adds:

  • Classifies the output as "billing-grade" under an org tiering standard, which dictates reconciliation and audit requirements
  • Asks whether the click stream should be a shared, contracted dataset other teams (fraud, reporting) consume rather than re-ingest

Drill 2: Watermarks Made Concrete#

Prompt: "What exactly is your watermark and how do you choose it?"

Staff Answer

"Per source partition, the watermark is max observed event time minus a bounded out-of-orderness — 30 s. The operator watermark is the min across inputs, excluding partitions idle for 60 s. A window [10:05, 10:06) fires when the watermark passes 10:06.

I choose 30 s from data: the lake records event_time and ingest_time, so I compute delay percentiles per producer. Server-side clicks: p99 = 4 s. Mobile: p99 = 25 s, p99.9 = 4 min. 30 s covers p99 for both; 5-min lateness covers p99.9 with upserts; the rest goes to batch. I alert when late_events_side_output_total exceeds 0.1% of volume — that means the distribution moved and the watermark needs revisiting."

Why this is L6:

  • Mechanism precise (min across partitions, idleness)
  • Parameters from measured distributions per producer
  • Alert tied to the assumption the watermark encodes

What L7 adds:

  • Makes producer delay a contract: each producer publishes a delay SLO; breaches are the producer's incident, not the pipeline's

Drill 3: Exactly-Once, End to End#

Prompt: "Your job crashes after writing to the database but before the checkpoint. What happens?"

Staff Answer

"On restart, Flink restores state and Kafka offsets from the last checkpoint, so events after it are replayed and the window results they produce are written again. Whether that's correct depends on the sink. With an append of deltas, it double-counts. With an upsert of the full window value keyed by (campaign, window_start), the replay recomputes the same value from the same inputs and overwrites — correct. That's why I use upserts.

If the output were a Kafka topic, I'd use the transactional sink: writes go into a transaction committed only on checkpoint completion, and downstream reads read_committed. Cost: output visible only every checkpoint interval (30 s). For side effects like 'pause campaign' calls, the call carries an idempotency key campaign:window:version, and the pacing service dedupes."

Why this is L6:

  • Explains the replay mechanically
  • Picks sink semantics that make replay safe
  • Handles side effects separately

What L7 adds:

  • Standardizes "sink contracts" in the platform: only upsert, transactional or idempotency-keyed sinks are allowed for billing-grade jobs; others are labeled at-least-once in the catalog

Drill 4: Kafka or Checkpoint Store Is Down#

Prompt: "S3, where checkpoints live, is degraded for 30 minutes."

Staff Answer

"Processing continues — checkpoints fail but the job keeps running. The risk grows with time: if a TaskManager dies, we restore from the last successful checkpoint, maybe 30 minutes old, and replay 30 minutes. That's fine only if Kafka retention covers it (it does, 7 days) and the sink is idempotent (it is). I'd set tolerable-failed-checkpoints so the job doesn't fail itself because of checkpoint failures, alert on last_successful_checkpoint_age_s > 5 min, and avoid deploys until checkpoints recover. If Kafka itself is down, producers buffer (bounded, e.g., 5 min in memory, then spill or drop per policy); the job idles; watermarks stall — pacing must then fall back to a conservative mode, e.g., throttle campaigns within 10% of budget."

Why this is L6:

  • Distinguishes degraded durability from outage
  • Ties replay safety to retention and sink design
  • Defines a downstream degraded mode

What L7 adds:

  • Makes "pipeline stale → consumer degraded mode" a required section in every consumer's design review

Drill 5: Hot Key#

Prompt: "One campaign is 40% of all clicks. Your job lags."

Staff Answer

"Keyed by campaign, all 400K events/s land on one subtask — capped at ~100–200K/s. More parallelism doesn't help. Two-phase aggregation: first key by (campaign_id, hash(click_id) % 32) and pre-aggregate per minute; then key by campaign_id and sum 32 partials. The hot campaign's work spreads across 32 subtasks; the merge stage sees 32 records per window instead of 24M. Dedup must happen before salting, or be keyed by click_id itself. I'd apply salting only to keys detected hot (top 100 by rate) to avoid 32× records for the long tail."

Why this is L6:

  • Knows parallelism doesn't split a key
  • Designs correct two-phase merge with dedup ordering
  • Limits overhead to hot keys

What L7 adds:

  • Provides skew handling as a platform operator so teams don't hand-roll it inconsistently

Drill 6: Multi-Tenant Streaming Platform#

Prompt: "40 teams want to run streaming jobs. How do you host them?"

Staff Answer

"Isolation first: per-job Flink clusters (application mode) on Kubernetes, so one team's OOM or checkpoint spiral can't take down another's. Shared: Kafka, schema registry, checkpoint storage with per-job prefixes and quotas, and standard metrics. Admission: jobs declare throughput, state size and tier; the platform sizes resources and sets alerts. Quotas on Kafka produce/consume per team. Billing-grade jobs get dedicated node pools and stricter change control."

Why this is L6:

  • Isolation by default; shared only where it's cheap to share
  • Declared resource contracts drive sizing and alerts

What L7 adds:

  • Chargeback per job; a deprecation policy for unowned jobs (no owner → paused after 30 days)

Drill 7: Build vs Buy#

Prompt: "Should we run Flink ourselves or use a managed service?"

Staff Answer

"Self-hosting Flink well needs ~2–4 engineers who understand checkpoints, RocksDB tuning and upgrades. For under ~20 jobs and no in-house expertise, a managed Flink or streaming SQL service is cheaper in total even at a 30–50% compute premium. Self-host when we have 50+ jobs, need custom operators or versions, or cost at scale dominates. Either way, the job code and state format should stay portable — avoid proprietary SQL extensions for billing-grade logic."

Why this is L6:

  • Frames in headcount and premium, not features
  • Protects portability

What L7 adds:

  • Negotiates exit terms (savepoint portability, data egress) and sets a 2-year review trigger

Drill 8: Changing Logic Without an Outage#

Prompt: "Product wants to change how spend is computed — exclude invalid clicks. Deploy it."

Staff Answer

"Is the state compatible? Adding a filter before aggregation doesn't change state schema, so: take a savepoint, deploy v2 from it — continuous results. But results change meaning at a point in time; consumers must know. For comparability, I'd run v2 in shadow from the same savepoint writing to a shadow table for 24 h, diff against v1 and batch, get product sign-off on the expected delta (e.g., −2.3%), then swap the read alias. If historical numbers must also change, that's a backfill: replay the retention window into the shadow table. If state is incompatible (new key), start v2 from Kafka offsets at a known time and run both until v2 catches up."

Why this is L6:

  • Savepoints + shadow + diff + alias swap
  • Treats semantic change as a product event

What L7 adds:

  • Versioned metric definitions in a catalog; every restated metric carries its definition version

Drill 9: Cost#

Prompt: "The streaming bill is $180K/month. Cut it."

Staff Answer

"Find the big levers: (1) state — dedup over 24 h is 160 GB; 99.9% of duplicates occur within 10 min, so a 10-min dedup window plus batch cleanup cuts state 100× and checkpoint I/O with it. (2) Over-provisioning — jobs sized for peak all day; autoscale on lag via savepoint rescale off-peak. (3) Sliding windows — 5-min windows sliding every 10 s means 30× assignments; switch to tumbling + client-side smoothing. (4) Kafka retention — tiered storage for days 2–7. (5) Kill streaming for dashboards nobody views more than daily. I'd expect 35–50%."

Why this is L6:

  • Levers tied to design choices with multipliers
  • Questions whether streaming is needed at all

What L7 adds:

  • Per-job cost in the catalog next to its consumers; jobs without an active consumer are auto-flagged

Drill 10: Multi-Region#

Prompt: "Clicks are served from 3 regions. Aggregate globally."

Staff Answer

"Aggregate regionally, merge globally. Each region runs the same job over its local Kafka, producing per-region window aggregates — no cross-region traffic on the hot path. A global job consumes the three regional aggregate topics (replicated with MirrorMaker or similar) and sums per window; its watermark is the min of regional watermarks, so a lagging region delays global finality — alert on it. Pacing uses local + last-known remote totals, accepting ~1–2 min of staleness for remote spend. If a region is cut off, the global job marks windows 'partial' rather than waiting forever."

Why this is L6:

  • Keeps raw events regional
  • Understands watermark propagation across regions
  • Defines partial-result behavior

What L7 adds:

  • Aligns with data residency: raw events never leave region; only aggregates cross, approved by legal

8. Deep Dive Scenarios#

Deep Dive 1: Peak-Traffic Lag Spiral#

Context: Black Friday. Click volume is 2.4× normal. The aggregation job's lag hits 25 minutes and climbing; checkpoints are timing out; budget pacing is flying blind. On-call escalates to you.

Questions to Surface First:

  • Is lag growing because of throughput (all subtasks busy) or skew (one subtask busy)?
  • Are checkpoints failing because of backpressure or storage?
  • What does pacing do when data is stale — does it fail open (keep serving ads) or closed?
  • How much budget overspend per minute is at stake?

Typical L5 Approach: Increases parallelism and restarts. The restart restores from a 20-minute-old checkpoint, replays 20 minutes at 2.4× load, and makes lag worse. Rescaling without a savepoint may not even be possible while checkpoints fail.

Staff Approach: Protect the business first: switch pacing into degraded mode (throttle campaigns above 85% of budget) because stale spend is worse than conservative pacing. Diagnose skew vs throughput via subtask busy time. If throughput: enable unaligned checkpoints to get one successful checkpoint, then savepoint and rescale 2×. If skew: enable hot-key salting. Don't restart from an old checkpoint under load.

Principal Approach: The job was sized for an average day in a business with a known peak calendar. Institute peak-readiness reviews for billing-grade pipelines (load test at 3× the last peak, 4 weeks prior), and make "stale data → consumer degraded mode" a contract every consumer implements. Budget: ~20% extra capacity for peak weeks is cheaper than one hour of uncapped ad spend.

Staff Approach — Full Reasoning
PhaseWhat to Do
Immediate (0–5 min)Pacing → conservative mode. Check subtask.busy_time_pct distribution and checkpoint_duration_ms.
TriageAll subtasks > 90% busy → throughput. One subtask pegged → hot key.
Quick fixUnaligned checkpoints → one success → savepoint → redeploy at 2× parallelism (5–10 min).
GuardrailsWatch consumer_lag_seconds slope; if still growing after rescale, shed the non-billing side outputs.
Post-mortemWhy was capacity 1.5× not 3× for a known peak? Why did pacing lack a staleness fallback?

Metrics to Watch: consumer_lag_seconds, watermark_lag_seconds, checkpoint_duration_ms, subtask.busy_time_pct, pacing.data_age_seconds

Organizational Follow-up: Peak calendar shared with the streaming platform; pre-scaled jobs; pacing team owns its degraded mode.

Ownership Question: "Who decides pacing goes conservative, costing advertisers delivery?" Staff answer: The pacing on-call, via a pre-approved runbook triggered automatically when pacing.data_age_seconds > 180; ads business owners agreed in advance that under-delivery beats over-delivery.

Key Takeaway: "Never restart a lagging stateful job from an old checkpoint under peak load. Get a fresh checkpoint first, then rescale."

What clears the Staff bar:

  • Protects the downstream decision before fixing the job
  • Distinguishes skew from throughput
  • Knows unaligned checkpoints and savepoint rescale

Deep Dive 2: Silent Undercount#

Context: Finance notices invoices from the batch path are 2.7% higher than the streaming dashboards advertisers see — for the last 11 days. No alert fired.

Questions to Surface First:

  • Is the diff uniform or concentrated by producer, region or campaign type?
  • Did the diff job exist, and why didn't it alert?
  • Did late-event volume change 11 days ago?
  • Which one is right?

Typical L5 Approach: Assumes the stream is buggy and replays the job over 11 days. May reproduce the same undercount if the cause is lateness policy, not code.

Staff Approach: Slice the diff: it's all from iOS clicks. 11 days ago an iOS SDK release changed batching from 10 s to 5 min; p99 delay jumped from 25 s to 6 min, so ~2.7% of events now exceed the 5-min lateness and go to the side output. The stream is working as designed; the design assumption broke. Fix: raise lateness to 15 min (state +3×, ~2 GB — fine), and alert on the side-output ratio. The diff job existed but its alert was routed to a deleted channel.

Principal Approach: Producer delay is an unmanaged dependency. Require SDK teams to publish delivery-delay SLOs and run a pre-release check against pipeline lateness budgets. Audit every alert route quarterly — "alert to nowhere" is an org failure, not a config typo.

Staff Approach — Full Reasoning
PhaseWhat to Do
ImmediateNotify advertiser-facing teams that dashboards under-report; batch is authoritative.
TriageDiff by dimension → iOS; event_delay_seconds by producer → shift at release date.
Quick fixLateness 5 → 15 min via savepoint redeploy; restate last 11 days from batch into serving store.
GuardrailsAlert: side-output ratio > 0.2%; diff > 0.1% pages the pipeline on-call (verified route).
Post-mortemContributing: SDK change without downstream review; alert routing untested.

Metrics to Watch: late_events_side_output_ratio, event_delay_seconds{producer} p99, reconcile.stream_vs_batch_diff_pct

Organizational Follow-up: SDK release checklist includes "delivery delay distribution change"; alert routing tested by synthetic fire monthly.

Ownership Question: "Who owns the 11 days of wrong dashboards?" Staff answer: The pipeline team owns detection and restatement; the SDK team owns the regression; advertiser comms is owned by the ads product team with finance approving any credits.

Key Takeaway: "A watermark is an assumption about producers. When producers change, the assumption silently breaks — so alert on the assumption, not just the output."

What clears the Staff bar:

  • Slices the diff before acting
  • Recognizes design-assumption drift vs code bug
  • Tests alert routing

Deep Dive 3: Onboarding a 10× Producer#

Context: A new product (connected-TV ads) will send 3M events/s — 3× current total — starting in 8 weeks, with events arriving up to 2 hours late due to device batching.

Questions to Surface First:

  • Does it need the same freshness and correctness as existing traffic?
  • Can its late data be absorbed without making every other window wait?
  • What's the key cardinality and skew?
  • Who pays for the capacity?

Typical L5 Approach: Adds partitions and parallelism to the existing job and raises allowed lateness to 2 hours for everyone — multiplying state 24× and delaying finality for all products.

Staff Approach: Separate topic and separate job instance (same code) with its own watermark and lateness profile: 5-min watermark, 2-h lateness for CTV, while web/mobile keep 30 s/5 min. Shared serving store with a product dimension. Capacity plan: 3M/s ≈ 20–30 subtasks with headroom; state for 2-h lateness ≈ 120 windows open per key — size it (e.g., 1M keys × 120 × 200 B ≈ 24 GB). Load test at 6M/s.

Principal Approach: One job per lateness profile is the pattern; codify "lateness class" as a producer attribute in the event contract so the platform routes new producers automatically. Chargeback the capacity to the CTV product line.

Staff Approach — Full Reasoning
PhaseWhat to Do
Week 1–2Measure CTV delay distribution from a pilot; set watermark/lateness; capacity model.
Week 3–4Topic with 512 partitions; job instance; shadow ingest of pilot traffic.
Week 5–6Load test 6M/s; checkpoint duration < 30% of interval; restore drill.
Week 7Reconciliation path for CTV in batch; diff threshold.
LaunchRamp 1% → 10% → 100% over a week; kill switch = stop CTV consumption without affecting others.

Metrics to Watch: consumer_lag_seconds{job=ctv}, state_size_bytes{job=ctv}, checkpoint_duration_ms, late_events_side_output_ratio{producer=ctv}

Organizational Follow-up: CTV team signs the event contract (schema, delay SLO, volume forecast).

Ownership Question: "If CTV traffic breaks the job, who's paged?" Staff answer: The CTV pipeline instance pages the aggregation on-call, but because it's isolated, the blast radius is CTV only — which is exactly why it's a separate instance.

Key Takeaway: "Don't let one producer's lateness set everyone's latency. Isolate by lateness profile."

What clears the Staff bar:

  • Isolates by lateness class
  • Sizes state from lateness × cardinality
  • Ramps with a kill switch

Deep Dive 4: Post-Mortem — The Backfill That Doubled Revenue#

Context: An engineer fixed a currency-conversion bug and backfilled 14 days by rerunning the job from old offsets — into the production serving store. Dashboards showed 2× revenue for 14 days; the pacing service throttled hundreds of campaigns. You lead the post-mortem.

Questions to Surface First:

  • Was the sink upsert or append?
  • Why was a backfill allowed to write to production tables?
  • What consumed the corrupted data in the meantime (pacing, exports, emails)?
  • Is there a clean copy?

Typical L5 Approach: Delete and recompute. Adds a "be careful" note to the wiki.

Staff Approach: Root cause: the backfill wrote to an append-mode sink under a different job ID, so its output added to existing rows. Standard: backfills never write to live tables — they write to versioned shadow tables (spend_v2_backfill_20260930), are diffed against batch, and are published by atomically swapping a read alias. Production sink credentials are unavailable to ad-hoc jobs. Pacing ignores restated windows older than 1 hour (it already acted on them).

Principal Approach: Backfill is a product, not a script. Provide a self-service backfill tool with shadow output, automatic diff and approval gate. Declare restatement policy: billing-grade metrics may be restated for 30 days with finance approval; beyond that, corrections go into an adjustments ledger instead of rewriting history.

Staff Approach — Full Reasoning
PhaseWhat to Do
ImmediateStop backfill; pacing to conservative mode; freeze exports.
TriageIdentify affected windows by job ID and write timestamp.
Quick fixRestore affected partitions from batch; rerun backfill into a shadow table; diff; swap alias.
GuardrailsSeparate credentials; alias-based publishing; diff gate.
Post-mortemContributing: append sink, shared credentials, no review for historical writes.

Metrics to Watch: reconcile.stream_vs_batch_diff_pct, sink.rows_written{job}, pacing.throttled_campaigns

Organizational Follow-up: Backfill runbook and tool owned by the platform; restatement policy owned by finance + data.

Ownership Question: "Who approves restating last month's advertiser spend?" Staff answer: Finance, with the pipeline owner attesting to the diff — because restated numbers can change invoices and advertiser trust.

Key Takeaway: "Backfills write beside production, never into it. Publishing is an alias swap after a diff."

What clears the Staff bar:

  • Identifies the append sink as the multiplier
  • Designs shadow + diff + alias swap
  • Defines a restatement policy

Deep Dive 5: Multi-Region Expansion#

Context: The company is adding EU and APAC ad serving. Raw click data from the EU must stay in the EU. Global advertisers want a single global spend view within 2 minutes.

Questions to Surface First:

  • Is it legal to replicate aggregates (non-personal) across regions?
  • What latency does global pacing need?
  • What happens to global finality if one region lags?
  • Where does batch reconciliation run?

Typical L5 Approach: Replicate all raw topics to one central region and run the existing job there — violating residency and adding cross-region dependency to the hot path.

Staff Approach: Regional jobs over regional Kafka produce per-region aggregates; only aggregates replicate to a global merge job. Global watermark = min of regional watermarks with a staleness cap: if a region lags > 3 min, global windows fire with partial=true and update later. Batch reconciliation runs per region over regional lakes; a global batch merges regional batch outputs.

Principal Approach: Define an org-wide classification: raw events are region-bound; aggregates meeting a k-anonymity threshold are global. Legal signs off once for the class, not per pipeline. Budget: roughly 2.5× infrastructure for 3 regions with per-region on-call coverage.

Staff Approach — Full Reasoning
PhaseWhat to Do
Q1Regional Kafka + job instances; regional serving stores.
Q2Aggregate replication; global merge job with partial-window semantics.
Q3Regional lakes + batch; global reconciliation of aggregates.
Q4Region-failure game day: cut APAC off for 30 min; verify partial flags and recovery.

Metrics to Watch: watermark_lag_seconds{region}, global.partial_windows_ratio, replication.lag_seconds, residency.raw_events_cross_region_total (must be 0)

Organizational Follow-up: Data classification for every topic; legal review of aggregate schemas.

Ownership Question: "Who decides global pacing uses partial data during a regional outage?" Staff answer: The pacing product owner, pre-agreed: partial data plus conservative throttling for global campaigns.

Key Takeaway: "Aggregate where the data lives; move answers, not events."

What clears the Staff bar:

  • Residency-first topology
  • Partial-result semantics for global finality
  • Game day for region loss

9. Level Expectations Summary#

After studying this case study, you should be able to:

  • Separate the decision a stream drives from the record it creates, and set a correctness bar for each
  • Explain event time, watermarks, allowed lateness and triggers with measured parameters
  • Defend exactly-once end to end, including sink semantics and side effects
  • Size keyed state and choose checkpoint settings; rescale via savepoints
  • Handle skew with two-phase aggregation
  • Design backfills as shadow + diff + alias swap, with retention sized for recovery
  • Reconcile streaming output against a batch source of truth, with owners for the diff

The Bar for This Question#

Mid-level (L4): Knows Kafka and a stream engine; builds a consumer that counts per key. Doesn't distinguish event and processing time.

Senior (L5): Uses event-time windows and knows watermarks exist; enables exactly-once; handles scaling with partitions. Late data and replay handled when asked, often with "drop it" or "rerun."

Staff+ (L6): Starts from the consumer's tolerance for late vs wrong. Sets watermark and lateness from data, designs idempotent sinks, builds reconciliation and backfill paths, sizes state, handles skew, and assigns ownership for schemas, timestamps and lateness policy. The interviewer should learn something from the answer.


10. Staff Insiders: Controversial Opinions#

10.1 "Exactly-Once Is a Property of Your Sink, Not Your Engine"#

ComponentWhat the engine guaranteesWhat you must add
SourceReplay from offsetsRetention ≥ recovery window
StateConsistent snapshotsTTLs, compatible upgrades
SinkNothing (or 2PC for some)Upserts / idempotency keys
Side effectsNothingReceiver-side dedup

The Staff position: Say "effectively-once to the serving store via upserts" instead of "exactly-once." It's more honest and more defensible.

Why this matters in interviews: The follow-up "your sink is an HTTP API" ends most exactly-once answers.

10.2 "Most Real-Time Dashboards Should Be Batch Jobs"#

If the output is viewed daily, streaming buys cost and pager load for nothing. A 15-minute scheduled SQL query is complete, cheap and debuggable.

The Staff position: Stream when a machine acts on the output within minutes. Humans reading dashboards rarely need it.

10.3 "Kappa Is Right Until Finance Asks How You Know"#

Replaying the stream is the best way to fix numbers; an independent computation is the only way to check them. Keep one logic definition, but run it twice for billing-grade outputs.

10.4 "The Watermark Is a Contract With Producers You Never Signed"#

Your completeness depends on producer delay distributions you don't control. Until producers publish delay SLOs, your pipeline's correctness is on someone else's roadmap.

10.5 "State Is a Database — Operate It Like One"#

Stateful jobs have schemas, migrations, backups (checkpoints), restores and capacity plans. Teams that treat them as stateless services discover this during their first incompatible upgrade.


11. The Principal Lens (L7)#

Why L7 Sees This Problem Differently#

A Staff engineer builds a correct aggregation pipeline. A Principal engineer sees 120 streaming jobs across 25 teams, each re-ingesting the same click stream, each with its own watermark heuristics, dedup logic and definition of "valid click" — producing five different spend numbers for the same campaign. The L7 problem is metric integrity across the org: shared event contracts, a small number of certified pipelines for money-grade metrics, and a platform that makes the right thing the default.

The Org-Level Fault Line#

Certified shared datasets vs team-owned derivations. Central, certified streams (one canonical "valid clicks" stream, one spend metric) ensure consistency but slow teams that need variants. Team-owned derivations are fast but multiply definitions and compute.

The L7 default: Certify a thin layer — canonical, deduped, validated event streams and the handful of money-grade metrics — owned by a data platform team with SLOs. Everything above is team-owned but must declare lineage to certified sources. Certification is a tier, not a gate.

Cost Model#

Assumptions: ~1 KB events; managed Kafka/self-hosted Flink on cloud VMs; ~$0.05/vCPU-hour; object storage for checkpoints and tiered log; engineer $25K/month fully loaded.

ScaleThroughputInfra ($/month)Platform HeadcountOn-call LoadNotes
Small50K events/s, 10 jobs~$15–25K1–2 FTE (shared)~2 pages/monthConsider managed Flink/streaming SQL; batch for most dashboards
Medium1M events/s, 100 jobs~$150–250K5–8 FTE platform~10 pages/monthTemplates, per-job clusters, reconciliation for billing tier
Large20M events/s, 1,000 jobs~$1.5–3M20–30 FTEDedicated rotations by tierCertified datasets, self-service SQL, chargeback

Biggest hidden costs: dedup/window state (checkpoint I/O scales with it), Kafka retention for replay, and duplicate ingestion by teams re-consuming the same raw topics. Consolidating five teams' parallel click-dedup jobs into one certified stream typically saves more than any tuning.

The 3-Year Evolution Path#

Diagram: The 3-Year Evolution Path

One-Way Doors vs Two-Way Doors#

DecisionDoorReversal Cost
Event schema + event_time semantics in producer contractsOne-wayEvery producer and consumer must migrate
Kafka partition key for core topicsOne-way-ishRepartition + reorder guarantees; weeks
Log retention (shortening)One-way for lost dataData gone; can't replay
Stream engine choiceTwo-way (painful)Rewrite jobs; state migration via replay
Window size / latenessTwo-waySavepoint redeploy; restatement
Sink store choiceTwo-wayDual-write + backfill; weeks
Metric definitions exposed to customers/invoicesOne-wayContractual and trust cost

The Standard I'd Write#

RFC-STREAM-003: Streaming Pipelines and Metric Integrity

Scope: All streaming jobs whose outputs are consumed by another team, a customer, or a financial process.

MUST:

  • Use producer-set event_time and event-time windows for any user-facing or billing output.
  • Declare watermark delay and allowed lateness, justified by measured producer delay (p99 / p99.9).
  • Write outputs via upsert, transactional or idempotency-keyed sinks; no append-delta sinks for aggregates.
  • Expose is_final (or equivalent) for outputs that can be revised.
  • Alert on consumer_lag_seconds, watermark_lag_seconds, checkpoint age and late-event ratio.
  • Billing-grade outputs: daily reconciliation against batch with diff ≤ 0.1%, owned by a named team.

SHOULD:

  • Consume certified event streams rather than raw topics.
  • Retain replayable input ≥ 7 days hot, ≥ 30 days tiered.
  • Perform backfills into shadow outputs published by alias swap.

Exceptions: Fraud/abuse jobs may use processing time and at-least-once with documented rationale. Exceptions reviewed by the Data Platform Council within 5 business days.

Success metrics: Zero unreconciled billing restatements; ≥ 90% of jobs on templates; median time-to-detect for pipeline correctness incidents < 1 hour; duplicate raw-topic consumers reduced 50% in a year.

What I'd Tell the VP#

"Our real-time numbers drive ad spend decisions worth millions per day, and today five teams compute 'spend' five different ways. Twice this year a streaming bug went unnoticed for over a week. I recommend we certify one source of truth for money-related metrics, check it daily against our batch system, and give teams a shared platform with safe defaults instead of bespoke pipelines. It costs roughly six engineers and should reduce our streaming infrastructure bill by around a quarter by eliminating duplicate processing. Most importantly, when an advertiser disputes an invoice, we'll be able to show exactly how the number was produced."

Principal Interview Signals#

SignalWhat It Sounds Like
Metric integrity framing"The problem isn't this job — it's five definitions of spend. I'd certify one."
Tiers correctness"Billing-grade and operational-grade metrics get different platforms, SLOs and review."
Prices retention as recovery"Seven days of hot retention is ~$20K/month; it's our insurance against late-found bugs."
Contracts with producers"Delay SLOs are part of the event contract; watermark breaches are the producer's incident."
Knows when not to stream"Half these jobs feed dashboards viewed daily; I'd retire them to batch."

Staff answers that L7 interviewers find insufficient:

  • A perfect watermark design that ignores that three other teams compute the same metric differently.
  • "We'll build a streaming platform" without the tiering, chargeback or retirement policy that keeps it from becoming 1,000 orphaned jobs.
  • Backfill tooling without a restatement policy agreed with finance.

Appendices

Appendix A: Windowing and Watermark Mechanics

A.1 Window Assignment#

tumbling(size=60s):   window_start = floor(event_time / 60s) * 60s
sliding(size=300s, slide=60s): event belongs to 5 windows → 5× state and output
session(gap=30min):   windows merge when a new event bridges two sessions (late data can merge sessions)

A.2 Watermark Generation#

per partition:  wm_p = max_event_time_seen_p − bounded_out_of_orderness (30 s)
operator:       wm = min(wm_p for p in active partitions)   # idle > 60 s excluded
fire window W when wm ≥ W.end
accept late event e for W while wm < W.end + allowed_lateness; re-fire W
drop / side-output when wm ≥ W.end + allowed_lateness

A.3 Triggers and Refinement Modes#

ModeLate firing emitsSink requirement
DiscardingOnly the new deltaSink must add (dangerous on replay)
AccumulatingFull new totalSink must upsert
Accumulating + retracting−old, +newDownstream aggregations stay correct

Default: accumulating + upsert sink.

Appendix B: Keys, Events and Contracts

B.1 Event Envelope#

{ event_id (uuid, dedup key), event_type, schema_version,
  event_time (producer clock, ms), ingest_time (broker),
  producer_id, partition_key, payload{...} }

B.2 Partition Keys#

Key by the aggregation key (campaign_id) to avoid a shuffle, unless skew dominates — then key by event_id for even ingestion and let the engine shuffle.

B.3 Clock Hygiene#

Producer clocks drift; reject or clamp events with event_time > ingest_time + 5 min (future-dated) — a single device with a clock set to 2030 can advance watermarks and cause everything else to be "late."

Appendix C: State and Checkpointing

C.1 Backends#

BackendMax practical stateAccessCheckpoint
Heap~10s of GB per TMnsFull
RocksDBTBsµs (SSD)Incremental
External KVUnlimited0.5–2 msN/A (store's durability)

C.2 Checkpoint Tuning#

Interval 30–60 s; timeout 10 min; min pause between checkpoints ≥ 50% of interval; unaligned checkpoints when backpressure is common; local recovery on for fast restarts.

C.3 Rescaling#

Keyed state is partitioned into key groups (max parallelism, e.g., 4,096 — a one-way door like logical shards). Rescale = savepoint → redeploy with new parallelism ≤ max parallelism.

C.4 Quick Comparison — Delivery Mechanisms#

MechanismGuaranteeLatency CostComplexity
Offsets committed before processingAt-most-onceNoneLow
Offsets after processingAt-least-onceNoneLow
Checkpoint + upsert sinkEffectively-onceNoneLow–medium
Checkpoint + 2PC sinkExactly-once≥ checkpoint intervalHigh
Appendix D: Output Contract and Consumer Behavior

D.1 Output Schema#

(campaign_id, window_start) PK, clicks, spend_micros, is_final, version, computed_at, pipeline_version

D.2 Consumer Rules#

  • Treat non-final windows as provisional; don't cache forever.
  • Idempotent on (key, window_start, version).
  • Staleness fallback: if max(computed_at) older than threshold, enter degraded mode.
Appendix E: Observability

E.1 Core Metrics#

Freshness:    consumer_lag_seconds, watermark_lag_seconds, output_age_seconds
Completeness: late_events_accepted_total, late_events_side_output_total, event_delay_seconds{producer}
Correctness:  reconcile.stream_vs_batch_diff_pct, dedup_hits_total, deserialization_errors_total
Health:       checkpoint_duration_ms, checkpoint_failures_total, last_successful_checkpoint_age_s,
              state_size_bytes, subtask.busy_time_pct, backpressure_ratio

E.2 Critical Alerts#

AlertThresholdSeverity
Consumer lag> 5 min for 5 minPage
Watermark lag> 2× watermark delay + 2 minPage
Checkpoint age> 3× intervalPage
Stream vs batch diff> 0.1% (billing)Page
Late side-output ratio> 0.2%Ticket
Deserialization errors> 0 sustained 5 minPage producer + pipeline

E.3 Debugging "Numbers Look Low"#

  1. Watermark stalled? 2) Late ratio up? 3) Dedup over-matching (key collision)? 4) Filter change deployed? 5) Upstream volume actually down? Check in that order — cheapest first.
Appendix F: Scale Evolution

F.1 What Works at Each Scale#

ScaleApproach
< 10K events/sConsumer + DB upserts, or scheduled SQL
10K–1M events/sKafka + Flink/Kafka Streams, templates
1M–20M events/sPlatform with per-job clusters, two-phase aggregation, tiered storage
Multi-regionRegional jobs, aggregate replication, partial finality

F.2 What You Don't Build on Day One#

  • A custom stream engine
  • Session windows and CEP unless the product needs them
  • Real-time invoicing
  • Cross-region raw event replication
Appendix G: Multi-Tenancy, Fairness and Cost

G.1 Tiers#

TierIsolationChange ControlReconciliation
Billing-gradeDedicated node poolReview + shadow runDaily diff ≤ 0.1%
Product-criticalPer-job clusterCanary via savepointWeekly spot checks
ExploratoryShared session clusterSelf-serviceNone; auto-expire in 30 days

G.2 Fairness on Shared Kafka#

Per-team produce/consume quotas; per-topic retention budgets; chargeback by bytes in, bytes retained and vCPU-hours.

  1. Loading the index…