Apache Flink
Flink is the stream-first engine: streaming is the native model and batch is the special case, with real managed state, event-time semantics and exactly-once guarantees. Interviewers bring it up whenever latency matters, usually as a comparison against Spark Structured Streaming.
on this pageshowhide
guide
overview
~1 minApache Flink is a distributed engine for stateful computation over streams. It treats an unbounded stream as the normal case and a bounded dataset as a stream that ends, and it keeps application state inside the engine, snapshotted with the job. Interviewers reach for it whenever a design needs low latency and correct numbers at once, and they often frame the conversation as a comparison with Spark Structured Streaming. A good Flink answer says which notion of time a result depends on, where its state lives, and what survives a restart. The hub follows the life of a job. The [DataStream API](/topics/data-flink-datastream-api) is how a pipeline is written. [Event time and watermarks](/topics/data-flink-event-time-watermarks) and [windowing](/topics/data-flink-windowing) decide when a result is ready. [State management](/topics/data-flink-state-management) and [checkpointing and exactly-once](/topics/data-flink-checkpointing-exactly-once) decide what is remembered and what a failure costs. [Deployment and scaling](/topics/data-flink-deployment-scaling) is the operator's view. [Flink SQL and the Table API](/topics/data-flink-table-api-sql) and [SQL joins and windows](/topics/data-flink-sql-streaming-operators) express the same ideas declaratively, as tables that keep changing. Junior rounds check the vocabulary: keyed streams, event time, window types, what a checkpoint is for. Senior rounds turn into incidents — a window that never fires, state that grows until checkpoints fail, a job that cannot keep up — and ask which signal you would read first. Start with the DataStream API and keyed streams, then event time. State, checkpoints and SQL all assume both.
primer
### Streams first, batch as a special case A Flink job is a graph of long-lived operators that records flow through. Bounded input runs on the same model, with optimisations for the fact that it ends. Answers that treat Flink as a faster batch scheduler miss what makes it different. ### Keys decide where work and memory go Partitioning a stream by key sends all records for that key to one parallel instance, and that instance owns the key's state and timers. Almost every stateful feature — keyed state, per-key windows, timers — hangs off this step, so the choice of key is also a choice about skew and memory. ### Time has two meanings **Event time** follows timestamps inside the data; **processing time** follows the wall clock of the machine. Event time is what lets a replay or a backfill reproduce the original answer. The cost is that the engine has to be told how far out of order data may arrive, and a **watermark** is that statement, carried through the pipeline alongside the records. Most "my window never fires" stories are watermark stories. ### State is part of the job Flink keeps counters, buffers and lookups in a managed **state backend** — on the heap or on local disk — rather than in an external database. Access stays local and fast, and state size, expiry and schema changes become design questions. State nobody expires only grows. ### A checkpoint is a consistent cut Periodically the engine snapshots operator state together with how far each source has read, all as of one logical point in the stream. Recovery rewinds to that point and replays. Inside Flink this gives exactly-once state; whether the **output** is exactly-once depends on what the sink can do with a replay. Interviewers expect you to draw that boundary before describing the protocol. ### SQL is the same engine, seen as tables A streaming query never finishes, and its result is a table that keeps changing. Some operators only append rows; others must retract or update earlier ones, and that decides which sinks can accept the output and how much state the query carries.
- DataStream
- Flink's core abstraction for a possibly unbounded sequence of records, transformed by chained operators into a dataflow graph that runs as a job.
- KeyedStream
- A stream partitioned by a key selector, so every record with a given key reaches the same parallel instance. Required for keyed state, timers and per-key windows.
- Event time
- Time taken from a timestamp inside each record, describing when the event happened. Results based on it are reproducible regardless of arrival order.
- Processing time
- Time taken from the clock of the machine running an operator. Simple and low-latency, but results change with load, delays and replays.
- Watermark
- A marker flowing with the records that advances an operator's event-time clock and tells it older records are no longer expected.
- Keyed state
- State scoped to the current record's key, such as ValueState, ListState or MapState, reachable only on a KeyedStream.
- State backend
- The component that stores working state and writes it into snapshots: on the JVM heap, or serialized in an embedded RocksDB on local disk.
- Checkpoint
- An automatic, periodic, consistent snapshot of all operator state and source positions, used to recover a job after a failure.
- Savepoint
- A snapshot triggered on purpose in a portable format, used to upgrade, rescale or move a job while keeping its state.
- Maximum parallelism
- The upper bound on an operator's parallelism, which fixes the number of key groups keyed state is divided into.
- Backpressure
- The condition in which a slow operator fills its input buffers and upstream operators are forced to slow down to match it.
- Changelog stream
- A stream of row changes — inserts, update pairs and deletes — through which Flink SQL represents a table that keeps changing.
Trace a job from the code to the cluster. A program builds a dataflow graph: sources, transformations, a key partition, windows or process functions, sinks. Nothing runs until the program submits it. The JobManager turns that graph into parallel subtasks and places them in TaskManager slots, and from then on records move directly between TaskManagers over the network. Each source assigns timestamps and emits watermarks; downstream, an operator's event-time clock follows the slowest of its inputs, and that clock fires windows and timers. Keyed operators read and write their state in the state backend. On a schedule, the JobManager asks the sources to inject checkpoint barriers, each operator snapshots its state as a barrier passes, and a checkpoint is complete once every task has reported. A transactional sink commits its pending output only then. A few lines of a DataStream job touch most of the hub: ```java env.enableCheckpointing(60_000); // consistent snapshot every minute env.fromSource(kafkaSource, WatermarkStrategy.<Click>forBoundedOutOfOrderness(Duration.ofSeconds(10)), "clicks") // event time, 10 s of tolerated disorder .keyBy(click -> click.userId) // one owner per user: state, timers, windows .window(TumblingEventTimeWindows.of(Duration.ofMinutes(1))) .aggregate(new CountClicks()) // one accumulator per user and window .sinkTo(kafkaSink); // end-to-end guarantee depends on this sink ``` Each line is a decision an interviewer can pull on: the bound trades latency for completeness, the key sets skew and state size, and the sink decides whether a replay after failure shows up as duplicates. Flink SQL compiles to the same runtime: a table definition declares the time attribute and watermark, the planner picks stateful operators, and the sink must accept whatever changelog the query produces. Rescaling either kind of job usually goes through a savepoint, and maximum parallelism caps it.
- DataStream API →
How a job is written and submitted: transformations, key partitioning, sources, sinks and parallelism. Every other section assumes it.
- Event Time & Watermarks →
The idea Flink questions return to most: which clock a result follows and how the engine decides data is complete.
- Windowing →
Windows turn event time into results; choosing a window type and handling late records follows directly from watermarks.
- State Management →
What keyed state is, where it is stored, and how it is bounded and evolved as a job changes.
- Checkpointing & Exactly-Once →
How state survives failure, and where the exactly-once guarantee stops once output leaves Flink.
- Flink SQL & Table API →
The declarative layer, read after the core: streams as changing tables, and what an updating result demands of a sink.
Treating the watermark as a timer: it moves only when data moves, so a silent input can hold downstream windows open until idleness is configured.
Claiming exactly-once end to end without naming the sink; checkpoints roll back Flink's own state, not rows already written to an external system.
Choosing processing time for a metric that must match a later backfill; replaying the same data can land records in different windows.
Confusing checkpoints with savepoints: one is the engine's recovery mechanism, the other is the operator's tool for upgrades and rescaling.
Restarting or adding slots to a job that is falling behind without first finding the operator that is actually the bottleneck.
Reading a streaming GROUP BY or regular join as if it returned final rows; its output updates earlier rows, and its state can grow without bound.
This guide assumes the Flink 2.x line. Much production code, and many answers found elsewhere, were written against 1.x, so an interviewer may probe what moved: - **1.12** let DataStream programs run in a batch execution mode over bounded input, the first step toward retiring the separate batch API. - **1.13** renamed the state backends to the heap-based and RocksDB-based ones used today, separating where working state lives from where checkpoints are stored. - **2.0** removed the DataSet API, the Scala APIs and the old SourceFunction/SinkFunction interfaces, so batch work goes through DataStream or SQL; Queryable State still ships but stays deprecated. When an answer depends on an API that changed — how a source is defined, how windows take a duration, how state is read from outside — say which line you are describing. Naming the era is part of a correct answer, not a detail.
Flink usually sits between a log and a store. Apache Kafka is the common source and sink; change-data-capture tools such as Debezium feed it database changes; table formats such as Apache Iceberg and Apache Paimon hold its output for querying. The comparison interviewers ask for most is Spark Structured Streaming. Both engines unify batch and streaming and both offer exactly-once state; the separating trade-off is that Spark processes streams as a series of small batches by default, while Flink processes records one at a time through long-running operators, which favours lower latency and fine-grained timers. Kafka Streams is the other common comparison: it is a library embedded in an application, with no cluster to run, and it fits when both input and output are Kafka topics. Flink fits when state, sources or event-time logic outgrow that.
explore
- DataStream API6 questions
- Event Time & Watermarks6 questions
- Windowing6 questions
- State Management7 questions
- Checkpointing & Exactly-Once6 questions
- Deployment & Scaling6 questions
- Flink SQL & Table API6 questions
- Flink SQL Joins & Windows6 questions
questions
page 2 of 2In Flink windowing, what is the difference between a Trigger returning FIRE and FIRE_AND_PURGE?
basics
~20 sBoth emit a result for the window. FIRE leaves the window's buffered elements in place, so a later firing recomputes over them again; FIRE_AND_PURGE clears the contents afterwards, so the next firing starts from empty. Neither removes the window itself.
When does a Flink DataStream pipeline need process() with a KeyedProcessFunction rather than flatMap()?
basics
~20 sReach for process() when the logic needs keyed state, timers, the record's event-time timestamp, or side outputs. flatMap() only transforms one record into zero or more records with no memory of what came before and no way to act at a future time.
In a Flink job, how do you find which operator is causing backpressure?
basics
~10 sWalk the Flink job graph downstream and find the first task that is busy but not backpressured — that is the bottleneck. Everything upstream of it shows backpressure, everything downstream shows idle time.
How do you change the parallelism of a running Flink job, and what caps it?
basics
~20 sStop the Flink job with a savepoint, then resubmit from that savepoint with a new -p value. The hard ceiling is the job's maximum parallelism, which fixes the number of key groups and cannot be changed without rewriting the state.
A Flink event-time job stops emitting window results when one Kafka partition goes quiet. Why?
basics
~20 sWatermarks are generated per input and combined by taking the minimum, so a silent partition's watermark never advances and pins the whole operator's event time. Configure withIdleness on the WatermarkStrategy so idle inputs are excluded from that minimum.
In Flink, what do allowedLateness and sideOutputLateData do with records behind the watermark?
basics
~20 sallowedLateness keeps a window's state alive past its watermark deadline, so a straggler triggers an extra firing with an updated result. sideOutputLateData routes records too late even for that grace period to a separate stream instead of dropping them silently.
A Flink SQL GROUP BY country over an orders stream backpressures because one country dominates. How do mini-batch and two-phase aggregation help?
basics
~20 sMini-batch (table.exec.mini-batch.enabled plus allow-latency and size) buffers rows so each key's state is touched once per batch. On top of it, table.optimizer.agg-phase-strategy = 'TWO_PHASE' pre-aggregates before the shuffle, so the hot key's task receives accumulators instead of raw rows.
How do you add a field to a POJO stored in Flink ValueState without losing the state?
basics
~20 sTake a savepoint, add the field to the POJO, and restore from that savepoint. Flink evolves POJO and Avro state schemas automatically: added fields get their Java default value, removed fields are dropped, but declared field types and the class name cannot change.
A Flink job's keyed state grows without bound across millions of keys. How do you bound it?
basics
~20 sAttach a StateTtlConfig to the state descriptors so entries expire, and back it with explicit cleanup: a KeyedProcessFunction timer that calls state.clear() when a key goes idle. Flink never garbage-collects a key just because it stopped appearing.
A Flink 2.3 INSERT INTO fails at planning because the query's upsert key differs from the sink's primary key - what is the risk, and how do you resolve it?
basics
~20 sSeveral result rows with different upsert keys can land on one sink key, so the stored value depends on arrival order. Flink 2.3 fails planning unless you fix the key mismatch or add ON CONFLICT DO ERROR, DO NOTHING or DO DEDUPLICATE.
A Flink job with SlidingEventTimeWindows of 24 hours and a 1-minute slide exhausts state — why?
basics
~20 sA 24-hour window sliding every minute keeps 1440 windows open per key at once, and every record is assigned to all of them. State is multiplied roughly 1440-fold per key, so checkpoints balloon and the job runs out of memory or disk.
When would you run a Flink pipeline at-least-once with idempotent sinks instead of exactly-once?
basics
~20 sWhen the sink is naturally idempotent — keyed upserts, partition overwrites — replayed records overwrite rather than duplicate, so at-least-once gives the same final result without alignment stalls, transaction machinery, or output that is invisible until the next checkpoint completes.
How do you choose a Flink job's maxParallelism, and why can't you change it later?
basics
~20 sMaximum parallelism fixes the number of key groups, the atomic unit by which Flink redistributes keyed state. Changing it changes the key-to-key-group mapping, so existing state can no longer be assigned and must be discarded. Set it once, generously, and never touch it.
Why does a Flink flatMap() written as a Java lambda throw InvalidTypesException?
basics
~10 sJava erases the generic parameter of the Collector in flatMap's signature, so Flink cannot infer the output type from the lambda. Supply it explicitly with .returns(...) or implement FlatMapFunction as a class.
Can an external service query a Flink job's keyed state directly, and should it?
basics
~20 sTechnically yes, via Queryable State, but it should not: Flink 2.3 ships it only as a deprecated, off-by-default leftover slated for removal. Publish readable state to an external store, or read a savepoint offline with the State Processor API.
In a Flink windowed aggregation, why does adding an Evictor disable pre-aggregation?
basics
~20 sAn evictor inspects and removes individual elements after the trigger fires, so every element must still exist when the window is evaluated. That forces Flink to buffer the whole window instead of folding records into a single accumulator on arrival.
A Flink job's checkpoints time out under heavy backpressure — how do unaligned checkpoints help?
basics
~20 sUnder backpressure a barrier crawls through full network queues and alignment stalls the snapshot. Unaligned checkpoints let the barrier overtake queued records and persist the in-flight buffers as part of the snapshot, so checkpoint duration stops depending on throughput.
In Flink, what problem does WatermarkStrategy.withWatermarkAlignment solve?
basics
~20 sIt stops one source split from racing far ahead of the others in event time. Flink pauses reading from splits whose watermark exceeds the group's minimum by more than the configured drift, bounding the state that accumulates while slow splits catch up.
When would you run a Flink job in reactive mode instead of rescaling it from savepoints?
basics
~20 sReactive mode suits a single standalone application-mode job with modest state and a strong daily traffic swing, where an external controller adds and removes TaskManagers. Steady jobs or very large state are better served by sizing for peak and rescaling rarely.
showing 31–49 of 49