When would you run a Flink job in reactive mode instead of rescaling it from savepoints?
answer
- the job takes whatever slots exist
- the mechanism, not the policy
- every resize is a restart from a checkpoint
- large state makes elasticity expensive
- CPU is the wrong signal for a stream
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.
solid answer
~50 sReactive mode (`scheduler-mode: reactive`, built on Flink's Adaptive Scheduler) drops the fixed-parallelism contract: the job uses whatever slots are registered, ignoring the configured parallelism so that only each operator's maximum parallelism caps it, and rescales itself when TaskManagers appear or disappear. You still need an external controller — a Kubernetes HPA on the TaskManager deployment, or a purpose-built autoscaler — to decide *when* to add capacity. On Flink 2.3 the limits are real: standalone application-mode clusters only, one job per cluster, and every rescale is a restart from a checkpoint. It pays off when idle capacity is expensive, load swings predictably, and state is small enough that a restart costs seconds; with multi-terabyte RocksDB state, sizing for peak is cheaper than flapping. For most platforms the Flink Kubernetes Operator's autoscaler, which reads busy time and lag and sets per-operator parallelism, is the more usable path.
code
yaml · 7 linesjobmanager.scheduler: adaptive
scheduler-mode: reactive
execution.checkpointing.interval: 10s
# configured parallelism is ignored in reactive mode;
# max parallelism is the ceiling, TaskManager count the lever
pipeline.max-parallelism: 128go deeper
Know that a Flink job normally runs at a fixed parallelism and that changing it requires a restart; reactive mode is an opt-in setting that removes the fixed part.
Explain what the Adaptive Scheduler and reactive mode change — parallelism follows available slots — and name the structural limits: standalone, application mode, one job per cluster.
Be ready to cost a rescale honestly: restart gap, reprocessing since the last checkpoint, and state reload, and to argue which metric should drive scaling for a stream.
Own the elasticity policy. Decide whether the platform autoscales streaming jobs at all, what signal drives it, what damping prevents flapping, and when overprovisioning for peak is simply the cheaper and safer engineering choice.
## What reactive mode actually is Normally a Flink job declares a parallelism, and the scheduler demands exactly that many slots before it will run. The **Adaptive Scheduler** (`jobmanager.scheduler: adaptive`) inverts that: it looks at the slots available and decides a parallelism that fits, up to the configured parallelism, then restarts the job at a new parallelism if the available slots change. **Reactive mode** is the configuration of the Adaptive Scheduler that takes this to its conclusion (`scheduler-mode: reactive`): the job simply uses *all* registered slots. The configured parallelism — set on the job or on individual operators — is ignored; the scheduler asks for each operator's **maximum parallelism**, which is the only ceiling you still control. Scaling the job means scaling the TaskManager pool — start another TaskManager and the job restarts wider; kill one and it restarts narrower, restoring from the latest checkpoint each time. ## What it does not do It does not decide when to scale. Reactive mode is the *mechanism*; the *policy* is external. On Kubernetes that is typically a HorizontalPodAutoscaler on the TaskManager Deployment, driven by CPU or a custom metric. That gap is the first thing to be honest about in an interview: reactive mode alone is not autoscaling, and CPU is a poor proxy for whether a streaming job is keeping up. It also carries hard structural limits: standalone deployments only, application mode only, and exactly one job per cluster. It does not apply to session clusters or to native Kubernetes/YARN deployments where Flink itself allocates containers. ## What a rescale costs Every rescale is a restart. The job stops, a new ExecutionGraph is scheduled at the new parallelism, and state is restored from the latest checkpoint. Concretely you pay: - the **gap** — no processing during restart, so lag grows exactly when you were scaling because lag was growing; - **reprocessing** of everything since the checkpoint it restores — on Flink 2.3 a resource-driven rescale is timed to follow a completed checkpoint (falling back after repeated checkpoint failures or a maximum delay), so a scale-up replays little, but a scale-down forced by a lost TaskManager is a failover and replays everything since the last checkpoint; - **state reload**, which with an embedded RocksDB backend means downloading and rebuilding state per subtask; small state is seconds, large state is minutes; - **key-group reshuffling**, still bounded by the job's maximum parallelism, which caps how far up any autoscaler can go. Multiply that by a flapping controller and the job spends its life restarting. Any policy needs a stabilization window and asymmetric thresholds — scale up quickly, scale down slowly and rarely. ## The judgment call Autoscaling a streaming job is worth it when three things hold at once: idle capacity is genuinely expensive, the load pattern has a large and predictable swing (a consumer-facing pipeline that is quiet overnight), and state is small enough that a restart is cheap. Miss any one and fixed parallelism wins. A pipeline at a steady rate simply has nothing to autoscale to; a pipeline with terabytes of keyed state pays more in reload than it saves in machines; a pipeline whose cost is dominated by the sink or by a fixed Kafka partition count cannot use extra subtasks anyway. The alternative is unglamorous and usually correct: size for peak plus headroom, alert on consumer lag and busy time, and rescale deliberately from a savepoint during a change window when the measured peak has moved. That is predictable, is easy to reason about during an incident, and costs one planned restart a quarter instead of several unplanned ones a day. ## Choosing a signal If you do autoscale, pick the signal carefully. CPU utilization is what an HPA offers by default and is close to meaningless for a job blocking on external I/O or on state access — a saturated pipeline can sit at 30% CPU. The signals that actually track "am I keeping up" are consumer lag and its derivative, and Flink's own busy/backpressured time per operator. The Flink Kubernetes Operator's autoscaler is built on exactly those: it estimates each operator's true processing rate, computes a target parallelism per operator, and applies it — on current Flink through the Adaptive Scheduler's `PUT /jobs/<job-id>/resource-requirements` endpoint, still a restart but with no savepoint round trip — a far better fit than a generic CPU-driven HPA, and the reason most teams reaching for elasticity today go through the operator rather than configuring reactive mode themselves. ## How to present the decision A strong answer names the mechanism accurately, then refuses to treat it as a default. Reactive mode is a real, narrow tool; the interesting content is the cost model of a restart, the fact that the scaling *policy* is yours to supply, and the observation that for most long-running stateful pipelines the cheapest form of elasticity is not scaling at all.
- Why is a CPU-based HorizontalPodAutoscaler a poor driver for a Flink job?Because CPU does not measure whether the job is keeping up. A subtask blocked on a synchronous lookup or on RocksDB disk reads is fully saturated at low CPU, so the autoscaler sits idle while lag grows; conversely a burst of cheap records spikes CPU without any backlog. Lag, lag derivative, and Flink's per-operator busy and backpressured time track the real question.
- What upper bound still applies to a job running in reactive mode?The job's maximum parallelism — the key-group count fixed when the job was created. Reactive mode ignores the configured parallelism and asks for every operator's maximum parallelism, so no amount of added TaskManagers takes an operator past it. Elasticity is therefore bounded by a decision made on day one, which is another reason to set maximum parallelism deliberately rather than accept the derived default.
- How would you keep an autoscaled Flink job from flapping?Asymmetric, damped policy: scale up on a sustained signal within a minute or two, scale down only after a long stabilization window and with a wide deadband, cap the number of rescales per hour, and never scale on a single noisy sample. Because each rescale costs a restart with state reload and reprocessing, the cost of a wrong decision is far higher downward than upward.
saying these in an interview costs you the question
- Calling reactive mode a complete autoscaler with no external controller
- Claiming rescaling happens without restarting the job
- Ignoring state reload cost when arguing for elasticity
- Assuming reactive mode works in a session cluster
- Autoscaling a streaming job on CPU utilization alone