skip to content

How do you change the parallelism of a running Flink job, and what caps it?

level: seniorimportance: should knowfreq 48%

answer

  1. not in place, snapshot and restart
  2. state is partitioned one level of indirection deeper
  3. the ceiling is fixed when the job is first born
  4. choose a parallelism that divides evenly
  5. stable operator identity or no restore

basics

~20 s

Stop 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.

solid answer

~50 s

With Flink's default scheduler a streaming job's parallelism is fixed for the life of a run, so rescaling means a restart: `flink stop --savepointPath <dir> <jobId>` to take a savepoint and stop cleanly, then `flink run -s <savepointPath> -p <new parallelism> job.jar`. Keyed state survives because Flink hashes keys into **key groups** and assigns contiguous ranges of key groups to subtasks; on restore the ranges are redivided. The number of key groups equals the job's **maximum parallelism**, which is baked into the savepoint — so parallelism can never exceed it, and changing it invalidates keyed state unless you rewrite the savepoint with the State Processor API. If you never set it, Flink derives a value (at least 128, capped at 32768) from the initial parallelism, which is why long-lived jobs should set `pipeline.max-parallelism` deliberately on day one. You also need the slots: raising parallelism means more TaskManagers or more slots each.

code

bash · 4 lines
bash
./bin/flink stop --savepointPath s3://bucket/savepoints 4f2a1c...
# Savepoint completed. Path: s3://bucket/savepoints/savepoint-4f2a1c-9b

./bin/flink run -s s3://bucket/savepoints/savepoint-4f2a1c-9b -p 16 orders.jar

go deeper

for a junior

Know that a Flink job's parallelism is set at submit time with -p or in the code, and that changing it means stopping and restarting the job rather than editing it live.

for a middle

Explain the stop-with-savepoint and restore-with-new-parallelism procedure, and why key groups exist as an indirection between keys and subtasks.

for a senior

Be ready to run the operation on a production job: choose the maximum parallelism up front, keep stable UIDs, size the slot requirement, and account for state-reload downtime in the change window.

for a principal

Own the standard. Decide the default maximum parallelism and UID discipline for every job on the platform, because both are one-way doors, and set the policy on how often rescaling is worth its downtime versus sizing for peak.

## Why it is a restart Under Flink's default scheduler the ExecutionGraph is built once, with a fixed parallelism per operator, and subtasks are deployed into slots for the life of the run. There is no in-place resize. Changing parallelism therefore means: snapshot the state, stop, and start a new run that restores from the snapshot at the new parallelism. ## The procedure ```bash # 1. snapshot and stop cleanly ./bin/flink stop --savepointPath s3://bucket/savepoints <jobId> # 2. restore at the new parallelism ./bin/flink run -s s3://bucket/savepoints/savepoint-xxxx -p 16 job.jar ``` `stop` with a savepoint is preferable to `cancel` because it stops the sources first and takes a final, consistent snapshot, so no work is repeated. Do **not** pass `--drain` when you intend to resume: draining emits a final maximum watermark, firing every remaining window and advancing event time to infinity, which corrupts the semantics of the resumed job. A retained checkpoint can also be used as a restore point when no savepoint exists, but the savepoint is the deliberate, portable one. ## What actually limits the new parallelism Keyed state is not partitioned by subtask directly — it is partitioned by **key group**. Flink hashes each key into one of *N* key groups, and each subtask owns a contiguous range of them. That indirection is what makes rescaling possible: on restore, the same *N* key groups are simply divided among a different number of subtasks, and each subtask reads exactly the key-group ranges it now owns. *N* is the operator's **maximum parallelism** (`setMaxParallelism(...)` or `pipeline.max-parallelism`), and it is fixed in the savepoint. Two consequences: - **Parallelism can never exceed maximum parallelism.** Ask for more and the job fails to start. - **Maximum parallelism cannot be changed on restore.** Changing it changes which key goes to which key group, so the restored state would be wrong; Flink refuses. Escaping the ceiling means rewriting the savepoint offline with the State Processor API — an expensive, rarely-run operation. If you never set it, Flink derives a maximum parallelism from the initial parallelism, choosing a value of at least 128 and at most 32768. A job first launched at parallelism 4 therefore usually gets 128, which is fine until the day you need 200 subtasks. Setting it explicitly at job creation — high enough to be a real ceiling, not so high that per-key-group overhead and metadata grow needlessly — is the cheap decision that avoids the expensive one. It is also worth choosing a parallelism that divides the key-group count evenly. With 128 key groups and parallelism 12, some subtasks own 11 groups and others 10, an ~10% built-in imbalance; parallelism 8, 16 or 32 divides cleanly. ## Non-keyed state and operator identity Operator (non-keyed) state is redistributed on restore too, either by splitting the list evenly across the new subtasks or by giving every subtask the union of all entries, depending on how the operator declared it. Kafka source partition offsets are the everyday example of the even-split form. All of this depends on Flink being able to match state in the savepoint to operators in the new job. It does that by operator ID, which is generated from the graph structure unless you assign one with `uid("...")`. Without explicit UIDs, adding or removing an operator can shift the generated IDs and the restore fails or silently drops state. Assign a stable `uid()` to every stateful operator from the first deployment; add `--allowNonRestoredState` only when you have deliberately removed an operator and accept discarding its state. ## The resource side A higher parallelism needs more slots — the requirement is the maximum operator parallelism under default slot sharing. On native Kubernetes or YARN, Flink's ResourceManager will request more TaskManager containers; in a standalone cluster you must add TaskManagers yourself, or the job sits in scheduling and times out. Raising parallelism also raises the source ceiling problem: a Kafka source cannot usefully run wider than the topic's partition count, so scaling the job without repartitioning the topic leaves idle source subtasks. ## Cost of the operation The rescale is a full stop and restart. Downtime is savepoint duration plus job startup plus state reload — and with a large embedded RocksDB state, reload means downloading and rebuilding many gigabytes per subtask, which can dominate. That cost is the practical reason teams size long-running jobs for peak and rescale rarely, rather than tracking load closely.

  • Why can't the maximum parallelism be raised when restoring from a savepoint?
    Because it *is* the key-group count, and the key group for a key is derived from a hash modulo that count. Changing it would remap keys to different key groups, so the state a subtask restores would no longer be the state for the keys it now handles. Flink refuses the restore rather than silently corrupting results; the only escape is rewriting the savepoint with the State Processor API.
  • What is the risk of passing --drain when stopping with a savepoint?
    `--drain` advances event time to its maximum before stopping, firing every pending window and flushing all timers. That is what you want when retiring a pipeline for good, but if you then restore from that savepoint the job resumes with event time already at the end of time, so subsequent windows behave nonsensically. For a rescale, stop without draining.
  • Your job restores fine at the new parallelism but one operator's state comes back empty. What do you check first?
    Operator UIDs. Flink matches savepoint state to operators by ID, generated from the graph topology unless you set `uid()` explicitly. If the job graph changed — an operator added, removed or reordered — the generated ID shifts and the state no longer matches. Check whether the operator has a stable `uid()`, and whether `--allowNonRestoredState` is quietly masking the mismatch.

Keys are sorted into a fixed number of pigeonholes when the job is created. Rescaling reassigns whole pigeonholes to more or fewer clerks; you can never have more clerks than pigeonholes, and you cannot change how many pigeonholes there are without re-sorting every letter.

saying these in an interview costs you the question

  • Claiming a running Flink job can be resized in place by default
  • Thinking maximum parallelism can be raised on restore
  • Using cancel instead of stop-with-savepoint before a rescale
  • Passing --drain when the job will be resumed
  • Omitting uid() on stateful operators and blaming the savepoint

context