Skip to content

How to find the bottleneck behind Flink backpressure

Flink
Chad Harris·October 3, 2026·11 min read

A backpressured Flink job usually shows up as three symptoms at once. The sources read more slowly than data arrives, end to end latency climbs, and checkpoints take longer until some expire. In the Flink web UI the job graph shows one or more tasks in black and a HIGH status on the Back Pressure tab.

Backpressure is what happens when an operator cannot process records as fast as they arrive, so its input buffers fill and the slowdown spreads upstream until it reaches the sources.

The back pressure monitoring documentation explains that a warning on a task means it produces data faster than the operators downstream can consume, so the cause sits further down the job graph and not on the task with the HIGH warning.

This page is about backpressure inside a Flink job. If the symptom is a Kafka consumer group falling behind and the consumer is not a Flink job, the guide to Kafka consumer lag is the better starting point. For how backpressure fits into running Flink in production more widely, see the complete Flink guide.

Confirm the job is backpressured

Every subtask splits each second into time spent backpressured, idle and busy. The three add up to about 1000 ms, and the back pressure documentation says they are averaged over the last couple of seconds.

backPressuredTimeMsPerSecond
idleTimeMsPerSecond
busyTimeMsPerSecond

The same documentation defines the status on the Back Pressure tab from the share of time a subtask spends backpressured.

OK:   0% to 10%
LOW:  above 10% up to 50%
HIGH: above 50% up to 100%

A subtask close to 1000 ms idle is starved of input, not slow. If every task is mostly idle, the job is not backpressured and the slowdown is somewhere else.

Two subtasks with different load patterns can report the same busy time of about 500 ms, one under constant 50% load and one alternating between fully busy and fully idle. A window operator that fires periodically can also move the bottleneck between two tasks every few seconds, as the Flink post on how to identify the source of backpressure shows with a sliding window example.

When a backpressured Flink job reads from Kafka, lag shown by Kafka tooling can move in steps. The Flink Kafka source commits its consuming offset only when a checkpoint completes, and it does not rely on committed offsets for fault tolerance. Tools that read committed offsets therefore see a backpressured job one checkpoint late, and see no progress at all while checkpoints fail. For lag on consumers that are not Flink jobs, the guide to Kafka consumer lag applies. The source’s own metric is the live signal.

pendingRecords

Find the busy task downstream

Backpressure and checkpoint duration are the two job-level signals to watch for a stuck job. The back pressure monitoring documentation defines the per-subtask metrics above, and the guide to checkpointing under backpressure says checkpointing durations become very high due to backpressure.

The web UI colours each task in the job graph. Idle tasks are blue, fully backpressured tasks are black and fully busy tasks are red, with shades in between. The Flink post on identifying the source of backpressure says the busiest (red) task downstream of the backpressured tasks will most likely be the bottleneck. In the rare case where the network exchange itself is the bottleneck, the same post notes that the downstream task has empty input buffers while the upstream output buffers are full.

Each task in the job graph shows the highest value across its subtasks, so one busy subtask colours the whole task. Click the task and open its BackPressure tab to see busy, backpressured and idle time for every subtask. That per-subtask view is where skew becomes visible.

To read the same numbers outside the UI, the Flink REST API exposes one endpoint per vertex. Its reference says the call may start back pressure sampling if necessary, so reading a whole job means one request for each vertex.

GET /jobs/:jobid/vertices/:vertexid/backpressure

For alerting, Flink’s metrics reference lists three per-task metrics that a metric reporter can send to an external system.

backPressuredTimeMsPerSecond
busyTimeMsPerSecond
isBackPressured
F1 The bottleneck task and the backpressured tasks upstream of it
Bottleneck task Backpressured task
Colour in the web UI job graph Red (busy) Black (waiting for output buffers)
Metric that proves it busyTimeMsPerSecond near 1000 on a subtask backPressuredTimeMsPerSecond high, busy time low
Where it sits Downstream of the backpressured tasks Upstream, often back to the sources
What to change Its parallelism, calls, keys or sink Nothing, it recovers with the bottleneck
Common mistake Missed because it shows no warning Scaled, which adds unneeded resources
Drawn from the Apache Flink documentation on monitoring back pressure and the Flink post How to identify the source of backpressure (flink.apache.org, 2021-07-07), both linked in this section. Back pressure monitoring documentation

Tell the four causes apart

Once the busy task is found, the next question is why it is busy. Four causes cover most cases, and each has a different fix.

A slow sink

The bottleneck is the last task, the one that writes to a database, a file system or another topic. Its subtasks are all busy at similar levels, and upstream tasks are black all the way to the sources.

The proof is the busy time of the sink subtasks together with the health of the system they write to. If the target system shows high write latency or throttling while the sink subtasks are busy, the sink is the bottleneck.

busyTimeMsPerSecond

Skewed keys

One or a few subtasks of the busy task sit near 1000 ms busy while the others are idle or close to it. The per-subtask view on the BackPressure tab shows it, and so does the per-subtask record count.

numRecordsInPerSecond

The Flink DataStream operators documentation says all records with the same key are assigned to the same partition. A single hot key therefore lands on one subtask no matter how much parallelism the operator has. The Flink SQL performance tuning guide names data skew as very common in production and as something that makes jobs easy to put under backpressure.

A slow external call

The busy task calls a database, a cache or an HTTP service for each record, often for enrichment. Its subtasks are busy, but CPU on the TaskManager is low, because the thread is waiting rather than working.

The flame graph documentation describes an Off-CPU view that visualises blocking calls found in the stack samples, and a wide block on a client’s request method in the Off-CPU view of the busy task points to the external call. Flame graphs are opt-in, because sampling affects the job being measured.

rest.flamegraph.enabled: true

Too little parallelism

Every subtask of the busy task is busy at a similar level, the work per record is CPU bound, and the On-CPU flame graph shows the operator’s own code. Input rate has grown past what the current number of subtasks can process. Busy time is spread evenly, which separates this case from skew, and CPU is high, which separates it from a slow external call.

F2 The four causes of a busy task and the proof for each
What shows first What proves it
Slow sink Last task busy, all upstream black High write latency or throttling in the target system
Skewed keys One subtask near 1000 ms busy, others idle Per-subtask view and numRecordsInPerSecond
Slow external call Busy subtasks, low TaskManager CPU Wide client block in the Off-CPU flame graph
Too little parallelism All subtasks evenly busy, CPU high On-CPU flame graph shows the operator's own code
Drawn from the Apache Flink documentation on monitoring back pressure and on flame graphs, and the SQL performance tuning guide on data skew. Flame graph documentation

Fix one cause at a time

Change one thing, then watch busy time on the task that was the bottleneck. If the fix worked, that task’s busy time drops and the black tasks upstream turn blue or show a lower backpressure share.

Relieve a slow sink

First check whether the target system can take more concurrent writes. If it can, raise the sink operator’s parallelism on its own. Connectors expose their own option, and the Kafka SQL connector documents one that otherwise defaults to the parallelism of the upstream chained operator.

sink.parallelism

If the target system is already the limit, more sink parallelism adds load to a system that is slow, so the target’s write capacity is what has to grow.

Add parallelism to the bottleneck operator

For a CPU-bound operator with evenly spread load, more parallelism on that operator is the direct fix. The Flink post on identifying the source of backpressure frames the choice as either adding resources or using existing resources better, by optimising the code, tuning configuration or avoiding data skew.

Stateful operators cannot scale past their maximum parallelism. The production readiness checklist says there is no way to change an operator’s maximum parallelism after the job has started without discarding that operator’s state, and that Flink picks 128 when none is set and parallelism is 128 or lower. Set it explicitly before the job’s first production run, using the key listed in the Flink configuration reference.

pipeline.max-parallelism

To rescale a running streaming job, the elastic scaling documentation describes stopping it with a savepoint and restarting it with a different parallelism. It also describes Reactive Mode and, from Flink 1.18, a REST endpoint that re-declares a running job’s resource requirements, which the documentation labels a minimum viable product feature. The guide to rescaling a Flink job that cannot keep up covers the decision to scale and the full procedure.

PUT /jobs/<job-id>/resource-requirements

Make external calls asynchronous

For a slow external call, more parallelism is the expensive fix. The async I/O documentation explains that very high parallelism means more tasks, threads, network connections to the database and buffers, while asynchronous requests let one parallel instance handle many requests at once.

Async I/O has three settings to choose deliberately.

  • Timeout: how long a request may take before it fails. By default a timed-out request throws an exception and restarts the job, which can turn into a restart loop.
  • Capacity: how many requests each subtask may have in flight. Once it is reached the operator triggers backpressure itself.
  • Retry strategy: when a failed request is tried again.
AsyncDataStream.unorderedWait(stream, new EnrichFunction(), 1000, TimeUnit.MILLISECONDS, 100)

The same documentation warns that the client must really be asynchronous, because a client whose query method blocks, or code that waits on the returned future inside asyncInvoke, voids the asynchronous behaviour. Unordered mode emits results as soon as each request finishes and, the documentation notes, has the lowest latency and overhead with processing time. Ordered mode keeps the stream order.

Change the key, not only the parallelism

For skew, parallelism alone does not help, because a hot key stays on one subtask. In Flink SQL and the Table API, the SQL performance tuning guide describes local-global aggregation, which aggregates on each upstream subtask before the keyed shuffle and depends on mini-batch being enabled.

table.exec.mini-batch.enabled: true
table.exec.mini-batch.allow-latency: 5 s
table.optimizer.agg-phase-strategy: TWO_PHASE

For a COUNT(DISTINCT ...) aggregation that is skewed, the same guide describes splitting the distinct aggregation into two levels.

table.optimizer.distinct-agg.split.enabled: true

In the DataStream API the key is the job’s own design. A finer key, or a first aggregation on a composite key followed by a second aggregation on the original key, spreads a hot key over more subtasks. Either change alters state layout, so it needs a new job graph or a savepoint migration rather than a configuration change. Growing keyed state from a poorly chosen key is covered in Flink state that keeps growing.

Use buffer debloating for checkpoint relief, not throughput

Buffer debloating shortens checkpoint barrier time under backpressure and leaves the bottleneck in place. The network memory tuning guide describes it as automatically adjusting the amount of in-flight data so that it can be consumed in a configured target time. Less in-flight data means checkpoint barriers reach the end of the job sooner under backpressure.

taskmanager.network.memory.buffer-debloat.enabled: true

The guide lists its limits. It caps the buffer size in use but does not reduce memory use, it can leave too much data on the slow input of a task with several inputs, and above a parallelism of about 200 it may need more floating buffers than the default.

taskmanager.network.memory.floating-buffers-per-gate

The guide to checkpointing under backpressure lists three responses in order: remove the backpressure source, reduce the in-flight data, and use unaligned checkpoints. When checkpoints are already failing, Flink checkpoints that fail or time out covers the checkpoint side.

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 inspect a backpressured job, and the comparison of tools for self-managed Apache Flink sets it beside the Flink web UI, the REST API and metrics stacks.

Flink has no separate client protocol for management tools. The 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, so Flex and the web UI read the same numbers. Flex connects through the endpoint set in its Flink cluster configuration.

That shared REST endpoint has an access question attached. 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. Giving an application team the web UI to read backpressure therefore also gives it the endpoints that stop and cancel jobs, unless a proxy or a tool with its own access control sits in between.

Diagnosis. Flex documents these views in its job inspection documentation.

  • The Topology tab shows the job’s dataflow graph, which the documentation describes as identical to the Flink Job Graph.
  • For each task selected in the graph, a Backpressures view shows the Status (OK, LOW or HIGH) and a backpressure level percentage.
  • Subtask metrics list values for each parallel instance, and aggregated metrics give Min, Max, Avg, Sum and percentiles of read-bytes and read-records across subtasks. A Max far above the Avg on the bottleneck task is the skew signal.
  • The Overview tab and the job Details view chart bytes and records read and written per second over the past hour, which shows when the slowdown started.

What stays in Flink. The job documentation does not describe the following, so they stay in the Apache Flink web UI, REST API or metric reporters.

  • Busy and idle time per subtask, and the red and black colouring of the job graph.
  • Flame graphs for the bottleneck operator.
  • Alerting on backpressure. The Flex metrics glossary lists job, cluster, JobManager and TaskManager metrics with no backpressure or busy time metric, so an alert on the backpressure time metric listed above comes from Flink’s own reporters.

Fixes. Flex documents these actions.

  • The Inspect view has quick actions to Stop, Cancel, take a Savepoint and trigger a Checkpoint.
  • Submitting a job sets its Parallelism and a Savepoint Path, which covers the stop with savepoint and restart that rescaling a job requires.

Flex ships as a single Docker container or JAR file, according to its GitHub repository, and connects through the Flink REST endpoint. To see backpressure status and per-subtask metrics for each job in Flex, point it at that endpoint.

FAQ

Which task is the bottleneck when Flink shows backpressure?

Usually the busiest task downstream of the backpressured ones, not the task with the HIGH warning. A backpressured task is waiting for output buffers, while the bottleneck is busy processing. In the web UI job graph that is the red task below the black ones.

Does buffer debloating fix backpressure?

No. It reduces the amount of in-flight data between tasks, which shortens checkpoint barrier propagation under backpressure, but the bottleneck operator is just as slow. The Flink documentation lists removing the backpressure source as the first response.

Is some backpressure normal in Flink?

Yes. The Flink post on identifying the source of backpressure notes that a lack of backpressure means the cluster is at least slightly under-utilised, and that minimising idle resources usually means accepting some. It becomes a problem when latency or checkpoint duration grows past what the job can tolerate.

Why does adding parallelism not remove the backpressure?

Either the load is skewed, so a hot key stays on one subtask, or the bottleneck is an external system that cannot take more concurrent requests. Check per-subtask busy time for skew and the Off-CPU flame graph for a blocking call before scaling again.

Related reading