Hiring BarSupport

Design with Flink & Stream Processing — Staff-Level Technology Guide

Technology guide36 min read6 diagrams

Why This Matters#

Stream processing is not a faster batch job. It is a distributed state machine that must keep answering while time itself is unreliable. Events arrive late, out of order, duplicated, and in bursts; the processor must hold gigabytes to terabytes of state, survive machine loss without double-counting, and decide when a window is complete even though the next event might still be in transit from a phone in a tunnel. Apache Flink is the reference implementation of that problem, which is why it's the default answer — and why the interesting questions are never about Flink's API.

It shows up constantly in Staff loops: ad click aggregation, fraud detection, real-time leaderboards, surge pricing, trending topics, metrics rollups, CDC-driven materialized views. The L5 candidate draws "Kafka → Flink → DB" and says "Flink gives exactly-once." The L6 candidate says "event-time tumbling windows of 1 minute, watermark at max-seen minus 10 seconds, late events within 5 minutes update the result via upsert, later than that go to a side output reconciled by the hourly batch job; RocksDB state, 30-second incremental checkpoints to S3, and a transactional Kafka sink so downstream read_committed consumers see each window exactly once." The L7 candidate asks whether the org should be running two pipelines — streaming and batch — that compute the same number two different ways, and who reconciles them when finance sees a 0.3% gap.

The gap between levels is not knowing what a watermark is. It is knowing that correctness in streaming is a negotiated contract about time, completeness, and cost, and that every one of those terms has a business owner.

The L5 → L6 → L7 Contrast#

BehaviorSenior (L5)Staff (L6)Principal (L7)
First move"Kafka into Flink, write results to Redis""Event time or processing time? How late can data be, and what happens to it?""Does this number already exist in batch? If yes, which one is the source of truth, and who reconciles?"
Correctness"Flink is exactly-once"Exactly-once state via checkpoints; end-to-end only with transactional or idempotent sinksDeclares a correctness tier per pipeline: approximate / exactly-once / auditable-with-reconciliation
TimeProcessing time windowsEvent time, watermarks, allowed lateness, side outputs for late dataSets the org-wide lateness contract with product: "dashboards final after 5 minutes, billing final after 24 hours"
State"Flink keeps state"Sizes state (keys × bytes × window), picks RocksDB, TTLs every state, incremental checkpointsTreats state as a migration liability: schema evolution, savepoint compatibility, 3-year upgrade path
Failure"It restarts from checkpoint"Recovery time = restore state + replay since checkpoint; sizes both; backpressure diagnosisPlatform-level: job isolation, shared cluster vs per-job clusters, blast radius of a Flink upgrade
Ownership"Data team owns it"Job owner owns output correctness and lag; platform owns the runtimeStreaming platform with SQL as the paved road; custom Java jobs only with justification
Why "Time" separates levels

Processing-time windows are simpler and have lower latency, and for "requests per second on this server" they are correct. But for "clicks per ad per minute" they are wrong: a mobile client that batches events and sends them 40 seconds later gets its clicks counted in the wrong minute, and a consumer restart that replays 10 minutes of Kafka puts 10 minutes of clicks into one window. The Staff answer switches to event time and immediately confronts the consequence — you must now decide when a window is "done," and that decision (the watermark) trades latency against completeness. Naming who accepts the incompleteness is what makes it Staff.

Why "Correctness" separates levels

Flink's checkpoints give exactly-once state: after a failure, operator state is restored as if every input was processed once. The output is a different story. If the sink writes to Redis with INCRBY on every record, a restart replays records since the last checkpoint and double-increments. End-to-end exactly-once requires either a transactional sink (two-phase commit tied to the checkpoint, e.g. Kafka transactions) or an idempotent sink (upsert the full window result by key). The Senior answer quotes the feature; the Staff answer designs the sink.

The 60-Second Pitch#

"I'd use Flink for the real-time aggregation. It's a true streaming engine — per-record processing with millisecond-level latency, not micro-batches — with first-class event time and watermarks, so late and out-of-order events land in the right window. It keeps keyed state locally in RocksDB, which scales to terabytes, and takes asynchronous incremental checkpoints to S3 every 30 seconds so a failure restores state and replays Kafka from the checkpointed offsets — exactly-once state. For exactly-once output I'd sink to Kafka with transactions, or upsert full window results by key into the serving store so replays are idempotent. I'd pick Kafka Streams instead if this is a single-service enrichment job and we don't want a cluster, and plain batch if nobody needs the answer in under an hour."

The Three Intents#

IntentConstraintStrategyFailure ModeCorrectness Bar
Real-time analytics / aggregation (dashboards, ad clicks, trending)Seconds of latency, high volume, late eventsEvent-time windows, watermarks, pre-aggregation, idempotent upsert sinkDouble counting on restart; late data silently droppedApproximately right now, exactly right later (reconciled)
Event-driven decisions (fraud, alerting, surge pricing)< 100ms–1s, stateful per entity, must not missKeyed process functions, timers, CEP, broadcast state for rulesBackpressure delays decisions; state loss forgets historyFresh > complete; missed decision is costly
Streaming ETL / materialized views (CDC → search/cache/lake)Ordering per key, schema evolution, completenessKafka/CDC source, stateless or keyed transforms, upsert sinksSchema drift breaks the job; reprocessing overwrites good dataEvery change applied once, in order per key

🎯 Staff Move: "I'll treat this as real-time aggregation with a reconciliation path: the streaming result is what users see within a minute, and a daily batch recomputation over the raw events is what finance bills from. That lets me set a 10-second watermark and move on, instead of trying to make streaming perfectly complete."

The Staff Positions#

PositionRationale
Event time by defaultProcessing time makes results depend on when the job ran, not what happened; replays corrupt them.
Every pipeline states its lateness contractWatermark delay + allowed lateness + what happens after — signed off by the consumer of the number.
Exactly-once is a sink design problemCheckpoints make state exactly-once; the sink must be transactional or idempotent.
RocksDB state with TTL on everythingHeap state caps you at JVM memory; state without TTL grows until the job dies months later.
Recovery time is a design parameterRTO ≈ state restore time + replay since last checkpoint; size both explicitly.
Pre-aggregate before the shuffleLocal combine cuts network 10–100× on skewed keys.
SQL first, DataStream API when justifiedFlink SQL covers most aggregations and joins; custom code is a maintenance and upgrade liability.

Architecture & Internals#

Five internals change design decisions: the job graph and parallelism, keyed state and state backends, checkpoints, event time and watermarks, and backpressure.

Job Graph, Parallelism, Slots#

A Flink job is a dataflow graph of operators (source → map → keyBy → window → sink). Each operator runs with a parallelism — N parallel subtasks. A keyBy hash-partitions records so all records for a key go to the same subtask, which owns that key's state. The JobManager coordinates scheduling and checkpoints; TaskManagers run subtasks in task slots. Adjacent operators with the same parallelism and no shuffle are chained into one thread, avoiding serialization.

Diagram: Job Graph, Parallelism, Slots

Max parallelism (the number of key groups, default 128 or derived from parallelism) is fixed at job creation and caps how far you can rescale keyed state later. Set it deliberately — e.g. 720 or 1024 — for jobs you expect to grow. Changing it requires state migration.

Keyed State and State Backends#

State is partitioned by key and scoped to an operator: ValueState, ListState, MapState, window contents, timers.

BackendWhere state livesSize ceilingAccess costUse when
HashMapStateBackendJVM heap objectsMemory of TaskManagers (tens of GB)~ns (object access)Small state, lowest latency
EmbeddedRocksDBStateBackendLocal SSD via RocksDB, off-heapLocal disk (TBs across cluster)~µs (serialize + LSM read)Default for production; large or growing state
Disaggregated state (Flink 2.x, ForSt)Remote object storage with local cacheEffectively unboundedHigher, cache-dependentVery large state, fast rescaling, cloud-native

State sizing is back-of-envelope arithmetic and belongs in the interview:

state = active_keys x bytes_per_key x (windows kept open)

ad click counts: 10M active ads x 64 B (count + sum + key) x 2 open windows
  = ~1.3 GB  -> trivial
dedupe by click_id for 24h: 2B clicks/day x 40 B
  = ~80 GB   -> RocksDB, TTL 24h, or a Bloom filter per hour (~1% FP)
session windows per user: 200M users x 1 KB
  = ~200 GB  -> RocksDB, incremental checkpoints mandatory

🎯 Staff Insight: "The state that kills Flink jobs isn't the window — it's the dedupe set and the join buffer nobody put a TTL on. I'd set StateTtlConfig on every piece of state and compute the steady-state size before launch, because unbounded state doesn't fail on day one; it fails on day 90 when checkpoints start timing out."

Checkpoints — Exactly-Once State#

Flink uses asynchronous barrier snapshotting, a variant of Chandy-Lamport. The JobManager injects checkpoint barriers into the sources; barriers flow with the data; when an operator has received the barrier on all inputs, it snapshots its state asynchronously and forwards the barrier. When every sink acknowledges, the checkpoint is complete and includes the source offsets.

Diagram: Checkpoints — Exactly-Once State
ConceptDefault / typicalWhy it matters
Checkpoint interval10s–60s typicalUpper bound on replay work and on transactional-sink output latency
Incremental checkpoints (RocksDB)Off by default; turn onUploads only new SST files — 100GB state, ~1–5GB per checkpoint
Aligned vs unaligned checkpointsAligned defaultUnaligned (1.11+) lets barriers overtake buffered data under backpressure, so checkpoints still complete
SavepointManual, full, portableFor upgrades, rescaling, migrations — the job's "backup"
Tolerable checkpoint failures0 by defaultRaise to a small number so one slow checkpoint doesn't restart the job

Recovery math: RTO ≈ scheduling (seconds) + state download (100GB at ~500MB/s ≈ 3–4 min, faster with local recovery) + replay since last checkpoint (30s interval ≈ ≤ 30s of data, processed at catch-up rate). Staff candidates put a number on it.

Event Time and Watermarks#

Every record carries an event timestamp (when it happened). A watermark is a special record asserting "no more events with timestamp ≤ W are expected." A window [12:00, 12:01) fires when the watermark passes 12:01.

Bounded-out-of-orderness watermark:
  W = max_event_time_seen - max_delay

max_delay = 10s:
  events: 12:00:58  12:01:03  12:00:55 (late but within 10s)  12:01:12
  watermark after 12:01:12 arrives = 12:01:02 -> window [12:00,12:01) fires
  an event stamped 12:00:40 arriving now is LATE:
     within allowedLateness (e.g. 5m) -> window re-fires with updated result
     beyond allowedLateness          -> side output (never silently dropped)
KnobIncrease itDecrease itWho pays
Watermark delayMore complete first resultLower latency to first resultDashboard users wait vs. see numbers change
Allowed latenessLater corrections acceptedLess state held openState size / checkpoint cost vs. downstream correctness
Idle source timeout—Prevents one idle Kafka partition from stalling the watermark foreverOtherwise all windows stop firing

The idle partition trap: a watermark is the minimum across all input partitions. One Kafka partition with no traffic holds the global watermark back and no window ever fires. Set withIdleness(Duration.ofMinutes(1)).

Backpressure#

Flink uses credit-based flow control: a downstream subtask grants credits (buffer space) to upstream; when a sink slows, buffers fill, credits hit zero, and the slowdown propagates back to the source, which stops reading Kafka. Nothing is dropped — lag accumulates in Kafka instead, which is exactly where you want it.

Diagnosis: the Flink UI/metrics show backPressuredTimeMsPerSecond and busyTimeMsPerSecond per subtask. The bottleneck is the first operator downstream of the backpressured ones that is busy ~100%. Common culprits: a slow sink, a hot key, a synchronous external call in a map function.


Core Usage — "The Entire Game": Keys, Windows, and Sinks#

Three decisions determine whether a Flink job is correct and cheap: what you key by, how you window, and how you write out.

Decision 1 — The Key#

keyBy decides parallelism and state locality — same trap as a Kafka partition key.

Use caseKeySkew riskMitigation
Ad click countsad_idViral ad: 1 key = 20% of trafficTwo-phase aggregation (salted key → merge)
Fraud per cardcard_idLow—
Trending hashtagshashtagExtreme (#WorldCup)Local pre-aggregation + salted keys
Sessionizationuser_idBots with 100K events/minBot filter before keyBy; cap per-key state
CDC materialized viewprimary keyCounter rowsMerge updates within a mini-batch

Two-phase aggregation for hot keys:

stage 1: keyBy(ad_id + "#" + (hash(click_id) % 16))  -> partial counts per 1m window
stage 2: keyBy(ad_id)                                 -> sum 16 partials
Hot key load spread 16 ways; stage 2 sees 16 records per ad per window, not millions.
Flink SQL does this automatically with table.optimizer.distinct-agg.split.enabled
and mini-batch + local-global aggregation settings.

Decision 2 — The Window#

WindowShapeState costUse for
TumblingFixed, non-overlapping (1m)1 accumulator per key per open windowPer-minute counts, billing buckets
Sliding / hoppingSize 10m, slide 1m → each event in 10 windows10× tumbling"Last 10 minutes" trending; prefer tumbling + query-time sum when possible
SessionGap-based (30m inactivity)Unbounded per active sessionUser sessions, engagement
Global + custom triggerYou decideYou decideCount-based or business-event-based emission
No window: keyed process + timersArbitraryExplicitFraud rules, timeouts, "no heartbeat in 60s"

Use aggregating functions, not full-window buffering. aggregate(AggregateFunction) keeps one accumulator per window; process(ProcessWindowFunction) alone buffers every record until the window fires — 1,000× the state for high-volume keys.

clicks
  .assignTimestampsAndWatermarks(
      WatermarkStrategy.<Click>forBoundedOutOfOrderness(Duration.ofSeconds(10))
          .withTimestampAssigner((c, ts) -> c.eventTimeMs)
          .withIdleness(Duration.ofMinutes(1)))
  .keyBy(c -> c.adId)
  .window(TumblingEventTimeWindows.of(Time.minutes(1)))
  .allowedLateness(Time.minutes(5))
  .sideOutputLateData(lateTag)
  .aggregate(new CountAndSum(), new EmitWithWindowBounds());

The same job in Flink SQL:

SELECT ad_id,
       window_start, window_end,
       COUNT(*)            AS clicks,
       SUM(bid_micros)     AS spend_micros
FROM TABLE(
  TUMBLE(TABLE clicks, DESCRIPTOR(event_time), INTERVAL '1' MINUTE))
GROUP BY ad_id, window_start, window_end;
-- with WATERMARK FOR event_time AS event_time - INTERVAL '10' SECOND in the DDL

Decision 3 — The Sink (Where Exactly-Once Lives or Dies)#

Sink patternEnd-to-end guaranteeLatency costExample
Transactional (2PC)Exactly-onceOutput visible only at checkpoint commit (= checkpoint interval)Kafka sink with DeliveryGuarantee.EXACTLY_ONCE
Idempotent upsertEffectively-onceNoneUPSERT window result keyed by (ad_id, window_start) into Postgres/Cassandra/Redis SET
Append / incrementAt-least-once → duplicatesNoneINCRBY, INSERT without key — avoid for counts
Write-ahead + dedupe downstreamEffectively-once if consumer dedupesSmallEmit with (key, window, checkpoint_id); consumer ignores repeats

🎯 Staff Move: "I'd write each window's final count as an upsert keyed by (ad_id, window_start). If the job replays after a failure, it rewrites the same value — idempotent, no transaction coordinator, no checkpoint-interval latency on the output. I'd only use the transactional Kafka sink when the consumer is another stream job that needs a clean changelog."

Kafka transactional sink gotcha: Kafka's transaction.max.timeout.ms (broker default 15 minutes) must exceed checkpoint interval + max recovery time, or transactions abort and data is lost; and downstream consumers must use isolation.level=read_committed, which means they see output only after each checkpoint — a 60s checkpoint interval is a 60s floor on end-to-end latency.

Joins#

JoinStateRiskUse
Interval join (click within 30m of impression)Both sides buffered for the intervalState = rate × intervalAttribution
Temporal / versioned table join (order × FX rate as of order time)Latest versions of the table sideNeeds a changelog with event timeEnrichment with slowly changing data
Lookup join (async call to DB/cache)None (cache optional)External latency and load; non-deterministic on replaySmall reference data; use AsyncDataStream with capacity limits
Regular unbounded joinBoth sides foreverUnbounded stateAlmost never without TTL
Broadcast stateSmall rule set on every subtaskMust fit in memory × parallelismFraud rules, feature flags, config

The Tunable Tradeoff — Latency × Completeness × Cost#

The formula every streaming design negotiates:

first_result_latency  ~= window_size + watermark_delay + processing + sink_commit
                         (sink_commit = checkpoint interval if transactional)
completeness_at_first ~= fraction of events with (arrival - event_time) < watermark_delay
correction_window     =  allowed_lateness  (state held open this long)
cost                  ~= state x checkpoint_frequency + parallelism x uptime

Example: 1m window, 10s watermark, 30s checkpoint, transactional sink
  first result ~= 60 + 10 + <1 + up to 30 = up to ~100s after window start
  switch to idempotent upsert sink: ~70s
  if p99 lateness is 45s (mobile), 10s watermark gives ~97% complete first result
ChoiceWhat WorksWhat BreaksWho Pays
Short watermark (5s)Fast dashboardsResults revised or incomplete for mobile/IoT sourcesDashboard consumers see numbers change
Long watermark (5m)Complete first answer5-minute-stale "real-time"Product: "real-time" becomes "near-real-time"
Long allowed latenessLate data correctedState and checkpoint size grow with itPlatform budget; recovery time
Frequent checkpoints (5s)Small replay, low transactional latencyCheckpoint overhead, S3 request cost, pressure on large stateJob throughput; S3 bill
Infrequent checkpoints (5m)Low overheadUp to 5m replay; 5m output latency on transactional sinksRecovery time; downstream freshness
Streaming-only (Kappa)One code pathReprocessing history requires replaying Kafka (retention!)Whoever needs a 6-month backfill
Streaming + batch (Lambda)Batch is the audited truthTwo implementations driftThe team reconciling mismatches

🎯 Staff Move: "I'll ask product what 'real-time' means for this number. If the answer is 'within a couple of minutes and it can wobble,' I take a 10-second watermark and upsert sink. If the answer is 'we bill from it,' then streaming is the preview and batch is the invoice."


1. State Without TTL#

Dedupe sets, join buffers, and per-user maps that grow forever. Checkpoints grow from 5GB to 800GB over three months, start timing out, the job can't complete a checkpoint, and the next failure replays hours of data. Fix: StateTtlConfig on every state descriptor, alert on checkpoint size trend.

2. Synchronous External Calls in Operators#

A map that calls a REST API for enrichment at 20ms per record caps each subtask at 50 records/s and backpressures the whole job. Fix: AsyncDataStream.unorderedWait with a capacity limit and timeout, a local cache, or turn the reference data into a CDC stream and join.

3. Processing-Time Windows for Business Metrics#

Results depend on when the job ran; any replay or restart smears data across windows. Fix: event time.

4. Non-Idempotent Sinks with "Exactly-Once" Claims#

INCRBY into Redis from a checkpointed job double-counts every record since the last checkpoint on every failure. Fix: upsert full window results, or a transactional sink.

5. Ignoring Watermark Stalls#

One idle partition, one source with a clock 1 hour in the future (pushes the watermark forward → everything else is "late"), or a source 1 hour behind (holds it back). Fix: idleness timeouts, timestamp sanity filters (drop/side-output events more than N minutes in the future), per-source watermark monitoring (currentInputWatermark).

6. One Giant Job for Everything#

Twenty unrelated pipelines in one job graph because "it's one deployment." One bad UDF restarts all twenty; upgrading one requires a savepoint of all. Fix: one job per business pipeline with its own owner and SLO.

7. Changing Operator Topology Without UIDs#

Savepoint restore maps state to operators by UID. Without explicit .uid("click-window"), Flink auto-generates IDs from the graph; adding a filter upstream changes them, and the upgraded job cannot restore state. Fix: set uid() on every stateful operator from day one.


The Technology Landscape — Head-to-Head Comparison#

DimensionApache FlinkKafka StreamsSpark Structured StreamingGoogle Dataflow (Beam)Streaming DBs (RisingWave, Materialize)
ModelTrue streaming, per recordLibrary inside your serviceMicro-batch (default)Beam model on managed runnerSQL materialized views, incrementally maintained
Latencyms–sub-secondms–sub-second~100ms–seconds per micro-batchSecondsSub-second to seconds
StateRocksDB/heap, TB scale, S3 checkpointsRocksDB + Kafka changelog topicsState store per micro-batch, checkpoint to HDFS/S3ManagedManaged, persisted
Event timeFirst-class, richestGood (grace periods)Good (watermarks)First-class (Beam originated much of the model)Supported
DeploymentCluster (K8s operator, YARN), or managed (AWS Managed Flink, Confluent)Just your app's podsSpark clusterFully managedManaged or cluster
SourcesAnything (Kafka, Kinesis, CDC, files)Kafka onlyManyManyKafka, CDC, Postgres
Ops burdenMedium–highLow (it's your service)Medium (usually already have Spark)LowLow–medium
Pick whenLarge state, complex event time, many sources, sub-secondKafka-in/Kafka-out logic owned by one service teamOrg is Spark-heavy; latency of seconds is fine; unify with batchGCP shop; want zero opsTeam speaks SQL; the output is a queryable view

The rule of thumb: Kafka Streams when the logic belongs to one service and both ends are Kafka; Flink when state, scale, or event-time complexity justify a separate runtime; Spark when the org already runs Spark and seconds are acceptable; managed (Dataflow, AWS Managed Flink) when nobody wants to be on call for a JobManager.


Patterns#

Pattern 1: Kappa — Streaming as the Only Path#

Kafka (long retention or tiered storage) → Flink → serving store. Reprocessing = deploy the new job version reading from an old offset, write to a new output table, cut over. Use when retention covers your backfill needs and correctness tolerates streaming semantics.

Pattern 2: Streaming Preview + Batch Truth (Lambda, Done Deliberately)#

Diagram: Pattern 2: Streaming Preview + Batch Truth (Lambda, Done Deliberately)

The batch job overwrites the streaming result for closed days. The reconciliation delta (e.g. alert if > 0.5%) is a correctness metric for the streaming job. Use for billing, ad spend, payouts — anything audited. Flink's unified batch/stream API and Flink SQL let one codebase serve both modes, which removes the classic Lambda drift problem.

Pattern 3: CDC → Materialized View#

Debezium/Flink CDC reads a database changelog; Flink joins and aggregates; upserts a denormalized view into Elasticsearch, Redis, or a wide table. Replaces dual writes and nightly ETL. Watch: schema changes in the source break the job; ordering per primary key must be preserved end to end.

Pattern 4: Stateful Event-Driven Application#

A keyed process function with timers implements a per-entity state machine: fraud scoring per card, "order not picked up in 10 minutes → alert," driver heartbeat timeouts. Rules arrive via broadcast state so they update without redeploying. Use when decisions depend on per-entity history and must fire on the absence of events.

Pattern 5: Pre-Aggregation Tier#

Flink rolls raw events to 1-minute aggregates before they reach the TSDB or OLAP store, cutting storage and query cost 10–1000×. See Time Series DBs.


Scaling#

Numbers#

ResourceRough figureNote
Stateless throughput per core~50K–500K records/sDepends on (de)serialization; chaining avoids copies
Keyed windowed aggregation per core~10K–100K records/s with RocksDBSerialization + LSM access dominate
RocksDB state per TaskManagerHundreds of GB comfortably on local NVMeCheckpoint upload bandwidth is the limit
Checkpoint interval10s–60s typical; minutes for TB stateMust complete well inside the interval
Incremental checkpoint delta~1–10% of state per checkpointWorkload dependent
Restore throughput from S3~200MB–1GB/s per job (parallel)1TB state → ~15–60 min cold restore without local recovery
Parallelism≤ max parallelism (key groups)Fixed at first launch; set 4–8× initial parallelism

Rescaling#

Rescaling keyed state is a savepoint → stop → restart with new parallelism operation (the reactive/adaptive schedulers automate parts of it). Key groups are redistributed across subtasks. Time to rescale ≈ savepoint time + restore time — minutes for GB, tens of minutes for TB. That is why autoscaling Flink is slower and coarser than autoscaling stateless services, and why large-state jobs are over-provisioned by 30–50% for peak.

Skew#

Parallelism doesn't help if one key is 20% of traffic — one subtask is pegged while 95 idle. Two-phase aggregation, salting, and filtering bots before keyBy are the fixes; adding TaskManagers is not.

Multi-Region#

Streaming jobs are almost always region-local: each region processes its own Kafka and writes regional results; a global view is either a second-level aggregation job reading regional outputs (mirrored topics) or a query-time merge. Running one Flink job across regions puts WAN latency inside every shuffle and checkpoint. For DR, keep savepoints/checkpoints replicated to the standby region and replay from the mirrored Kafka offsets — RPO = mirror lag, RTO = restore time.


Failure Modes & Recovery#

1. Checkpoint Timeouts → Restart Loop#

  • Symptom: Checkpoints fail repeatedly; job restarts; each restart replays more data.
  • Root cause: Backpressure delaying aligned barriers, state growth (no TTL), or S3 throttling.
  • Detection: lastCheckpointDuration trending toward checkpointTimeout; numberOfFailedCheckpoints; lastCheckpointSize growth week over week.
  • Fix: Enable unaligned checkpoints; raise tolerable failures; fix the backpressure source; add TTL.
  • Prevention: Alert when checkpoint duration > 50% of interval; state size budget per job.

2. Watermark Stall — Windows Stop Firing#

  • Symptom: Output stops, but the job is "RUNNING" and consuming; no errors.
  • Root cause: An idle partition, a stuck source, or a producer with a clock far behind holds the min watermark back.
  • Detection: currentOutputWatermark not advancing vs wall clock (alert when now - watermark > 5 × watermark_delay); output record rate = 0.
  • Fix: Idleness timeout; filter/side-output events with absurd timestamps.
  • Prevention: Watermark lag is a first-class SLO metric on every event-time job.

3. Backpressure → Kafka Lag#

  • Symptom: Consumer lag on the source topic grows; outputs increasingly stale.
  • Root cause: Slow sink, hot key, synchronous external calls, GC pressure.
  • Detection: backPressuredTimeMsPerSecond high upstream; busyTimeMsPerSecond ~1000 on the bottleneck operator; Kafka records-lag-max.
  • Fix: Scale the bottleneck (if not skew), async I/O, batch sink writes, two-phase aggregation.
  • Prevention: Load test at 2–3× peak with production key distribution.

4. Duplicate Output After Recovery#

  • Symptom: Counts in the serving store are higher than the batch truth after an incident.
  • Root cause: Non-idempotent sink (increments/appends) replayed from the last checkpoint.
  • Detection: Stream-vs-batch reconciliation delta > threshold; spike correlated with numRestarts.
  • Fix: Recompute affected windows from raw data; overwrite.
  • Prevention: Upsert-by-key or transactional sinks only; ban increment sinks in review.

5. Savepoint Incompatibility on Upgrade#

  • Symptom: New job version cannot restore state; deploy blocked or forced to start empty.
  • Root cause: Missing operator UIDs, changed state serializer (POJO field types, Kryo), or Flink major-version state format change.
  • Detection: Restore failure in CI when testing against a production savepoint.
  • Fix: allowNonRestoredState for dropped operators; State Processor API to migrate state; last resort — backfill from Kafka/lake.
  • Prevention: UIDs on every stateful operator; Avro/Protobuf state types with schema evolution; CI job that restores the latest prod savepoint into the new build.
Diagram: 5. Savepoint Incompatibility on Upgrade

Operational Reality Matrix#

FailureDetection SignalBlast RadiusMitigationOwner
Checkpoint timeoutslastCheckpointDuration > 50% intervalOne job, growing replayUnaligned, TTL, fix backpressureJob owner
Watermark stallnow - currentOutputWatermark > thresholdAll windowed outputs of the jobIdleness, timestamp filterJob owner; source team for bad clocks
Backpressurebusy/backpressured time; Kafka lagFreshness of every downstream consumerScale bottleneck, async I/OJob owner
Duplicate outputReconciliation deltaDownstream reports, possibly billingRecompute from rawJob owner + data consumer
Upgrade breaks stateCI savepoint restore testOne job's historyState migration or backfillJob owner
Cluster-wide failure (shared session cluster)JobManager down, all jobs restartingEvery job on the clusterApplication mode, per-job clustersStreaming platform

When to Use vs. Alternatives#

NeedPickWhy
Results needed within hours; daily reportsBatch (Spark/SQL on the warehouse)Cheaper, simpler, easier to reprocess
Kafka-in/Kafka-out enrichment owned by one teamKafka StreamsNo cluster; deploy with the service
Large keyed state, event-time windows, multiple sources, sub-secondFlinkBest-in-class state and time semantics
Already all-in on Spark, seconds OKSpark Structured StreamingReuse platform and skills
Output is a continuously updated SQL viewFlink SQL or a streaming databaseDeclarative, less code to own
Simple per-event transform, no stateConsumer service or serverlessA stream processor is overkill
Billing-grade totalsBatch as truth, streaming as previewAuditable recomputation
  • Nobody needs the answer in under an hour. Batch is 3–10× cheaper to run and far cheaper to debug.
  • Stateless transforms. A Kafka consumer in a normal service does it with less machinery.
  • The team can't own a stateful distributed runtime (upgrades, savepoints, RocksDB tuning) — use a managed offering or Flink SQL on a platform.
  • The "stream" is 50 events/s. A cron job and a SQL query win.
Diagram: When NOT to Use Flink

Operational Concerns#

What the On-Call Actually Does#

  1. Watches four signals per job: source lag (seconds), watermark lag, checkpoint duration/size, restart count. Those four catch ~90% of incidents.
  2. Upgrades via savepoints: take savepoint → deploy new version from savepoint → verify watermark advances and output resumes. Keep the previous savepoint for rollback.
  3. Tunes RocksDB memory (managed memory fraction, block cache, write buffers). Most "Flink is slow" tickets are RocksDB under-provisioned on memory or on network-attached disks.
  4. Runs backfills as separate job instances reading from an old offset or the lake, writing to a shadow table, then swapping — never by rewinding the production job's offsets.
  5. Uses application mode on Kubernetes (one JobManager per job via the Flink Kubernetes Operator) so one job's failure can't take down others.

Key Metrics & Alerts#

MetricHealthyAlert
Source consumer lag (seconds)< pipeline SLO> SLO for 5m
now - currentOutputWatermark≈ watermark delay> 5× delay
lastCheckpointDuration< 30% of interval> 50% of interval
numberOfFailedCheckpoints rate0> 2 in 10m
numRestartsFlatAny increase (ticket), > 3/h (page)
lastCheckpointSize trendStable ±10%/week+50%/week (missing TTL)
Late records dropped / side-output rate< 0.1%> 1% (watermark too tight)
Stream vs batch reconciliation delta< 0.5%> 1%

Interview Application — Staff-Level Plays#

Case StudyHow Stream Processing Is UsedKey Pattern
Stream ProcessingAd click aggregation — the canonical questionEvent-time windows, idempotent upsert, batch reconciliation
LeaderboardReal-time score aggregation feeding a sorted setTwo-phase aggregation for hot keys
Metrics & MonitoringPre-aggregation before the TSDBRollups reduce cardinality
Search IndexingCDC → denormalized documentsTemporal joins, upsert sink
Ride Hailing & DeliverySupply/demand per geo cell for surge pricingSliding windows per cell, broadcast config
Payment ProcessingReal-time fraud scoringKeyed state per card + timers
Data Pipeline PatternsKappa vs Lambda decisionsStreaming preview, batch truth

Every System Design Question Has a Stream Processing Moment#

  • URL shortener: "Click analytics are a Flink job: 1-minute tumbling windows per link, upserted by (link_id, minute). The redirect path never waits on it."
  • Notification system: "Digest emails are a session-window job — group a user's events with a 15-minute gap, then emit one notification."
  • News feed: "Trending topics use a hopping window with two-phase aggregation, because #WorldCup is one key."
  • Rate limiter: "Abuse detection runs asynchronously in Flink over request logs; the inline limiter stays local and fast. Flink updates the blocklist via broadcast state."

What Interviewers Probe#

After You Say...They Will Ask...What They're Evaluating
"Flink does exactly-once""What does your sink do on replay?"Sink idempotency / 2PC
"Windows of 1 minute""Event time or processing time? What about late events?"Time semantics, lateness contract
"Flink keeps state""How big? What if a node dies?"State sizing, recovery time
"Key by ad_id""One ad gets 30% of clicks."Skew handling
"We'll reprocess if there's a bug""From where? How far back?"Retention, backfill design
"Real-time dashboard""Finance says it doesn't match the invoice."Reconciliation, source of truth

Common Interview Mistakes#

What Candidates SayWhat Interviewers HearWhat Staff Engineers Say
"Exactly-once end to end with Flink"Double counts incoming"Exactly-once state; upsert sink makes output effectively-once."
"Drop late events"Silent data loss"Late within 5m updates the result; later goes to a side output reconciled by batch."
"Store counts with INCRBY"Non-idempotent under replay"Write the full window value keyed by window."
"Add parallelism for the hot key"Doesn't understand keyed routing"Salt the key, aggregate twice."
"Streaming replaces batch"Ignores audit and backfill"Streaming previews, batch is the truth for money."
"Call the user service in the map"Backpressure incoming"Async I/O with a cache, or join a CDC stream."

L5 vs L6 vs L7 Responses#

ScenarioL5 AnswerL6 / Staff AnswerL7 / Principal Answer
"Ad click counts per minute"Flink window, write to RedisEvent time, 10s watermark, 5m lateness, side output, dedupe on click_id with 24h TTL, upsert by (ad_id, minute), nightly batch overwrite + reconciliationPlus: define the correctness tier with finance — streaming for pacing, batch for invoices — and one shared definition of a "click" across both paths
"The job keeps restarting"Increase resourcesCheck checkpoint duration/size, backpressure, state TTL; unaligned checkpointsWhy is a single job's failure paging a human? Platform SLOs, self-service diagnostics, restart budgets
"Upgrade Flink 1.x → 2.x"RedeploySavepoint compatibility tests per job, UIDs, staged rolloutFleet upgrade program across 200 jobs: compatibility matrix, owners, deadline, what gets rewritten in SQL
"Should we adopt streaming?"Yes, it's modernFor the 3 pipelines with sub-hour needs; batch for the restCost and headcount model; platform vs managed; SQL as the paved road

The Staff Stream Processing Checklist#

  1. Latency need: "Product needs this within 2 minutes; batch can't do it, so streaming is justified."
  2. Time semantics: "Event time, 10s watermark, 5m allowed lateness, side output for the rest."
  3. Key and skew: "Key by ad_id with two-phase aggregation for hot ads."
  4. State budget: "~1.3GB for windows, ~80GB for 24h dedupe — RocksDB, TTLs, 30s incremental checkpoints."
  5. Sink semantics: "Upsert by (ad_id, window_start) — replays are harmless."
  6. Truth and recovery: "Nightly batch overwrites closed days; reconciliation delta alerts at 0.5%; RTO ~5 minutes."

🎯 Staff Insight: Don't use Flink for stateless transforms, for anything nobody needs within the hour, or as the sole source of truth for money. The strongest signal is stating the lateness contract and who signed it.

Evaluation Rubric#

DimensionSenior (L5)Staff (L6)Principal (L7)
TimeWindowsEvent time, watermarks, lateness policyOrg lateness contracts per data product
Correctness"Exactly-once"State vs sink guarantees; idempotent/2PC sinks; reconciliationCorrectness tiers; single metric definitions across batch and stream
State"Flink handles it"Sized, TTL'd, backend chosen, recovery time computedState as a long-lived migration liability
Operations"Monitor the job"Lag, watermark, checkpoint, restart SLOsPlatform model: application mode, self-service, upgrade programs
ChoiceFlink for everything streamingFlink vs Kafka Streams vs batch by needBuild vs managed vs retire; SQL paved road

Strong hire signals

SignalWhat It Sounds Like
Time contract"Dashboards final after 5 minutes, invoices after 24 hours."
Sink-level exactly-once"Replays rewrite the same row."
Quantified state"80GB of dedupe state with a 24h TTL."
Recovery math"30s checkpoints, ~3 minutes to restore — RTO under 5."
Knows when to batch"This report is daily; I wouldn't stream it."

Lean no-hire signals

SignalWhy It Misses the Bar
Processing time for business metricsResults change on replay
No late-data policySilent loss or unbounded state
Increment sinks with "exactly-once" claimsDouble counts on every failure
No reconciliation for moneyUnauditable

Common false positives

  • Knowing Chandy-Lamport ≠ designing the sink. The algorithm is solved; the sink isn't.
  • Flink API fluency ≠ judgment — can they argue for batch?
  • "We process a billion events a day" ≠ correctness — ask what happens on replay.

The Principal Lens#

Why L7 Sees This Problem Differently#

At Staff level a streaming job is a correct, well-operated pipeline. At Principal level, streaming is a second computation of numbers the org already computes in batch, and every such duplicate is a future argument about which number is right. The L7 question is: which data products genuinely need sub-hour freshness, what does it cost per pipeline to provide it, and how do we make sure a "click" means the same thing in the stream, the warehouse, and the invoice? Streaming platforms tend to accrete hundreds of hand-written jobs whose authors have left, each holding state that makes it painful to upgrade. The Principal job is to keep that fleet small, declarative, and owned.

The Org-Level Fault Line#

A central streaming platform with SQL as the paved road vs. every team writing and running its own jobs.

OptionWhat WorksWhat BreaksWho Pays
Central platform, Flink SQL self-serviceConsistent ops, upgrades once, metric definitions sharedComplex logic doesn't fit SQL; platform team becomes a bottleneckPlatform headcount
Platform runtime, teams write DataStream jobsFlexibility200 bespoke jobs; upgrade program takes a yearJob owners at every upgrade; platform
Each team runs its own Flink/Kafka StreamsAutonomyDuplicate definitions, uneven ops qualityIncident responders; data consumers who see conflicting numbers
Managed service (Dataflow, AWS Managed Flink, Confluent)No runtime opsPer-unit pricing at scale; feature lag; lock-in on runner specificsFinance

The Principal default: one platform, Flink SQL (or a streaming database) as the default for aggregations and joins, DataStream API by exception with an owner and upgrade commitment, and metric definitions shared with the batch warehouse through a semantic layer.

Cost Model#

Assumptions: ~$0.04/vCPU-hr, ~$0.005/GB-RAM-hr, local NVMe for RocksDB, S3 checkpoints, loaded engineer ~$250K/yr; running 24/7. Directional.

ScaleThroughputJobsInfra/monthHeadcountOn-call
Startup5K ev/s3–5~$2–4K (managed Flink or small K8s)0.5 FTEOwned by data team, rare pages
Growth200K ev/s30–60~$25–50K2–3 FTE platformWeekly incidents; platform rotation
Enterprise5M+ ev/s300+~$300K–1M8–15 FTEDedicated platform rotation + job-owner rotations

The comparison that matters: a streaming pipeline typically costs 3–10× its batch equivalent (always-on compute, state, checkpoints, on-call). Every pipeline should earn that multiple with a real freshness requirement.

The 3-Year Evolution Path#

Diagram: The 3-Year Evolution Path

One-Way Doors vs Two-Way Doors#

DecisionReversibilityCost to Reverse
Max parallelism of a stateful jobOne-way without state migrationState Processor API rewrite or full backfill
State serializer (Kryo vs Avro/Protobuf)One-way-ishState migration per job
Metric definitions exposed to financeOne-way once invoices depend on themRestatements, customer communication
Engine choice (Flink vs Spark vs Dataflow) for 200 jobsOne-way at fleet scaleYear-long rewrite
Checkpoint interval, parallelism, watermark delayTwo-wayConfig + savepoint restart
Managed vs self-hosted FlinkTwo-way-ishSavepoints are portable across same-version Flink

🧭 Principal Move: "I'd guard metric definitions and state formats, not the engine version. If every job uses Avro state and SQL definitions from the semantic layer, engine upgrades are a quarter. If 200 jobs use Kryo-serialized POJOs, every upgrade is a negotiation."

The Standard I'd Write#

RFC-STR-003: Streaming Pipelines Standard

Scope: All production stream processing jobs on the shared platform.

MUST
  1. Declare a freshness SLO, a lateness contract (watermark delay, allowed
     lateness, late-data handling) and a correctness tier
     (approximate | exactly-once state | reconciled-with-batch).
  2. Use event time for any business metric.
  3. Write via idempotent upsert or transactional sinks; increment/append sinks
     require platform approval.
  4. Set uid() on every stateful operator and TTL on every state descriptor.
  5. Pass a CI check that restores the latest production savepoint.
  6. Pipelines feeding billing or finance MUST have batch reconciliation with an
     alert threshold <= 1%.
SHOULD
  7. Be written in Flink SQL; DataStream jobs need a named owner and a written
     reason SQL is insufficient.
  8. Run in application mode, one job per deployment.

Exceptions: Streaming platform lead + consuming data owner; reviewed quarterly.

Rollout: new jobs immediately; existing jobs by criticality tier over two
quarters (shadow report -> owner notified -> enforced at next deploy).

Success metrics: upgrade of the full fleet to a new minor version in < 1 quarter;
reconciliation deltas < 0.5%; < 1 page per job per quarter.

What I'd Tell the VP#

"Streaming gives us minute-fresh numbers, but each pipeline costs 3 to 10 times what the same report costs in the nightly batch, and we now have two systems that sometimes disagree about revenue. I'm proposing we keep streaming for the dozen things that truly need it — fraud, ad pacing, live dashboards — and move the rest back to batch. The ones that stay use one shared definition of each metric and are checked against batch nightly, so finance never sees two answers. That should cut our streaming spend by roughly a third and remove a recurring source of incidents."

Principal Interview Signals#

SignalWhat It Sounds Like
Justifies streaming economically"This costs 5× batch; the fraud use case earns it, the weekly report doesn't."
Owns definitions"A click is defined once and used by both paths."
Plans for upgrades"State formats and UIDs are fleet policy because upgrades are where streaming hurts."
Chooses the paved road"SQL by default; custom code needs an owner."
Retires things"Year three goal: fewer jobs, not more."

Staff answers that L7 interviewers find insufficient:

  • "We'll reconcile with batch" — without saying which number the business uses and who resolves mismatches.
  • "Use Flink SQL" — without a plan for the 150 existing DataStream jobs.
  • "Autoscale the jobs" — without acknowledging stateful rescaling cost and the 3–10× cost multiple.

In the Wild#

Alibaba adopted Flink (via its internal fork, Blink, later contributed back upstream) for real-time search ranking, recommendations, and the live GMV dashboards during Singles' Day, publicly reporting peak processing in the billions of records per second. Much of Flink SQL's maturity and its batch/stream unification came from that investment.

Staff insight: The workloads that justify Flink combine huge state, strict freshness, and business visibility. Say which of those three your problem has.

Uber — Streaming Platform and SQL Self-Service#

Uber built its real-time analytics on Kafka and Flink and has written about exposing streaming via SQL (the AthenaX project) so teams could build pipelines without writing Flink code, plus operating Flink at scale for use cases like surge pricing and marketplace metrics.

Staff insight: Large orgs converge on SQL as the streaming interface. Proposing a self-service SQL layer is a Principal-level answer to "how do 50 teams use streaming safely."

Netflix — Keystone Routing and Stream Processing#

Netflix's Keystone platform uses Kafka and Flink to route and process event data at trillions of events per day, offering a managed, self-service experience for routing jobs while supporting custom stream processing apps for more complex needs.

Staff insight: Two tiers — a declarative self-service tier for common patterns and a custom tier for hard problems — is how a platform stays small while serving many teams.


Practice Drill#

Prompt: "Build real-time ad spend tracking so advertisers see spend within 1 minute and campaigns stop when they hit budget. 500K impression/click events per second at peak. Finance invoices from the same data."

Staff Answer

Two outputs with different correctness bars. Pacing (stop campaigns at budget) needs freshness and can tolerate small overspend; invoicing needs exactness and can wait a day. Pipeline: Kafka (ad-events, 256 partitions, keyed by campaign_id) → Flink. Dedupe on event_id with a 24h TTL (~500K/s × 86,400 × 40B ≈ 1.7TB — too big, so dedupe within 1h in state (~72GB, RocksDB) and rely on batch for the long tail). Event-time 1-minute tumbling windows per campaign_id, 10s watermark, 5-minute allowed lateness, side output beyond. Hot campaigns get two-phase aggregation with 16 salts. Sink: upsert (campaign_id, minute) → spend into the serving store, plus a keyed process function holding cumulative spend per campaign that emits a budget_exhausted event to Kafka when spend ≥ 98% of budget (2% buffer for late data). Checkpoints every 30s, incremental, to S3; RTO ≈ 3–5 minutes, during which pacing falls back to the last known spend and a conservative overspend cap. Invoices come from a nightly batch job over the raw events in the lake; it overwrites closed days and alerts if stream vs batch differs by > 0.5%.

Why this is L6:

  • Separates pacing (fresh, approximate) from invoicing (exact, delayed) and designs each for its bar.
  • Sizes state and chooses a dedupe horizon deliberately instead of "dedupe everything."
  • Idempotent upsert sink + reconciliation instead of relying on exactly-once claims.
  • Defines degraded mode during recovery with a bounded overspend.

What L7 adds:

  • Gets finance and the ads product owner to sign the overspend tolerance (2%) and the reconciliation threshold.
  • One definition of "billable click" shared by the stream job and batch via a semantic layer.
  • Prices the streaming path against a 5-minute micro-batch alternative and keeps it only because pacing loss exceeds its cost.

Quick Reference Card#

Model:              per-record streaming, keyed state, async barrier checkpoints
State backend:      RocksDB for production; TTL on every state
Checkpoints:        10-60s, incremental, to S3; unaligned under backpressure
Exactly-once:       state yes; output only with 2PC (Kafka txn) or idempotent upsert
Kafka txn sink:     output visible per checkpoint; transaction timeout > ckpt + recovery
Watermark:          max_event_time - delay (typ. 5-30s); withIdleness(1m)
Late data:          allowedLateness (typ. 1-10m) + side output; never silent drop
Hot keys:           two-phase / local-global aggregation, salting
Max parallelism:    fixed at first launch; set 4-8x initial parallelism
Recovery (RTO):     schedule + state restore (~GB/s) + replay since checkpoint
Upgrades:           savepoints + uid() on stateful operators + CI restore test
Streaming cost:     ~3-10x equivalent batch
Truth for money:    batch recomputation; stream is preview

RED FLAGS
  - Processing time for business metrics
  - INCRBY / append sinks with "exactly-once" claims
  - State without TTL
  - Synchronous external calls in operators
  - No uid() on stateful operators
  - Streaming a report nobody reads within the hour
  1. Loading the index…