skip to content

When should a Flink DataStream job run with execution.runtime-mode set to BATCH instead of STREAMING?

level: middleimportance: must knowfreq 62%

answer

  1. one property of the sources decides eligibility
  2. stages instead of everything online at once
  3. the state backend stops mattering
  4. sorting replaces the watermark heuristic
  5. one fault-tolerance mechanism simply is not there

basics

~20 s

Use BATCH only when every source is bounded. It lets Flink run the job stage by stage with materialized shuffles, sort by key instead of keeping all state live, and recover by backtracking. STREAMING is the default and the only option for unbounded input.

solid answer

~50 s

`execution.runtime-mode` takes `STREAMING` (the default), `BATCH`, or `AUTOMATIC`, which picks based on source boundedness. BATCH is legal only if **all** sources are bounded; one unbounded source makes the job unbounded and BATCH unavailable. When it applies, BATCH is the more efficient choice: Flink executes the job in stages separated by shuffles, materializes intermediate results instead of pipelining them, and can therefore run on fewer slots and backtrack to the last completed stage on failure rather than restarting everything from a checkpoint. Keyed operations sort by key and hold only one key's state at a time, so the configured state backend is ignored. The tradeoffs are real: checkpointing is unsupported in BATCH, so anything depending on it breaks, and rolling operators like `reduce()` emit only the final value per key instead of an update per record.

code

bash · 3 lines
bash
# preferred: choose the mode at submission, keep the JAR mode-agnostic
bin/flink run -Dexecution.runtime-mode=BATCH job.jar
bin/flink run -Dexecution.runtime-mode=AUTOMATIC job.jar

go deeper

for a junior

Know that Flink runs the same DataStream code in two modes, that BATCH needs bounded sources, and that STREAMING is the default. Naming the execution.runtime-mode setting is enough at this level.

for a middle

Explain the mechanics: staged execution with materialized shuffles versus pipelined ones, sorting by key instead of live keyed state, perfect watermarks, and the fact that checkpointing is unavailable in BATCH.

for a senior

Show you have hit the limits in production — a sink that silently stopped committing because it depended on checkpoints, or a downstream consumer that expected rolling reduce updates and got one final value per key.

for a principal

Own the platform decision: whether one artifact serves both backfill and live paths with the mode chosen at submission, and what that implies for sink selection, since only unified-Sink-API sinks behave transactionally in both modes.

## Two runtime modes over one API Flink's DataStream API is deliberately unified: the same program can execute in `STREAMING` mode — the classic behaviour, continuous incremental processing, expected to stay online indefinitely — or in `BATCH` mode, which executes more like a batch framework. Flink's promise is that a DataStream application over *bounded* input produces the same **final** results in either mode. The word final matters: a STREAMING run may emit incremental updates along the way (think upserts), while a BATCH run emits one result at the end. The mode is set with `execution.runtime-mode`, which accepts three values: `STREAMING` (default), `BATCH`, and `AUTOMATIC`, where Flink decides from the boundedness of the sources. ```bash bin/flink run -Dexecution.runtime-mode=BATCH job.jar ``` You can also call `env.setRuntimeMode(RuntimeExecutionMode.BATCH)` in code, but the documentation explicitly recommends against it: keeping the application configuration-free lets the same JAR run in either mode depending on how it is submitted. ## The boundedness rule Boundedness is a property of a *source*: does it know all its input up front, or may new data keep arriving indefinitely? A job is bounded only if **every** source is bounded. BATCH mode can only be used for bounded jobs. STREAMING works for both, which is why it is the safe default. A bounded file source is bounded. A Kafka source configured with an end offset is bounded; the same source reading to the end of the topic and staying is not. This is why `AUTOMATIC` is attractive for pipelines whose source configuration varies between environments. ## What changes underneath **Scheduling and shuffles.** In STREAMING mode all tasks are online simultaneously, and shuffles are *pipelined*: records go to downstream tasks immediately, with buffering at the network layer. That is required for low latency, and it means the cluster must have enough slots to run every task at once. In BATCH mode the job splits into stages at the shuffle boundaries. Flink fully processes one stage, materializes its output to non-ephemeral storage, and then runs the next — so upstream tasks can go offline and the job can complete on fewer slots than it has tasks. **State.** In STREAMING mode the configured state backend governs how keyed state is stored. In BATCH mode the state backend is **ignored**. Instead the input to a keyed operation is grouped by key using sorting, and all records for one key are processed in turn, so only one key's state needs to be resident; that key's state is discarded when the run moves to the next key. That is why a BATCH run can process a key space far larger than would fit in a streaming job's state. **Time and watermarks.** STREAMING mode uses watermarks as a heuristic about out-of-orderness. BATCH mode can sort by timestamp instead, so it behaves as if watermarks were perfect. Flink only needs a `MAX_WATERMARK` at the end of each key's input; user `WatermarkGenerator` implementations are ignored. A `WatermarkStrategy` is still worth supplying, because its `TimestampAssigner` still assigns record timestamps. Processing-time timers likewise all fire at the end of input. **Recovery.** STREAMING mode recovers by restarting the running tasks from a checkpoint. BATCH mode backtracks to previous stages whose intermediate results still exist, so potentially only the failed task and its predecessors restart. That is cheaper, and it is one of the stated reasons to prefer BATCH when the job allows. ## What breaks in BATCH mode The documentation is blunt about the limits. **Checkpointing does not work in BATCH mode**, and neither does anything that waits for a checkpoint: a `CheckpointListener` never hears `notifyCheckpointComplete()`, and a legacy `SinkFunction` that commits on that callback — one built on the deprecated `TwoPhaseCommitSinkFunction`, say — never gets the signal to commit. A transactional sink that must work in BATCH has to be built on the unified Sink API instead. In Flink 2.3 the file sink (`FileSink`) and the Kafka sink (`KafkaSink`) already are: at end of input the sink writer flushes and the committer commits everything it holds, so their output still lands. What changes is the timing — a `FileSink` with an on-checkpoint rolling policy never rolls mid-run, and each open part file is closed and committed once, when input ends. Rolling operations change shape. In STREAMING mode `reduce()` or `sum()` on a keyed stream emits an updated value for every arriving record; in BATCH mode they are not rolling and emit only the final result per key. A downstream consumer written to expect the update stream will see far fewer records. Custom operators are the sharp edge. Because BATCH processes key by key, the watermark switches from `MAX_VALUE` back to `MIN_VALUE` between keys — so an operator that caches the last seen watermark and assumes it only ascends will misbehave. Timers fire in key order first and timestamp order within a key. Changing a key manually inside an operator is unsupported. For most use cases the guidance is to use a process function rather than a custom operator, precisely to avoid these assumptions. ## Choosing The rule of thumb is simple: if the program is bounded, use BATCH, because it is more efficient; if it is unbounded, you must use STREAMING. The notable exception is deliberately running a bounded job in STREAMING mode to bootstrap state — run it, take a savepoint, restore that savepoint into the unbounded job — and running bounded sources in STREAMING mode when writing tests for code destined for unbounded input.

  • Why does a Flink BATCH-mode job ignore the configured state backend?
    BATCH sorts the input of a keyed operation by key and processes all records for one key before moving on, so only that key's state has to be resident and it is dropped when the run advances. There is nothing to store across the whole key space, so the state backend's job — managing large live keyed state and snapshotting it — does not apply.
  • A bounded job needs to bootstrap state for a long-running streaming job. Which mode should it use?
    STREAMING, despite the input being bounded. BATCH mode has no checkpointing and therefore no savepoint to hand off. Running the bounded job in STREAMING mode lets you take a savepoint at the end and restore it into the unbounded job. It is one of the explicit exceptions to the otherwise-use-BATCH rule.
  • Why does the documentation recommend setting the runtime mode at submission rather than in code?
    Keeping the application configuration-free means the same JAR can run in either mode. `bin/flink run -Dexecution.runtime-mode=BATCH` decides at submission time, so a backfill and the live pipeline can share one artifact. Hard-coding `env.setRuntimeMode(...)` pins the job to one mode and forces a rebuild to change it.
  • What is the AUTOMATIC value for, and why not always use it?
    AUTOMATIC lets Flink pick based on whether all sources are bounded, which suits a pipeline whose source configuration differs per environment. The cost is that the execution model — and therefore whether checkpointing exists and whether rolling operators emit updates — becomes a property of configuration rather than an explicit choice, which can surprise operators reading the job graph.

saying these in an interview costs you the question

  • Says BATCH mode works fine with an unbounded Kafka source
  • Claims checkpointing still runs in BATCH mode
  • Thinks BATCH and STREAMING need different application code
  • Assumes reduce emits per-record updates in BATCH mode
  • Believes the state backend still governs BATCH keyed state

context