skip to content

In a Flink job, how do you find which operator is causing backpressure?

level: seniorimportance: should knowfreq 58%

answer

  1. the symptom travels upstream
  2. victims above, starved below
  3. find the first one that is busy and not blocked
  4. one hot subtask means one hot key
  5. busy is not the same as CPU-bound

basics

~10 s

Walk the Flink job graph downstream and find the first task that is busy but not backpressured — that is the bottleneck. Everything upstream of it shows backpressure, everything downstream shows idle time.

solid answer

~40 s

Flink exposes per-subtask `busyTimeMsPerSecond`, `backPressuredTimeMsPerSecond` and `idleTimeMsPerSecond`, which the Web UI renders on the job graph and the Backpressure tab. Backpressure propagates upstream through credit-based flow control, so every task above the real bottleneck looks backpressured and every task below it looks idle. The diagnosis rule is: **the bottleneck is the first task, going downstream, that is busy and not backpressured**. Then ask *why* it is busy by opening the per-subtask breakdown. If one subtask is at 100% and its siblings are idle, you have key skew from `keyBy`. If all subtasks are busy, the operator is genuinely under-parallelized, its user function is expensive, it is blocking on an external call, or it is thrashing state on disk. Cross-check with record rates, checkpoint duration, GC logs and Kafka consumer lag before changing anything.

code

text · 7 lines
text
Task                       Busy   BackPressured   Idle
Source: kafka (4/4)         12%            88%      0%
map -> filter (4/4)         20%            80%      0%
KeyedProcess (1/4)         100%             0%      0%
KeyedProcess (2/4)          15%             0%     85%
KeyedProcess (3/4)          14%             0%     86%
Sink: jdbc (4/4)             5%             0%     95%

go deeper

for a junior

Know that backpressure means a downstream operator cannot keep up and that Flink's Web UI shows it per operator. Being able to point at the tab and say the pipeline is saturated is enough.

for a middle

Explain credit-based flow control and why backpressure propagates upstream, and use the busy/backpressured/idle metrics to name the bottleneck rather than the loudest victim.

for a senior

Demonstrate the full diagnosis: identify the operator, split skew from under-parallelization from blocking I/O from state or GC using the per-subtask breakdown and corroborating metrics, and pick a fix that matches the cause.

for a principal

Own what the platform measures by default. Decide which signals — lag, busy time, checkpoint duration — page a team, and set the expectation that pipelines are sized against a measured per-subtask rate rather than tuned reactively during incidents.

## What backpressure is in Flink Flink's network stack uses **credit-based flow control**: a downstream subtask announces to its upstream how many network buffers (credits) it has free, and the upstream only sends what it has credit for. When a downstream operator cannot keep up, its buffers fill, credits go to zero, and the upstream subtask blocks trying to hand off a record. That blocking is *backpressure*, and it walks up the chain until it reaches the sources, which then stop polling — so Kafka consumer lag grows and end-to-end latency rises while throughput settles at the slowest stage's rate. Backpressure is therefore a **symptom that travels**. The operator showing backpressure is almost never the problem; it is a victim of something further downstream. ## The three metrics Every task reports, per subtask, how it spent each second: - `busyTimeMsPerSecond` — doing work in user code or the runtime. - `backPressuredTimeMsPerSecond` — blocked trying to emit downstream. - `idleTimeMsPerSecond` — waiting for input. They sum to roughly 1000 ms per second, and the Web UI colours each job-graph vertex from them. The signature of a healthy pipeline is nearly everything idle. The signature of a saturated pipeline is a chain of backpressured tasks ending at one busy one. ## The diagnosis rule Start at the source and walk downstream. Skip every task that is backpressured — it is waiting on someone else. **The first task that is busy and not backpressured is the bottleneck.** Downstream of it, tasks will be idle because they are starved. A log of a real job graph makes this obvious: ``` Source: kafka (4/4) busy 12% backpressured 88% idle 0% map -> filter (4/4) busy 20% backpressured 80% idle 0% KeyedProcess (1/4) busy 100% backpressured 0% idle 0% KeyedProcess (2..4/4) busy 15% backpressured 0% idle 85% Sink: jdbc (4/4) busy 5% backpressured 0% idle 95% ``` The `KeyedProcess` operator is the bottleneck, and because only subtask 1 of 4 is busy, the cause is **key skew** — one hot key routed by `keyBy` to a single subtask. Raising parallelism will not help; every extra subtask still sends that key to one place. ## Why is the bottleneck busy? Once the operator is identified, open its subtask breakdown and separate the cases: - **One busy subtask, siblings idle** → skew. Fixes are semantic: pre-aggregate with a salted key and then combine, split the hot key into sub-keys, or move a global aggregate into a two-phase local/global shape. - **All subtasks busy at similar rates** → the stage is genuinely under-provisioned or the work per record is expensive. Raise that operator's parallelism (up to the maximum parallelism, and never above what the input partitioning allows for a source), or make the function cheaper. - **Busy but low CPU on the host** → the subtask is blocking, not computing. Synchronous lookups against a database or REST service are the usual cause; the fix is Flink's Async I/O operator, which keeps many requests in flight per subtask, or a cached/lookup-join approach. - **Busy with heavy disk I/O** → state access. A keyed operator on the embedded RocksDB backend pays a disk round trip per state access; check that the state directory is on local SSD, look at the RocksDB metrics, and consider whether the state per key is larger than it needs to be. - **Periodic stalls across all subtasks of a TaskManager** → garbage collection. Check GC logs and heap sizing; this shows up as busy time that does not correlate with record rate. ## What to look at alongside `numRecordsInPerSecond` and `numRecordsOutPerSecond` tell you what the pipeline is actually achieving and where the rate drops. Checkpoint duration is a second signal: under backpressure, barriers travel slowly and checkpoints take longer or start timing out — a job whose checkpoint duration climbed at the same moment its lag started growing is telling you the same story twice. Kafka consumer lag confirms that the source has been throttled rather than that the input dried up. And the source's own metrics matter: if the source is the busy-not-backpressured task, the ceiling may simply be the number of input partitions, since a Kafka source cannot usefully run at higher parallelism than there are partitions. ## Sources of confusion "Busy" is not CPU utilization — a subtask blocked in a synchronous HTTP call counts as busy while using almost no CPU. A 100% backpressured source with everything downstream idle is not possible; if the whole graph is idle, the input is the limit, not the job. And relieving backpressure by increasing buffer sizes only hides it: more in-flight data means the same throughput with worse latency and slower checkpoints, which is why buffer debloating exists to shrink in-flight data rather than grow it.

  • Every task in the job graph shows high idle time and consumer lag is flat. Is the job backpressured?
    No. An all-idle graph means the pipeline is waiting for input, so the limit is upstream of Flink — fewer records are arriving than the job can process, or the source cannot read faster because of partition count, fetch settings or broker-side throttling. Adding parallelism or tuning operators changes nothing here; investigate the source and the input topic instead.
  • How does backpressure show up in checkpoint metrics?
    Checkpoint duration climbs and checkpoints may start timing out. Barriers flow with the records, so when in-flight data is queued up behind a bottleneck the barrier takes far longer to reach downstream operators. A job whose checkpoint duration rose at the same time as its lag is giving you a second, independent signal that the pipeline is saturated rather than merely slow to snapshot.
  • Would raising the network buffer sizes relieve backpressure?
    No — it moves the queue, not the bottleneck. More in-flight buffers absorb bursts briefly, then settle at the same throughput with higher latency and slower checkpoints because barriers travel behind more data. That is the reasoning behind buffer debloating, which tunes in-flight data *down* to a target time rather than up. Fix the busy operator instead.

A traffic jam is longest behind the accident and the road ahead is empty. You do not diagnose it by measuring the queue; you drive forward until you reach the first car that is moving slowly with nothing in front of it.

saying these in an interview costs you the question

  • Blaming the operator that shows the most backpressure
  • Reading busy time as CPU utilization
  • Raising parallelism to fix a single hot key
  • Enlarging network buffers to make backpressure disappear
  • Ignoring checkpoint duration as a corroborating signal

context