In Flink, which parallelism setting wins — parallelism.default, env.setParallelism(), or an operator's own?
answer
- four places can say it, one of them wins
- more specific beats more general
- the CLI flag is not the last word
- a separate ceiling exists that you cannot raise later
- key groups are the reason for that ceiling
basics
~10 sThe most specific setting wins. An operator's setParallelism() overrides the environment's env.setParallelism(), which overrides the submission-time -p flag, which overrides the cluster-wide parallelism.default. Maximum parallelism is a separate, state-bound ceiling.
solid answer
~40 sFlink resolves parallelism from most specific to least. `setParallelism()` on an individual operator, source or sink wins for that operator. Failing that, the job-wide `env.setParallelism(n)` set in the program applies. Failing that, the value submitted at the client — `bin/flink run -p 10` — applies. Last comes the cluster-wide `parallelism.default` from the Flink configuration file. Separately, `setMaxParallelism()` bounds how far a keyed operator can ever be rescaled, because it fixes the number of key groups that keyed state is partitioned into. Its default is roughly `parallelism + parallelism/2`, with a floor of 128 and a ceiling of 32768, and it can be set per operator, per environment, or program-wide through the `pipeline.max-parallelism` option in the configuration file or a `-D` flag at submission. Changing it on an existing job makes previously saved state incompatible.
code
java · 7 linesfinal StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(8); // job-wide default
stream.map(new Enrich()).setParallelism(24) // CPU-heavy stage
.keyBy(Event::getUserId)
.sum("count") // POJO field, by name
.sinkTo(jdbcSink).setParallelism(4); // limited by the connection poolgo deeper
Be ready to name the places parallelism can be set — operator, environment, CLI -p, and parallelism.default — and to say the more specific one wins.
Explain the full resolution order and why maximum parallelism is a different setting: it fixes the key-group count that keyed state is partitioned into, defaulting to roughly parallelism plus half, floored at 128 and capped at 32768.
Show judgment about per-operator parallelism — sources capped by partition count, sinks capped by a connection pool — and the cost that every parallelism change breaks the operator chain and adds a serialized exchange.
Own the irreversible part: maximum parallelism is chosen once and bounds the job's scaling horizon for the life of its state, so it belongs in the design review alongside the key choice, not in an incident response.
## Parallelism, concretely A Flink job is a graph of tasks — transformations, sources, sinks. Each task is split into parallel instances, and the number of instances is that task's *parallelism*. Every instance processes a subset of the task's input. Parallelism is therefore the main throughput knob, and it is set independently per operator if you want. ## The four levels **Operator level.** Call `setParallelism()` directly on the returned stream, right after the transformation it refers to: ```java DataStream<Tuple2<String, Integer>> counts = text .flatMap(new LineSplitter()) .keyBy(value -> value.f0) .window(TumblingEventTimeWindows.of(Duration.ofSeconds(5))) .sum(1).setParallelism(5); ``` This is the most specific setting and it wins for that operator alone. **Execution-environment level.** `env.setParallelism(3)` sets the default for every operator, source and sink the environment executes. Any operator that sets its own value overrides this. **Client level.** The CLI accepts `-p`: `./bin/flink run -p 10 job.jar`. This supplies the job's parallelism at submission. **System level.** `parallelism.default` in the Flink configuration file gives a cluster-wide default for all execution environments. The resolution order runs from the most specific to the most general: operator, then environment, then client, then system. In practice this is why a job submitted with `-p 32` can still run part of its graph at parallelism 5 — an explicit `setParallelism(5)` in the code is more specific and stands. That specificity ordering also explains a common operational complaint: an operator team tries to scale a job by raising `-p` at submission, and nothing changes, because the program itself calls `env.setParallelism(...)`. Leaving parallelism out of the application code is the usual fix, for the same reason the docs prefer choosing the runtime mode at submission — a configuration-free artifact is one a platform can tune without a rebuild. ## Why some operators should differ from the job default Uniform parallelism is rarely optimal. A source is often capped by the external system: a Kafka source cannot usefully run more subtasks than the topic has partitions, and extra subtasks simply idle. A sink writing to a database may need *lower* parallelism than the compute stages so it does not exhaust the connection pool or trip a rate limit. A CPU-heavy enrichment step in the middle may want more than either. Setting these individually is what the operator-level knob is for. The other side of that coin: every parallelism change between adjacent operators forces a redistribution and breaks the operator chain, so records get serialized and shipped between subtasks. A chain of `map().filter().map()` at the same parallelism runs in a single thread with no serialization; inserting one differing `setParallelism()` in the middle turns that into two tasks and a network exchange. Vary parallelism where the bottleneck actually justifies it, not by reflex. ## Maximum parallelism is a different thing `setMaxParallelism()` looks like a sibling but plays a different role. Flink partitions keyed state into *key groups*; the number of key groups equals the operator's maximum parallelism, and each subtask owns a contiguous range of them. That indirection is what lets you restore a savepoint at a different parallelism — key-group ranges are reassigned, and each subtask reconstructs its state from the groups it now owns. So max parallelism is the hard ceiling on how far a keyed operator can ever be scaled out for the life of its state. It can be set on an operator or on the environment with `setMaxParallelism()`, and in Flink 2.3 also program-wide without touching code, through the `pipeline.max-parallelism` option — in `config.yaml` or as `-Dpipeline.max-parallelism=...` at submission. There is no dedicated CLI flag like `-p`, and the parallel-execution docs page still says client and system level are excluded; the option is what the runtime actually reads. Its default is roughly `operatorParallelism + (operatorParallelism / 2)`, with a lower bound of 128 and an upper bound of 32768. Two warnings come with it. Setting it very large is not free: some state backends keep internal structures that scale with the key-group count, so an unnecessarily huge value costs performance. And changing it explicitly when recovering a job from its original state leads to state incompatibility — the key-group mapping no longer matches what the savepoint recorded. That is why a long-lived job's max parallelism is a decision made deliberately at design time, not adjusted later during an incident. ## Reading the result The web UI shows each job vertex with its parallelism, which is the fastest way to confirm what the resolution order actually produced — more reliable than reading the submission command, because the code may have overridden it. Naming operators with `.name("...")` makes that view legible. If one vertex shows high backpressure while its neighbours are idle, the parallelism of that vertex, or the key distribution feeding it, is the first thing to check.
- A job is submitted with -p 64 but one vertex still shows parallelism 4 in the web UI. Why?Something more specific overrode the client flag — almost always an explicit `setParallelism(4)` on that operator in the code, or a job-wide `env.setParallelism(4)`. The resolution order is operator, then environment, then client, then `parallelism.default`, so any in-code setting beats `-p`. Removing the in-code call restores the platform's ability to tune parallelism at submission.
- Why can't you raise a keyed operator's parallelism above its maximum parallelism after the job has state?Maximum parallelism fixes the number of key groups keyed state is partitioned into, and each subtask owns a range of them. You cannot have more subtasks than key groups, and changing the key-group count changes how every key maps to a group, so a savepoint taken under the old value no longer matches. The ceiling is set once, before the state exists.
- Why not just set maximum parallelism to 32768 on every job to keep options open?Some state backends maintain internal structures that scale with the key-group count, so a needlessly large maximum costs memory and performance for a scaling headroom you will never use. Pick a value that covers a realistic scaling horizon for the job, then leave it alone, since changing it later invalidates existing state.
saying these in an interview costs you the question
- Says the -p flag always overrides in-code parallelism
- Confuses maximum parallelism with the number of task slots
- Thinks max parallelism can be raised later on a running job's state
- Believes every operator must run at the same parallelism
- Assumes more source subtasks than Kafka partitions adds throughput