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.
| Mode | Time | What to Read |
|---|---|---|
| Quick Review | 15 min | Executive Summary → Interview Walkthrough → Fault Lines table → Active Drills 1–3 |
| Targeted Study | 1–2 hrs | Executive Summary → Walkthrough → Section 3 (Fault Lines) → Section 4 (Failures) → weak-spot Deep Dives |
| Deep Dive | 3+ hrs | Everything, 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
| Concept | How It Works | Pros | Cons |
|---|---|---|---|
| Tumbling window | Fixed, non-overlapping (e.g., each minute) | Simple, one result per key per window | Boundary effects; bursty emission |
| Sliding / hopping window | Fixed size, advancing by a smaller step (5 min every 1 min) | Smooth trends | Each event in size/step windows → 5× state/work |
| Session window | Gap-based; closes after N min of inactivity | Models user behavior | Unbounded length; merges on late data |
| Global window + triggers | One window per key; custom emission | Flexible (running totals) | You own all correctness logic |
| At-most-once | Commit offsets before processing | Lowest latency | Loses data on failure |
| At-least-once | Commit after processing | No loss | Duplicates on replay |
| Exactly-once (effectively-once) | Checkpointed state + offsets, transactional/idempotent sinks | Correct counts | Latency ≈ 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#
| Behavior | Senior (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?" |
| Time | Processing-time windows | Event-time windows with watermarks; states allowed lateness with a number | Standardizes 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 breaks | Defines which metrics are "billing-grade" (reconciled) vs "operational-grade" (approximate), with different platforms and SLOs |
| Late data | Drops it | Allowed lateness + retractions/upserts + side output + batch reconciliation | Sets 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 strategy | Prices state and replay capacity; retention on the log as a recovery budget |
| Ownership | Pipeline team owns everything | Producers own schemas + timestamps; platform owns engine; consumers own lateness policy | Redraws 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#
| Position | Rationale |
|---|---|
| Event time, always, for anything user-facing or billed | Processing time makes results depend on outages and replays |
| Idempotent sinks over transactional magic | Upsert by (key, window) makes replay safe regardless of engine guarantees |
| State allowed lateness explicitly (e.g., 5 min) and route the rest | Waiting forever is unbounded state; dropping silently is wrong numbers |
| Batch is the source of truth for money | Stream for speed, batch (or re-stream) for reconciliation; publish the diff |
| Keep the log long enough to replay | Kafka retention ≥ 7 days (or tiered storage) is your recovery budget |
| Pre-aggregate at the edge of the hot key | Two-phase (local then global) aggregation for skewed keys |
| Version stateful jobs like databases | Savepoints, state schema evolution, blue/green for incompatible changes |
The Three Intents#
| Intent | Constraint | Strategy | Failure Mode | Correctness Bar |
|---|---|---|---|---|
| Real-time analytics / billing aggregation | Counts must reconcile with money | Event-time windows, idempotent upserts, batch reconciliation | Double/under-counting on replay or late data | ≤ 0.1% diff vs batch; finalized windows |
| Low-latency detection / alerting (fraud, anomaly) | Decide in < 1 s | Per-key stateful rules/CEP, processing-time timers where needed | False negatives during lag; alert storms on replay | Latency SLO p99 < 1 s; recall over precision |
| Continuous materialized views / CDC propagation | Downstream view mirrors source | Ordered per-key changelog, upsert/delete semantics | Out-of-order updates overwrite newer state | Per-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 Line | The Tension |
|---|---|---|
| 1 | Latency vs Completeness | Emit fast with incomplete data, or wait for stragglers? (watermarks, allowed lateness) |
| 2 | Exactly-Once vs Throughput & Simplicity | Transactional sinks and aligned checkpoints, or at-least-once + idempotency? |
| 3 | Local State vs External State | Embedded RocksDB state (fast, hard to rescale) or remote store (simple, slow)? |
| 4 | Streaming-Only vs Stream + Batch Reconciliation | One code path (Kappa), or a batch source of truth (Lambda-ish)? |
| 5 | Platform vs Team-Owned Pipelines | Shared 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.
Uber / Netflix — Flink at Scale#
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#
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#
| Topic | The L5 Answer | The 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#
| Metric | Value | Why 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/s | 1M events/s ≈ 10–20 slots before headroom |
| Checkpoint interval | 10 s – 5 min (30–60 s typical) | Transactional sinks commit at checkpoints → end-to-end latency floor |
| Checkpoint duration alert | > 50% of interval | Backpressure or state growth; risk of never completing |
| Watermark delay | 5 s – 2 min typical | The price of completeness paid in latency |
| Allowed lateness | 1 min – 24 h | Every unit of lateness is retained state |
| Mobile event delay | p99 minutes, p99.9 hours–days | Why a batch correction path exists |
| Dedup state size | keys/day × ~50–100 B | 2B click IDs/day ≈ 100–200 GB of state |
| Log retention for replay | ≥ 7 days hot; 30–90 days tiered | Recovery budget for bugs found late |
| Stream vs batch diff tolerance | ≤ 0.1% billing, ≤ 1–2% ops dashboards | Defines "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)#
Walk the flow:
- Ad servers stamp
event_timeat the edge and produce to Kafka keyed bycampaign_id. - Flink dedupes by
click_id(at-least-once producers retry). - Event-time tumbling windows fire when the watermark passes
window_end; late events within 5 minutes re-fire an updated result. - Results upsert into the serving store by
(campaign_id, window_start); pacing reads running totals. - 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#
| Mistake | Time Lost | Fix |
|---|---|---|
| Explaining Kafka internals (ISR, segments) | 5–8 min | "Kafka, RF=3, 256 partitions, 7-day retention." |
| Enumerating window types | 3–5 min | Pick one and justify it |
| Designing the dashboard | 3–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#
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#
| Situation | Better Alternative | Why |
|---|---|---|
| Consumers read results daily | Hourly/daily batch (Spark, warehouse SQL) | Cheaper, simpler, naturally complete |
| Freshness need ≥ 15 min | Micro-batch or scheduled SQL on the warehouse | No watermarks, no stateful ops |
| Simple per-event transform, no state | Stateless consumer / serverless function | Engine overhead isn't worth it |
| Exact answers over all history | Batch or OLAP query | Unbounded state in a stream |
| Low volume (< 1K events/s) with a DB already | DB triggers / outbox consumers | One fewer distributed system |
| Team has no one to own stateful jobs | Managed streaming SQL (e.g., cloud-managed Flink/ksqlDB) or batch | Stateful 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#
| Gap | Why It Matters | What to Say |
|---|---|---|
| Lateness distribution | Sets watermark and lateness | "Do we know p99/p99.9 delay per producer? I'll measure it from the lake." |
| Consumer of the output | Sets finality semantics | "Does anyone act on a number that may change?" |
| Duplicate rate from producers | Sets dedup design | "Are producers at-least-once with retries?" |
| Key cardinality and skew | Sets state size and hot-key handling | "How many campaigns? Top campaign's share?" |
| Correction window | Sets restatement policy | "Can yesterday's numbers change? Last month's?" |
| Source of truth | Sets reconciliation | "Is the batch/ledger the truth, or is the stream?" |
2.4 Precise Terminology#
| Term | Precise Meaning | Common Confusion |
|---|---|---|
| Event time | When the event happened (set by producer) | Confused with ingest or processing time |
| Processing time | Wall clock of the operator processing it | Seems simpler; results depend on lag |
| Watermark | Assertion that no events with timestamp ≤ W are expected | Not a guarantee — a heuristic |
| Allowed lateness | How long after the watermark passes a window it still accepts updates | Not the same as watermark delay |
| Trigger | When a window emits (on watermark, early every N s, on late event) | Assumed to be "once at window end" |
| Accumulating vs retracting | Late firing emits the new total vs emits (−old, +new) | Downstream sums double-count with accumulating + append sinks |
| Checkpoint | Consistent snapshot of state + source offsets | Not the same as a savepoint (user-triggered, portable) |
| Exactly-once | Each event affects state/output exactly once as observed | Not "processed once" — replays happen; effects are deduped |
| Backpressure | Slow operator throttles upstream | Symptom, not a failure; lag grows |
| Kappa vs Lambda | Reprocess by replaying the stream vs a separate batch layer | Presented 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.
| Strategy | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Processing-time windows | Lowest latency, trivial | Results depend on lag/replays; wrong after outages | Consumers (silently wrong numbers) |
| Event time, tight watermark (5 s), no lateness | Fast, final on emit | Drops stragglers (~1–5%) | Advertisers / finance (under-count) |
| Event time, 30 s watermark + 5 min lateness + upserts | Converges, ~99.9% captured | Numbers change after first emit; state × 6 windows | Consumers (handle updates) + platform (state) |
| Long watermark (1 h) | Nearly complete | 1-hour latency | Pacing (useless for real-time) |
| Early + on-time + late triggers | Speculative early results, refined | Consumers must understand refinement | Product (UX for provisional) |
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_totalso 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.
| Approach | What Works | What Breaks | Who Pays |
|---|---|---|---|
| At-most-once | Fast | Data loss on crash | Everyone downstream |
| At-least-once, append sink | Simple, fast | Duplicates on replay (e.g., 30 s of events re-added) | Consumers (inflated counts) |
| At-least-once + idempotent upsert sink | Replays overwrite; no coordination | Needs deterministic keys; not for side effects | Sink design (key schema) |
| Two-phase commit sink (Kafka transactions) | True end-to-end EOS into Kafka | Output visible only after checkpoint (30 s+); read_committed consumers; transaction timeouts | Latency budget |
| Side effects with idempotency keys | Safe external calls | Receiver must support dedup | Receiving team |
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.
| Approach | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Heap state | Fastest | GC pauses, bounded by memory | On-call (OOM) |
| RocksDB + incremental checkpoints | TB-scale state, ~µs reads | Restore/rescale time (minutes per 100 GB); compaction tuning | Platform (state ops) |
| External KV (Redis/Cassandra) | Stateless workers, easy rescale | 0.5–2 ms per access; 1M events/s → 1M+ store ops/s; exactly-once harder | Store owners + latency |
| Hybrid: local cache + external | Balance | Cache coherence on rebalance | Job 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.
| Approach | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Stream only | One codebase; replay to fix | No independent check; bugs undetected | Finance (silent errors) |
| Lambda (separate batch + stream code) | Independent truth | Two implementations drift; double maintenance | Pipeline team (2× code) |
| Same logic, two runners (Beam/Flink batch mode, SQL) | One definition, two executions | Engine semantic differences at edges | Platform |
| Stream + batch diff only for billing-grade outputs | Truth where it matters | Diff pipeline to own | Pipeline 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#
| Option | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Each team runs its own Flink cluster | Autonomy | 30 clusters, 30 upgrade schedules, inconsistent checkpointing | Every team's on-call |
| Central platform, teams write jobs | Shared ops, standard alerts | Platform bottleneck; noisy neighbors on shared clusters | Platform team |
| Self-service SQL / templates on platform | Fast onboarding, guardrails | Limited expressiveness | Power users |
| Managed cloud service | No cluster ops | Cost; version lag; lock-in | Budget |
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#
| Failure | Detection Signal | Blast Radius | Mitigation | Owner |
|---|---|---|---|---|
| Stalled watermark | watermark_lag_seconds > 2× delay | All windowed outputs | Idleness config; alert on watermark lag | Platform + pipeline |
| Replay double-count | stream_vs_batch_diff_pct > 0.1% | Windows since last checkpoint | Upsert sinks; recompute | Pipeline owner |
| Checkpoint spiral | checkpoint_duration_ms > 50% interval | Whole job; hours of lag | Unaligned checkpoints; scale via savepoint | Platform |
| Hot key | subtask busy-time skew > 3× | Whole job via backpressure | Two-phase aggregation | Pipeline owner |
| Schema break | deserialization_errors_total | Affected fields/records | DLQ; registry enforcement | Producer team |
| Late surge | late_events_side_output_total | Past windows under-counted | Batch restatement | Pipeline + finance |
| Consumer lag | consumer_lag_seconds > SLO | Freshness | Scale out; shed | Pipeline owner |
| State explosion | state_size_bytes growth > 20%/day | Checkpoints, restores | TTL; reduce dedup horizon | Pipeline owner |
| Kafka retention exceeded before fix | oldest_needed_offset < log_start_offset | Unrecoverable from stream | Rebuild from lake | Platform |
5. Evaluation Rubric#
5.1 Level-Based Signals#
| Dimension | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| Time semantics | Windows by arrival | Event time, watermark from measured delay, allowed lateness, side outputs | Org-wide timestamp contract; producers own event_time accuracy with SLAs |
| Correctness | "Exactly-once enabled" | End-to-end chain with idempotent sink; names where it breaks | Billing-grade vs operational-grade tiers with different platforms, SLOs and audit |
| State | Uses keyed state | Sizes state, TTLs, checkpoint budget, rescale strategy | Prices state + retention as recovery budget; capacity model per tier |
| Reprocessing | "Rerun the job" | Shadow output + diff + alias swap; retention sized for it | Self-service backfill product; restatement policy with finance/legal |
| Operations | Consumer lag alert | Watermark lag, checkpoint health, late-data, diff vs batch — with owners | Error budgets per pipeline tier; game days for replay |
| Ownership | Pipeline team owns all | Producers own schemas/timestamps; platform owns engine; consumers own lateness | Data contracts as org standard; platform-as-product with templates |
5.2 Strong Hire Signals#
| Signal | What 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#
| Signal | Why It Misses the Bar |
|---|---|
| Processing-time windows for billing | Results depend on outages |
| "Exactly-once" with an HTTP sink | Doesn't understand end-to-end boundary |
| No late-data policy | Silently wrong numbers |
| No replay plan | Can't fix bugs found after the fact |
| Unbounded state | Job dies weeks later |
5.4 Common False Positives#
- Flink API fluency ≠ streaming judgment — knowing
ProcessFunctiondoesn'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#
| Phase | Time | Goal |
|---|---|---|
| Framing | 0–3 min | Decision vs record; freshness, correctness bar, lateness |
| Entities + API | 3–5 min | Events, window aggregates with is_final |
| High-level design | 5–10 min | Kafka, dedup, event-time windows, upsert sink, lake + batch |
| Transition | 10–11 min | Offer time / exactly-once / state & backfill |
| Deep dives | 11–40 min | Watermarks & late data → sink correctness → skew & state → backfill |
| Wrap-up | 40–45 min | Evolution, ownership, what not to build |
6.2 How Interviewers Pivot — And What They're Testing#
| Interviewer Says | What They're Testing | Where to Go |
|---|---|---|
| "Events arrive hours late" | Completeness policy | Lateness, side output, batch correction |
| "The job crashed" | Exactly-once reasoning | Checkpoint + replay + idempotent sink |
| "We found a bug from last week" | Reprocessing | Retention, shadow run, alias swap |
| "One key is 40% of traffic" | Skew | Two-phase aggregation |
| "Make it sub-second" | Latency tradeoffs | Early triggers, drop lateness, processing-time timers |
| "Why not just batch?" | Judgment | When NOT to stream |
6.3 What to Deliberately Skip#
| Topic | Why L5 Goes Here | What L6 Says Instead |
|---|---|---|
| Kafka broker internals | Familiar | "RF=3, min ISR 2, acks=all." |
| Chandy–Lamport algorithm detail | Sounds impressive | "Aligned barriers snapshot state + offsets consistently." |
| Window type catalog | Easy | "Tumbling 1 min; sliding would cost 5× state for no user benefit." |
| Serving-store internals | Comfortable | "Any KV/OLAP supporting upsert and range by campaign." |
6.4 Follow-Up Questions to Expect#
- "How do you choose the watermark delay, and what happens when it's wrong?"
- "Walk me through a crash between writing to the sink and checkpointing."
- "How do you change the aggregation logic without losing state?"
- "How do you backfill 30 days after a bug?"
- "How do you handle a hot key?"
- "How do you know the numbers are right?"
- "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
| Phase | What to Do |
|---|---|
| Immediate (0–5 min) | Pacing → conservative mode. Check subtask.busy_time_pct distribution and checkpoint_duration_ms. |
| Triage | All subtasks > 90% busy → throughput. One subtask pegged → hot key. |
| Quick fix | Unaligned checkpoints → one success → savepoint → redeploy at 2× parallelism (5–10 min). |
| Guardrails | Watch consumer_lag_seconds slope; if still growing after rescale, shed the non-billing side outputs. |
| Post-mortem | Why 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
| Phase | What to Do |
|---|---|
| Immediate | Notify advertiser-facing teams that dashboards under-report; batch is authoritative. |
| Triage | Diff by dimension → iOS; event_delay_seconds by producer → shift at release date. |
| Quick fix | Lateness 5 → 15 min via savepoint redeploy; restate last 11 days from batch into serving store. |
| Guardrails | Alert: side-output ratio > 0.2%; diff > 0.1% pages the pipeline on-call (verified route). |
| Post-mortem | Contributing: 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
| Phase | What to Do |
|---|---|
| Week 1–2 | Measure CTV delay distribution from a pilot; set watermark/lateness; capacity model. |
| Week 3–4 | Topic with 512 partitions; job instance; shadow ingest of pilot traffic. |
| Week 5–6 | Load test 6M/s; checkpoint duration < 30% of interval; restore drill. |
| Week 7 | Reconciliation path for CTV in batch; diff threshold. |
| Launch | Ramp 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
| Phase | What to Do |
|---|---|
| Immediate | Stop backfill; pacing to conservative mode; freeze exports. |
| Triage | Identify affected windows by job ID and write timestamp. |
| Quick fix | Restore affected partitions from batch; rerun backfill into a shadow table; diff; swap alias. |
| Guardrails | Separate credentials; alias-based publishing; diff gate. |
| Post-mortem | Contributing: 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=trueand 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
| Phase | What to Do |
|---|---|
| Q1 | Regional Kafka + job instances; regional serving stores. |
| Q2 | Aggregate replication; global merge job with partial-window semantics. |
| Q3 | Regional lakes + batch; global reconciliation of aggregates. |
| Q4 | Region-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"#
| Component | What the engine guarantees | What you must add |
|---|---|---|
| Source | Replay from offsets | Retention ≥ recovery window |
| State | Consistent snapshots | TTLs, compatible upgrades |
| Sink | Nothing (or 2PC for some) | Upserts / idempotency keys |
| Side effects | Nothing | Receiver-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.
| Scale | Throughput | Infra ($/month) | Platform Headcount | On-call Load | Notes |
|---|---|---|---|---|---|
| Small | 50K events/s, 10 jobs | ~$15–25K | 1–2 FTE (shared) | ~2 pages/month | Consider managed Flink/streaming SQL; batch for most dashboards |
| Medium | 1M events/s, 100 jobs | ~$150–250K | 5–8 FTE platform | ~10 pages/month | Templates, per-job clusters, reconciliation for billing tier |
| Large | 20M events/s, 1,000 jobs | ~$1.5–3M | 20–30 FTE | Dedicated rotations by tier | Certified 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#
One-Way Doors vs Two-Way Doors#
| Decision | Door | Reversal Cost |
|---|---|---|
| Event schema + event_time semantics in producer contracts | One-way | Every producer and consumer must migrate |
| Kafka partition key for core topics | One-way-ish | Repartition + reorder guarantees; weeks |
| Log retention (shortening) | One-way for lost data | Data gone; can't replay |
| Stream engine choice | Two-way (painful) | Rewrite jobs; state migration via replay |
| Window size / lateness | Two-way | Savepoint redeploy; restatement |
| Sink store choice | Two-way | Dual-write + backfill; weeks |
| Metric definitions exposed to customers/invoices | One-way | Contractual 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_timeand 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#
| Signal | What 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#
| Mode | Late firing emits | Sink requirement |
|---|---|---|
| Discarding | Only the new delta | Sink must add (dangerous on replay) |
| Accumulating | Full new total | Sink must upsert |
| Accumulating + retracting | −old, +new | Downstream 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#
| Backend | Max practical state | Access | Checkpoint |
|---|---|---|---|
| Heap | ~10s of GB per TM | ns | Full |
| RocksDB | TBs | µs (SSD) | Incremental |
| External KV | Unlimited | 0.5–2 ms | N/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#
| Mechanism | Guarantee | Latency Cost | Complexity |
|---|---|---|---|
| Offsets committed before processing | At-most-once | None | Low |
| Offsets after processing | At-least-once | None | Low |
| Checkpoint + upsert sink | Effectively-once | None | Low–medium |
| Checkpoint + 2PC sink | Exactly-once | ≥ checkpoint interval | High |
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#
| Alert | Threshold | Severity |
|---|---|---|
| Consumer lag | > 5 min for 5 min | Page |
| Watermark lag | > 2× watermark delay + 2 min | Page |
| Checkpoint age | > 3× interval | Page |
| Stream vs batch diff | > 0.1% (billing) | Page |
| Late side-output ratio | > 0.2% | Ticket |
| Deserialization errors | > 0 sustained 5 min | Page producer + pipeline |
E.3 Debugging "Numbers Look Low"#
- 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#
| Scale | Approach |
|---|---|
| < 10K events/s | Consumer + DB upserts, or scheduled SQL |
| 10K–1M events/s | Kafka + Flink/Kafka Streams, templates |
| 1M–20M events/s | Platform with per-job clusters, two-phase aggregation, tiered storage |
| Multi-region | Regional 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#
| Tier | Isolation | Change Control | Reconciliation |
|---|---|---|---|
| Billing-grade | Dedicated node pool | Review + shadow run | Daily diff ≤ 0.1% |
| Product-critical | Per-job cluster | Canary via savepoint | Weekly spot checks |
| Exploratory | Shared session cluster | Self-service | None; 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.