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#
| Behavior | Senior (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 sinks | Declares a correctness tier per pipeline: approximate / exactly-once / auditable-with-reconciliation |
| Time | Processing time windows | Event time, watermarks, allowed lateness, side outputs for late data | Sets 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 checkpoints | Treats 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 diagnosis | Platform-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 runtime | Streaming 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#
| Intent | Constraint | Strategy | Failure Mode | Correctness Bar |
|---|---|---|---|---|
| Real-time analytics / aggregation (dashboards, ad clicks, trending) | Seconds of latency, high volume, late events | Event-time windows, watermarks, pre-aggregation, idempotent upsert sink | Double counting on restart; late data silently dropped | Approximately right now, exactly right later (reconciled) |
| Event-driven decisions (fraud, alerting, surge pricing) | < 100ms–1s, stateful per entity, must not miss | Keyed process functions, timers, CEP, broadcast state for rules | Backpressure delays decisions; state loss forgets history | Fresh > complete; missed decision is costly |
| Streaming ETL / materialized views (CDC → search/cache/lake) | Ordering per key, schema evolution, completeness | Kafka/CDC source, stateless or keyed transforms, upsert sinks | Schema drift breaks the job; reprocessing overwrites good data | Every 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#
| Position | Rationale |
|---|---|
| Event time by default | Processing time makes results depend on when the job ran, not what happened; replays corrupt them. |
| Every pipeline states its lateness contract | Watermark delay + allowed lateness + what happens after — signed off by the consumer of the number. |
| Exactly-once is a sink design problem | Checkpoints make state exactly-once; the sink must be transactional or idempotent. |
| RocksDB state with TTL on everything | Heap state caps you at JVM memory; state without TTL grows until the job dies months later. |
| Recovery time is a design parameter | RTO ≈ state restore time + replay since last checkpoint; size both explicitly. |
| Pre-aggregate before the shuffle | Local combine cuts network 10–100× on skewed keys. |
| SQL first, DataStream API when justified | Flink 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.
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.
| Backend | Where state lives | Size ceiling | Access cost | Use when |
|---|---|---|---|---|
| HashMapStateBackend | JVM heap objects | Memory of TaskManagers (tens of GB) | ~ns (object access) | Small state, lowest latency |
| EmbeddedRocksDBStateBackend | Local SSD via RocksDB, off-heap | Local 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 cache | Effectively unbounded | Higher, cache-dependent | Very 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
StateTtlConfigon 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.
| Concept | Default / typical | Why it matters |
|---|---|---|
| Checkpoint interval | 10s–60s typical | Upper bound on replay work and on transactional-sink output latency |
| Incremental checkpoints (RocksDB) | Off by default; turn on | Uploads only new SST files — 100GB state, ~1–5GB per checkpoint |
| Aligned vs unaligned checkpoints | Aligned default | Unaligned (1.11+) lets barriers overtake buffered data under backpressure, so checkpoints still complete |
| Savepoint | Manual, full, portable | For upgrades, rescaling, migrations — the job's "backup" |
| Tolerable checkpoint failures | 0 by default | Raise 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)
| Knob | Increase it | Decrease it | Who pays |
|---|---|---|---|
| Watermark delay | More complete first result | Lower latency to first result | Dashboard users wait vs. see numbers change |
| Allowed lateness | Later corrections accepted | Less state held open | State size / checkpoint cost vs. downstream correctness |
| Idle source timeout | — | Prevents one idle Kafka partition from stalling the watermark forever | Otherwise 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 case | Key | Skew risk | Mitigation |
|---|---|---|---|
| Ad click counts | ad_id | Viral ad: 1 key = 20% of traffic | Two-phase aggregation (salted key → merge) |
| Fraud per card | card_id | Low | — |
| Trending hashtags | hashtag | Extreme (#WorldCup) | Local pre-aggregation + salted keys |
| Sessionization | user_id | Bots with 100K events/min | Bot filter before keyBy; cap per-key state |
| CDC materialized view | primary key | Counter rows | Merge 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#
| Window | Shape | State cost | Use for |
|---|---|---|---|
| Tumbling | Fixed, non-overlapping (1m) | 1 accumulator per key per open window | Per-minute counts, billing buckets |
| Sliding / hopping | Size 10m, slide 1m → each event in 10 windows | 10× tumbling | "Last 10 minutes" trending; prefer tumbling + query-time sum when possible |
| Session | Gap-based (30m inactivity) | Unbounded per active session | User sessions, engagement |
| Global + custom trigger | You decide | You decide | Count-based or business-event-based emission |
| No window: keyed process + timers | Arbitrary | Explicit | Fraud 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 pattern | End-to-end guarantee | Latency cost | Example |
|---|---|---|---|
| Transactional (2PC) | Exactly-once | Output visible only at checkpoint commit (= checkpoint interval) | Kafka sink with DeliveryGuarantee.EXACTLY_ONCE |
| Idempotent upsert | Effectively-once | None | UPSERT window result keyed by (ad_id, window_start) into Postgres/Cassandra/Redis SET |
| Append / increment | At-least-once → duplicates | None | INCRBY, INSERT without key — avoid for counts |
| Write-ahead + dedupe downstream | Effectively-once if consumer dedupes | Small | Emit 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#
| Join | State | Risk | Use |
|---|---|---|---|
| Interval join (click within 30m of impression) | Both sides buffered for the interval | State = rate × interval | Attribution |
| Temporal / versioned table join (order × FX rate as of order time) | Latest versions of the table side | Needs a changelog with event time | Enrichment with slowly changing data |
| Lookup join (async call to DB/cache) | None (cache optional) | External latency and load; non-deterministic on replay | Small reference data; use AsyncDataStream with capacity limits |
| Regular unbounded join | Both sides forever | Unbounded state | Almost never without TTL |
| Broadcast state | Small rule set on every subtask | Must fit in memory × parallelism | Fraud 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
| Choice | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Short watermark (5s) | Fast dashboards | Results revised or incomplete for mobile/IoT sources | Dashboard consumers see numbers change |
| Long watermark (5m) | Complete first answer | 5-minute-stale "real-time" | Product: "real-time" becomes "near-real-time" |
| Long allowed lateness | Late data corrected | State and checkpoint size grow with it | Platform budget; recovery time |
| Frequent checkpoints (5s) | Small replay, low transactional latency | Checkpoint overhead, S3 request cost, pressure on large state | Job throughput; S3 bill |
| Infrequent checkpoints (5m) | Low overhead | Up to 5m replay; 5m output latency on transactional sinks | Recovery time; downstream freshness |
| Streaming-only (Kappa) | One code path | Reprocessing history requires replaying Kafka (retention!) | Whoever needs a 6-month backfill |
| Streaming + batch (Lambda) | Batch is the audited truth | Two implementations drift | The 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."
Anti-Patterns — What Kills Flink Deployments#
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#
| Dimension | Apache Flink | Kafka Streams | Spark Structured Streaming | Google Dataflow (Beam) | Streaming DBs (RisingWave, Materialize) |
|---|---|---|---|---|---|
| Model | True streaming, per record | Library inside your service | Micro-batch (default) | Beam model on managed runner | SQL materialized views, incrementally maintained |
| Latency | ms–sub-second | ms–sub-second | ~100ms–seconds per micro-batch | Seconds | Sub-second to seconds |
| State | RocksDB/heap, TB scale, S3 checkpoints | RocksDB + Kafka changelog topics | State store per micro-batch, checkpoint to HDFS/S3 | Managed | Managed, persisted |
| Event time | First-class, richest | Good (grace periods) | Good (watermarks) | First-class (Beam originated much of the model) | Supported |
| Deployment | Cluster (K8s operator, YARN), or managed (AWS Managed Flink, Confluent) | Just your app's pods | Spark cluster | Fully managed | Managed or cluster |
| Sources | Anything (Kafka, Kinesis, CDC, files) | Kafka only | Many | Many | Kafka, CDC, Postgres |
| Ops burden | Medium–high | Low (it's your service) | Medium (usually already have Spark) | Low | Low–medium |
| Pick when | Large state, complex event time, many sources, sub-second | Kafka-in/Kafka-out logic owned by one service team | Org is Spark-heavy; latency of seconds is fine; unify with batch | GCP shop; want zero ops | Team 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)#
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#
| Resource | Rough figure | Note |
|---|---|---|
| Stateless throughput per core | ~50K–500K records/s | Depends on (de)serialization; chaining avoids copies |
| Keyed windowed aggregation per core | ~10K–100K records/s with RocksDB | Serialization + LSM access dominate |
| RocksDB state per TaskManager | Hundreds of GB comfortably on local NVMe | Checkpoint upload bandwidth is the limit |
| Checkpoint interval | 10s–60s typical; minutes for TB state | Must complete well inside the interval |
| Incremental checkpoint delta | ~1–10% of state per checkpoint | Workload 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:
lastCheckpointDurationtrending towardcheckpointTimeout;numberOfFailedCheckpoints;lastCheckpointSizegrowth 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:
currentOutputWatermarknot advancing vs wall clock (alert whennow - 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:
backPressuredTimeMsPerSecondhigh upstream;busyTimeMsPerSecond~1000 on the bottleneck operator; Kafkarecords-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:
allowNonRestoredStatefor 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.
Operational Reality Matrix#
| Failure | Detection Signal | Blast Radius | Mitigation | Owner |
|---|---|---|---|---|
| Checkpoint timeouts | lastCheckpointDuration > 50% interval | One job, growing replay | Unaligned, TTL, fix backpressure | Job owner |
| Watermark stall | now - currentOutputWatermark > threshold | All windowed outputs of the job | Idleness, timestamp filter | Job owner; source team for bad clocks |
| Backpressure | busy/backpressured time; Kafka lag | Freshness of every downstream consumer | Scale bottleneck, async I/O | Job owner |
| Duplicate output | Reconciliation delta | Downstream reports, possibly billing | Recompute from raw | Job owner + data consumer |
| Upgrade breaks state | CI savepoint restore test | One job's history | State migration or backfill | Job owner |
| Cluster-wide failure (shared session cluster) | JobManager down, all jobs restarting | Every job on the cluster | Application mode, per-job clusters | Streaming platform |
When to Use vs. Alternatives#
| Need | Pick | Why |
|---|---|---|
| Results needed within hours; daily reports | Batch (Spark/SQL on the warehouse) | Cheaper, simpler, easier to reprocess |
| Kafka-in/Kafka-out enrichment owned by one team | Kafka Streams | No cluster; deploy with the service |
| Large keyed state, event-time windows, multiple sources, sub-second | Flink | Best-in-class state and time semantics |
| Already all-in on Spark, seconds OK | Spark Structured Streaming | Reuse platform and skills |
| Output is a continuously updated SQL view | Flink SQL or a streaming database | Declarative, less code to own |
| Simple per-event transform, no state | Consumer service or serverless | A stream processor is overkill |
| Billing-grade totals | Batch as truth, streaming as preview | Auditable recomputation |
When NOT to Use Flink#
- 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.
Operational Concerns#
What the On-Call Actually Does#
- Watches four signals per job: source lag (seconds), watermark lag, checkpoint duration/size, restart count. Those four catch ~90% of incidents.
- Upgrades via savepoints: take savepoint → deploy new version from savepoint → verify watermark advances and output resumes. Keep the previous savepoint for rollback.
- 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.
- 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.
- 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#
| Metric | Healthy | Alert |
|---|---|---|
| Source consumer lag (seconds) | < pipeline SLO | > SLO for 5m |
now - currentOutputWatermark | ≈ watermark delay | > 5× delay |
lastCheckpointDuration | < 30% of interval | > 50% of interval |
numberOfFailedCheckpoints rate | 0 | > 2 in 10m |
numRestarts | Flat | Any increase (ticket), > 3/h (page) |
lastCheckpointSize trend | Stable ±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#
Which Case Studies Use Flink & Stream Processing#
| Case Study | How Stream Processing Is Used | Key Pattern |
|---|---|---|
| Stream Processing | Ad click aggregation — the canonical question | Event-time windows, idempotent upsert, batch reconciliation |
| Leaderboard | Real-time score aggregation feeding a sorted set | Two-phase aggregation for hot keys |
| Metrics & Monitoring | Pre-aggregation before the TSDB | Rollups reduce cardinality |
| Search Indexing | CDC → denormalized documents | Temporal joins, upsert sink |
| Ride Hailing & Delivery | Supply/demand per geo cell for surge pricing | Sliding windows per cell, broadcast config |
| Payment Processing | Real-time fraud scoring | Keyed state per card + timers |
| Data Pipeline Patterns | Kappa vs Lambda decisions | Streaming 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 Say | What Interviewers Hear | What 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#
| Scenario | L5 Answer | L6 / Staff Answer | L7 / Principal Answer |
|---|---|---|---|
| "Ad click counts per minute" | Flink window, write to Redis | Event time, 10s watermark, 5m lateness, side output, dedupe on click_id with 24h TTL, upsert by (ad_id, minute), nightly batch overwrite + reconciliation | Plus: 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 resources | Check checkpoint duration/size, backpressure, state TTL; unaligned checkpoints | Why is a single job's failure paging a human? Platform SLOs, self-service diagnostics, restart budgets |
| "Upgrade Flink 1.x → 2.x" | Redeploy | Savepoint compatibility tests per job, UIDs, staged rollout | Fleet upgrade program across 200 jobs: compatibility matrix, owners, deadline, what gets rewritten in SQL |
| "Should we adopt streaming?" | Yes, it's modern | For the 3 pipelines with sub-hour needs; batch for the rest | Cost and headcount model; platform vs managed; SQL as the paved road |
The Staff Stream Processing Checklist#
- Latency need: "Product needs this within 2 minutes; batch can't do it, so streaming is justified."
- Time semantics: "Event time, 10s watermark, 5m allowed lateness, side output for the rest."
- Key and skew: "Key by ad_id with two-phase aggregation for hot ads."
- State budget: "~1.3GB for windows, ~80GB for 24h dedupe — RocksDB, TTLs, 30s incremental checkpoints."
- Sink semantics: "Upsert by
(ad_id, window_start)— replays are harmless." - 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#
| Dimension | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| Time | Windows | Event time, watermarks, lateness policy | Org lateness contracts per data product |
| Correctness | "Exactly-once" | State vs sink guarantees; idempotent/2PC sinks; reconciliation | Correctness tiers; single metric definitions across batch and stream |
| State | "Flink handles it" | Sized, TTL'd, backend chosen, recovery time computed | State as a long-lived migration liability |
| Operations | "Monitor the job" | Lag, watermark, checkpoint, restart SLOs | Platform model: application mode, self-service, upgrade programs |
| Choice | Flink for everything streaming | Flink vs Kafka Streams vs batch by need | Build vs managed vs retire; SQL paved road |
Strong hire signals
| Signal | What 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
| Signal | Why It Misses the Bar |
|---|---|
| Processing time for business metrics | Results change on replay |
| No late-data policy | Silent loss or unbounded state |
| Increment sinks with "exactly-once" claims | Double counts on every failure |
| No reconciliation for money | Unauditable |
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.
| Option | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Central platform, Flink SQL self-service | Consistent ops, upgrades once, metric definitions shared | Complex logic doesn't fit SQL; platform team becomes a bottleneck | Platform headcount |
| Platform runtime, teams write DataStream jobs | Flexibility | 200 bespoke jobs; upgrade program takes a year | Job owners at every upgrade; platform |
| Each team runs its own Flink/Kafka Streams | Autonomy | Duplicate definitions, uneven ops quality | Incident responders; data consumers who see conflicting numbers |
| Managed service (Dataflow, AWS Managed Flink, Confluent) | No runtime ops | Per-unit pricing at scale; feature lag; lock-in on runner specifics | Finance |
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.
| Scale | Throughput | Jobs | Infra/month | Headcount | On-call |
|---|---|---|---|---|---|
| Startup | 5K ev/s | 3–5 | ~$2–4K (managed Flink or small K8s) | 0.5 FTE | Owned by data team, rare pages |
| Growth | 200K ev/s | 30–60 | ~$25–50K | 2–3 FTE platform | Weekly incidents; platform rotation |
| Enterprise | 5M+ ev/s | 300+ | ~$300K–1M | 8–15 FTE | Dedicated 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#
One-Way Doors vs Two-Way Doors#
| Decision | Reversibility | Cost to Reverse |
|---|---|---|
| Max parallelism of a stateful job | One-way without state migration | State Processor API rewrite or full backfill |
| State serializer (Kryo vs Avro/Protobuf) | One-way-ish | State migration per job |
| Metric definitions exposed to finance | One-way once invoices depend on them | Restatements, customer communication |
| Engine choice (Flink vs Spark vs Dataflow) for 200 jobs | One-way at fleet scale | Year-long rewrite |
| Checkpoint interval, parallelism, watermark delay | Two-way | Config + savepoint restart |
| Managed vs self-hosted Flink | Two-way-ish | Savepoints 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#
| Signal | What 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 — Flink at Singles' Day Scale#
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