How do you choose a Flink job's maxParallelism, and why can't you change it later?
answer
- State does not rescale key by key
- There is a fixed number of buckets
- That number is set by one job setting
- Change it and the mapping no longer lines up
- 2^15 is the documented ceiling
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.
solid answer
~50 sFlink organises keyed state into **key groups**, and there are exactly as many key groups as the operator's maximum parallelism. Each parallel subtask owns a contiguous range of key groups, so rescaling is just reassigning ranges — that is what makes stateful rescaling cheap. A key's group is derived from its hash modulo the number of key groups, so changing maximum parallelism changes where every key belongs; the persisted key groups no longer line up and there is no way to restore that operator's state. The constraint is `0 < parallelism <= maxParallelism <= 2^15`. If you never call `setMaxParallelism`, Flink picks one at first start — `nextPowerOfTwo(p + p/2)`, floored at 128 and capped at 2^15 — and freezes it into the savepoint. Set it explicitly, well above any parallelism you expect to need, but not absurdly high: rescaling metadata grows linearly with it.
code
java · 10 linesStreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(8);
// Irreversible once the first savepoint exists -- set it deliberately,
// above any parallelism this job could ever need.
env.setMaxParallelism(1024);
stream.keyBy(Event::userId)
.process(new ProfileUpdater())
.uid("profile-updater");go deeper
Recall that a Flink job has both a parallelism and a maximum parallelism, that the second is a ceiling on the first, and that it cannot be changed once the job has run.
Explain the mechanics: maximum parallelism equals the number of key groups, key groups are the atomic unit of keyed-state redistribution, and a subtask owns a contiguous range of them.
Show that you have hit this in production — the silently defaulted 128, the uneven key-group split when the numbers do not divide, and the State Processor API rewrite that gets a job out of the corner.
Own it as capacity planning: bound the lifetime worst-case parallelism including backfills, pick a divisible value with headroom, mandate that it and operator uids are set explicitly before the first deployment, and define the migration and rollback plan for when the number turns out to be wrong.
## Key groups: the unit of rescaling When you `keyBy` a stream, Flink does not track state per individual key for redistribution purposes. It hashes each key into one of a fixed number of **key groups**, and the key group — not the key — is the atomic unit by which keyed state can be redistributed. There are exactly as many key groups as the *maximum parallelism* defined for the operator, and at runtime each parallel instance of a keyed operator works with the keys of one or more whole key groups. That indirection is what makes stateful rescaling tractable. Going from parallelism 4 to 8 does not mean rehashing billions of keys; it means splitting contiguous key-group ranges and shipping the corresponding state files. A subtask's assignment is a range, and the range boundaries are the only thing that moves. ## Why the number is frozen A key's group is a deterministic function of its hash and the total number of key groups. Change the total and every key potentially lands in a different group. The savepoint on disk is organised by key group, so after such a change Flink cannot map the stored groups onto the new layout — the mapping is simply a different function. Hence the documented rule: there is currently **no way to change the maximum parallelism of an operator after a job has started without discarding that operator's state**. This is not a warning about degraded behaviour; the restore fails or the state is dropped. It is one of the few genuinely irreversible decisions in a Flink deployment, which is why it belongs on a production-readiness checklist rather than being discovered later. ## The constraints and the default The rule is `0 < parallelism <= maxParallelism <= 2^15`. You set it explicitly with `setMaxParallelism(int)`, at job or operator granularity. If you set nothing, Flink decides at first start as a function of the operator's parallelism. Flink 2.3 computes `MIN(MAX(nextPowerOfTwo(parallelism + (parallelism / 2)), 128), 2^15)`, which works out as: - `128` for any parallelism up to 85; - `256` for parallelism 86 to 171; - beyond that, the smallest power of two at or above one and a half times the parallelism, capped at 2^15. The production-readiness page still says 128 for any parallelism up to 128; the code already gives 256 at parallelism 86. The trap is obvious once stated: a job that launched at parallelism 4 quietly received a ceiling of 128 and froze it into every savepoint. If that job later needs parallelism 200, it cannot get there with its state intact. Nobody chose 128; it was chosen for them. ## Why not just set it to 32768 everywhere Because it is not free. Flink maintains metadata for its ability to rescale state, and that metadata grows linearly with maximum parallelism. Very high values inflate checkpoint metadata and add per-checkpoint bookkeeping that a low-parallelism job gains nothing from. There is a second, subtler cost. Key groups are assigned to subtasks as contiguous ranges. If maximum parallelism is not an integer multiple of the running parallelism, some subtasks receive one more key group than others. With 128 groups over 5 subtasks, three subtasks get 26 and two get 25 — a 4% imbalance, harmless. With a maximum parallelism of 10 over 3 subtasks, the split is 4, 3 and 3: one subtask carries a third more key groups than the others, and if keys are evenly distributed it carries a third more state and a third more work. The rule of thumb is to pick a maximum parallelism that is comfortably larger than any parallelism you will run, and ideally a round number that your realistic parallelisms divide into cleanly — powers of two are popular for exactly this reason. ## How to actually choose it This is a capacity-planning judgment, not a formula: 1. **Bound the horizon.** What is the largest parallelism this job could plausibly need over its lifetime — peak traffic, plus a multiple for backfills and replays, which are usually where the highest parallelism is genuinely needed? 2. **Multiply for headroom**, because the decision is irreversible and the cost of being too high is small while the cost of being too low is a state-losing migration. 3. **Round to something divisible.** 128, 256, 512, 1024 are the values you see in practice. A job that will never exceed parallelism 32 does not need 32768. 4. **Set it explicitly**, in code, from the very first deployment — even when the default would have given you the same number. Explicit means the next engineer can see the decision instead of reverse-engineering it. 5. **Set stable operator uids at the same time.** Uid and maximum parallelism are the two things that must be right before the first savepoint exists, and both are painful to retrofit. ## What to do when it is already wrong The honest answer is that you migrate the state. The State Processor API reads a savepoint offline — `SavepointReader` in a batch job — lets you extract the keyed state as ordinary records, and `SavepointWriter` writes a fresh savepoint for a job configured with the new maximum parallelism. Flink re-hashes the keys into the new key-group layout as it writes. This is a real batch job with real cost proportional to state size, and it needs a maintenance window and a tested rollback, but it beats reprocessing history from the source. The alternative — replay from the source into a fresh job — is only viable when the source retains enough history and the job's output is idempotent or can be rebuilt. ## Where this sits relative to other scaling levers Maximum parallelism is not how you scale; it is the ceiling on how far you can scale. Day-to-day scaling is changing parallelism and restoring from a savepoint, where Flink redistributes key-group ranges for you. Note also that operator (non-keyed) state does not use key groups at all — it redistributes by even-split or union over list elements — so this ceiling constrains keyed state specifically. ## The one-sentence version for an interview Maximum parallelism is the number of key groups; key groups are the atomic unit of keyed-state redistribution; changing the count changes the key-to-group mapping and orphans the state; so set it once, explicitly, above your worst-case parallelism, and treat any later change as a savepoint-rewrite migration.
- What maximum parallelism does Flink pick if you never set one?It decides at first start from the operator's parallelism: `nextPowerOfTwo(parallelism + parallelism / 2)`, never below `128` and never above `2^15` — so 128 up to parallelism 85 and 256 from 86 to 171. The number is then frozen into every savepoint. A job launched at parallelism 4 therefore silently carries a ceiling of 128 — fine until someone needs parallelism 200, at which point the state cannot come along.
- Why can parallelism change freely when maximum parallelism cannot?Because parallelism only changes how key groups are *assigned*, not how many exist. Each subtask owns a contiguous range of key groups, so rescaling means splitting or merging ranges and moving the corresponding state files — the key-to-key-group function is untouched. Changing maximum parallelism changes that function itself, so the persisted groups no longer correspond to anything in the new layout.
- How can uneven key-group distribution show up as skew?Key groups are handed out as contiguous ranges, so when maximum parallelism is not a multiple of the running parallelism some subtasks get one more group than others. At 128 groups over 5 subtasks that is a 4% imbalance and irrelevant; at 10 groups over 3 subtasks the split is 4, 3 and 3, so one subtask carries a third more state and work. Choose a maximum parallelism your realistic parallelisms divide into cleanly.
- What is the migration path if maximum parallelism was set too low?Rewrite the savepoint with the State Processor API: read it in a batch job with `SavepointReader`, extract the keyed state as records, and write a fresh savepoint with `SavepointWriter` for a job configured with the new maximum parallelism, letting Flink re-hash keys into the new key-group layout. It costs a batch run proportional to state size plus a maintenance window, but it is far cheaper than replaying history from the source.
Key groups are like numbered post-office boxes: you can hire more or fewer clerks and reassign blocks of boxes between them, but renumbering the boxes themselves invalidates every address already printed on the mail.
saying these in an interview costs you the question
- Says maxParallelism can be raised on the next restart
- Confuses maximum parallelism with the running parallelism
- Thinks there is one key group per distinct key
- Believes operator state is also partitioned into key groups
- Assumes setting it to 32768 everywhere is free