Traffic has grown and a Flink job no longer keeps up. The sources read more slowly than data arrives, the Back Pressure tab in the Flink web UI shows HIGH on upstream tasks, and one operator’s subtasks report busy time close to 100 percent. The team wants to give the job more parallelism and needs to know whether that will work and how to do it without losing state.
Parallelism is the number of parallel instances a task is split into, each processing a subset of the task’s input, as the Flink parallel execution documentation defines it. Rescaling means changing that number for a running job, which for a stateful job means redistributing its state across the new instances.
More parallelism helps when the busy operator’s subtasks are all close to 1000 ms busy and the load is spread evenly, within the limits set by max parallelism, source partitions and free task slots. The standard way to rescale a streaming job is to stop it with a savepoint and restart it from that savepoint with a new parallelism.
This page covers the decision to scale and the procedure. Finding which operator is the bottleneck and why is covered in the guide to the bottleneck behind Flink backpressure. If the lagging consumer is a Kafka consumer group and not a Flink job, the guide to Kafka consumer lag is the better starting point. For how scaling fits into running Flink in production more widely, see the complete Flink guide.
Check that more parallelism will help
Busy time and backpressure are the per-job signals that show a job falling behind. The back pressure monitoring documentation says every subtask reports three metrics that add up to about 1000 ms per second.
busyTimeMsPerSecond
backPressuredTimeMsPerSecond
idleTimeMsPerSecond
More parallelism helps when the busy operator’s subtasks are all close to 1000 ms busy and the work is spread evenly across them. Each new subtask then takes a share of the same kind of work.
It does not help for a hot key, for a sink or external call limited by another system, or for a source already at its Kafka partition count, as the table below shows. The Flink Kubernetes Operator autoscaler documentation says a vertex limited by an external system gains coordination overhead and no throughput from additional subtasks. The backpressure guide shows how to tell these cases apart with per-subtask metrics and flame graphs.
Know the limits before you rescale
Three limits decide how far a job can scale: the max parallelism of its stateful operators, the partition count of its sources and the task slots in the cluster.
Max parallelism fixes the number of key groups
Flink partitions keyed state into key groups, and max parallelism sets how many there are. The parallel execution documentation says it is the upper bound on parallelism when a job is restored from a savepoint with a new parallelism. The production readiness checklist adds that there is currently no way to change an operator’s max parallelism after the job has started without discarding that operator’s state.
pipeline.max-parallelism
setMaxParallelism(int maxParallelism)
A job that never set it has the default Flink chose when it first started, which the savepoint upgrade guide explains along with how to read the stored value from a savepoint. Read that value before planning a rescale. If the target parallelism is above it, the job cannot be scaled that far without rewriting its state.
Pick the max parallelism for new jobs with rescaling in mind. The autoscaler documentation recommends a number with many divisors, such as 120, 180, 240, 360 or 720, because key groups only spread evenly across a parallelism that divides them. The checklist warns against going too high, since the metadata Flink keeps for rescaling grows linearly with max parallelism.
Kafka partitions cap source parallelism
In Flink’s Kafka source, each split is one Kafka partition. The Kafka connector documentation says the source’s split enumerator assigns partitions to readers itself, distributing them across subtasks round-robin, and that the source does not rely on committed offsets for fault tolerance. Partition ownership is Flink’s decision, not a consumer group rebalance, so a source subtask beyond the partition count gets nothing to read.
The same documentation warns that the Kafka source does not go idle on its own when parallelism is higher than the number of partitions. Those empty subtasks can hold back watermarks downstream. The fix is to lower the source parallelism or add an idle timeout to the watermark strategy.
WatermarkStrategy.withIdleness(...)
Adding partitions to the topic is the way to raise the source’s ceiling. Partition discovery picks up new partitions without a restart, every 5 minutes by default.
partition.discovery.interval.ms
Operators downstream of the source are not bound by the partition count. A job can keep its source at the partition count and raise the parallelism of the busy operator after it.
Task slots and TaskManager sizing
A Flink cluster runs a JobManager and one or more TaskManagers, and the task slot is the unit of resource scheduling, as the Flink architecture documentation describes. With slot sharing, which is the default, a cluster needs exactly as many task slots as the highest parallelism used in the job. Raising one operator from 8 to 16 therefore needs 16 slots for the job, even if every other operator stays at 8.
taskmanager.numberOfTaskSlots
The number of slots per TaskManager is a sizing choice. The architecture documentation says a TaskManager with three slots dedicates a third of its managed memory to each slot, and that slots separate managed memory only, with no CPU isolation. Adding slots to existing TaskManagers splits the same memory and CPU more ways, so a CPU-bound operator usually needs more TaskManagers, not more slots per TaskManager.
In a standalone cluster the ResourceManager can only hand out slots from TaskManagers that are already running and cannot start new ones, so the TaskManagers have to be added before the rescale.
Rescale with a savepoint
For a streaming job, the elastic scaling documentation describes the standard way to rescale as stopping the job with a savepoint and restarting it with a different parallelism. The savepoints documentation confirms a program can be restored from a savepoint with a new parallelism.
1. Check uids, max parallelism and free slots
Every stateful operator needs a stable uid, which the savepoints documentation strongly recommends setting by hand so state maps back to the right operator after the restart. Confirm the target parallelism is at or below the stored max parallelism, and that the cluster has enough free slots for the new highest parallelism. The savepoint upgrade guide covers the uid and max parallelism checks in detail.
2. Stop the job with a savepoint
The Flink CLI documentation describes stop as a graceful stop that takes a savepoint and then finishes the job.
./bin/flink stop --savepointPath <savepoint-dir> <job-id>
Leave out the drain flag, shown below. The CLI documentation says draining emits a final watermark and should only be used to end a job permanently, because resuming a drained job can give incorrect results.
--drain
3. Restart with the new parallelism
Submit the same job from the savepoint path printed by the stop command. The client parallelism option sets the default for every operator that has not pinned its own parallelism in code, so this command changes all of those operators together.
./bin/flink run --fromSavepoint <savepoint-path> -p <new-parallelism> <job.jar>
The parallel execution documentation says an operator’s parallelism is set by calling a method on that operator, and that this overwrites the execution environment’s default. To change one operator and leave the rest alone, change that operator’s own setting in code, rebuild the JAR and submit with the client option left at the value the job already runs with.
operator.setParallelism(<new-parallelism>)
Changing one operator at a time where possible makes the effect of each change visible.
4. Confirm the job caught up
Watch the same metric that showed the problem. The busy operator’s busyTimeMsPerSecond should fall below 1000 on every subtask, and source lag should shrink. If one subtask stays busy while the new ones idle, the load is skewed and more parallelism will not fix it.
Rescale automatically with Reactive Mode
The elastic scaling documentation describes Reactive Mode, an experimental mode of the adaptive scheduler that makes a job use every slot in the cluster and rescales it from the latest completed checkpoint when TaskManagers are added or removed, with no savepoint step. Its limitations section says the supported deployments are standalone in application mode, Docker in application mode and standalone Kubernetes application clusters, with one job per application. Active resource providers such as native Kubernetes and YARN are explicitly not supported.
scheduler-mode: reactive
The documentation recommends periodic checkpointing for stateful jobs in Reactive Mode and warns that without a restart strategy, Reactive Mode fails the job instead of scaling it.
Rescale on Kubernetes with the operator autoscaler
Teams running Flink with the Flink Kubernetes Operator can let its built-in autoscaler set parallelism per job vertex, a task made of one operator or a chain of operators. The autoscaler documentation says it measures how busy each vertex is and adjusts parallelism up or down. It scales streaming jobs, and scaling a source requires the FLIP-27 Source API, the source interface that exposes the busy-time metric, with Kafka sources the most complete.
job.autoscaler.enabled: "true"
A safe first step is the metrics-only mode, which evaluates everything and reports the parallelism the autoscaler would apply without scaling the job.
job.autoscaler.scaling.enabled: "false"
How a decision is applied depends on the scheduler. With the adaptive scheduler, the autoscaler documentation says decisions are applied in place without restarting the job, although Flink’s elastic scaling documentation says scaling events under the adaptive scheduler trigger job and task restarts. Without the adaptive scheduler, every decision is rolled out as a full upgrade of the job, carrying state across the restart according to the job’s upgrade mode. The operator’s job management documentation lists the supported values as stateless, last-state and savepoint.
spec.job.upgradeMode
jobmanager.scheduler: adaptive
The autoscaler respects the same limits as a manual rescale. Its documentation says the max parallelism fixes the number of key groups and is a hard ceiling no scaling can exceed.
pipeline.max-parallelism
The autoscaler’s own upper bound defaults to 200 and never exceeds the max parallelism configured for the vertex.
job.autoscaler.vertex.max-parallelism
It also aligns each parallelism to the key groups of a keyBy vertex or the partitions of a source, so the data spreads evenly. After each scale-up the autoscaler compares the expected gain with the gain that materialized. The autoscaler documentation says a scale-up that delivers less than a tenth of the expected gain is marked ineffective and reported with an IneffectiveScaling event, which is the sign of a vertex limited by an external system. The operator configuration reference describes the detection switch below as enabling detection of ineffective scaling and allowing the autoscaler to block further scale-ups, and lists it as false by default. The threshold and the switch, with their defaults:
job.autoscaler.scaling.effectiveness.threshold: "0.1"
job.autoscaler.scaling.effectiveness.detection.enabled: "false"
When the rescaled job does not come back
A rescale fails in one of three places: the restore, the scheduling of the new subtasks, or the throughput after the job is running again.
The restore fails
A restore that fails with a max parallelism exception, or one that cannot map state to an operator, is a savepoint problem. The savepoint upgrade guide matches each exception to its cause and fix, including how to read the stored max parallelism from the savepoint.
The job waits for slots
If the job restores but its tasks stay in scheduling, the cluster does not have enough free slots for the new highest parallelism. Under the default scheduler, slot requests time out after the period set in the configuration reference, 5 minutes by default.
slot.request.timeout
Add TaskManagers, or lower the parallelism to fit the free slots. With the adaptive scheduler, the elastic scaling documentation describes a resource wait timeout that controls how long the job waits for enough TaskManagers before it stops.
jobmanager.adaptive-scheduler.resource-wait-timeout
The job runs but still falls behind
Check the busy operator’s subtasks again. One subtask still busy while the new ones idle means skew. Source subtasks with no partitions mean the source parallelism is above the partition count. Busy subtasks with no gain mean an external system is the limit. Each of these points back to the backpressure guide.
A larger job can also take longer to checkpoint. If checkpoints start to slow or time out after the rescale, the guide to Flink checkpoints that fail or time out covers how to find the slow stage. If the state itself keeps growing as traffic grows, see Flink state that keeps growing.
Where Flex fits
This page is published by Factor House, which makes Flex, its Flink job management and monitoring product, so weigh this section with that in mind. Flex is one of several ways to watch and restart a Flink job, and the comparison of tools for self-managed Apache Flink sets it beside the Flink web UI, the REST API and the Kubernetes operator.
The Flink REST API reference says the monitoring API is used by Flink’s own dashboard and is designed to be used by custom monitoring tools as well. Flex reads the job’s parallelism and the cluster’s slots through that API, so the numbers it shows are the ones Flink reports.
What Flex shows. The Flex jobs documentation lists the Parallelism set for each job and a Stoppable flag that says whether the job can be stopped gracefully via a savepoint. The Flex task managers documentation shows total slots, free slots, total CPU and total heap for the cluster, free slots over time, and total and free slots for each TaskManager. Free slots is the number to check before a rescale.
What Flex does. The Flex jobs documentation says the Inspect view has quick actions to Stop, Cancel, take a Savepoint and trigger a Checkpoint, and that submitting a JAR sets its Parallelism and a Savepoint Path, with an option to allow non-restored state. Together these cover the stop, change and restore sequence above.
What stays in Flink. The Flex documentation describes no action that changes the parallelism of a running job and no autoscaling, so a rescale in Flex is a stop with savepoint followed by a new submission. It does not describe showing max parallelism or key groups, so the stored max parallelism comes from the savepoint, as the savepoint upgrade guide shows. Busy time per subtask comes from Flink’s web UI or metric reporters.
The Flink REST endpoint also carries the actions a rescale depends on, stopping a job with a savepoint and submitting it again. The Flink SSL setup documentation says the REST endpoint accepts connections from any client by default and does not authenticate the client, and recommends an authenticating proxy in front of it. Anyone who can reach the REST port can therefore stop and resubmit a job at a new parallelism, unless a proxy sits in between. Flex has its own user authentication and role-based access control, and its role-based access control documentation describes roles being granted ALLOW, DENY or STAGE on actions for specific resources.
To check each job’s parallelism and the cluster’s free slots before the next rescale, Flex reads them from the Flink REST endpoint.
FAQ
Can a Flink job be rescaled without a restart?
Not without any restart. The elastic scaling documentation says scaling events under the adaptive scheduler trigger job and task restarts, and that Reactive Mode restarts the job from the latest completed checkpoint. What the adaptive scheduler removes is the manual savepoint and resubmission. From Flink 1.18 the same documentation describes re-declaring a running job’s resource requirements through a REST endpoint, and the Kubernetes operator autoscaler applies decisions in place when the adaptive scheduler is enabled.
PUT /jobs/<job-id>/resource-requirements
What parallelism should a Flink Kafka source have?
No more than the number of partitions it reads. The Kafka connector documentation says each split is one partition, so extra source subtasks have nothing to read and can hold back watermarks unless an idle timeout is set.
How many task slots does a rescaled Flink job need?
With default slot sharing, as many as the highest parallelism of any operator in the job. The Flink architecture documentation states this rule, so check free slots against the new highest parallelism before restarting.
Can max parallelism be raised as part of a rescale?
Not by configuration alone. The production readiness checklist says max parallelism cannot change after a job has started without discarding the affected operator’s state, so raising it means rewriting the savepoint, which the savepoint upgrade guide covers.