Backpressure appears when a downstream stage consumes slower than upstream arrival; backlog grows until capacity, latency or retention runs out.
Stream backpressure and catch-up budget
Measure the rate mismatch
If a source emits 1,700 records per second and the sink sustains 1,300, backlog grows by about 400 per second. A 720,000-record backlog then needs more than mere parity to clear. At 2,100 per second of sustainable processing against 1,700 new records, net catch-up is 400 per second, or 30 minutes. Use measured sustained rates, not a brief peak benchmark.
Find the slow boundary
Inspect source lag, per-operator busy time, output buffer occupancy, checkpoint duration, sink acknowledgements and largest keyed partition. A slow sink can make every upstream operator look busy because flow control propagates backward. Hot keys can create the same symptom on one task. Scale the actual bottleneck or reduce its work.
Budget for recovery
A consumer that only matches normal arrival rate can never catch up after an outage. Specify maximum tolerated lag and outage duration, then reserve enough spare throughput to drain backlog within that window. Include restore time before useful processing resumes. Source retention must exceed outage plus catch-up plus a safety margin, or the missing records will expire.
Avoid destructive emergency fixes
Increasing parallelism can help only when partitions, keys and downstream service limits allow it. Disabling checkpoints may improve momentary throughput but removes a reliable recovery point. Dropping events to make a dashboard green violates the data contract. Shed optional enrichment or route overflow to a durable queue only if the consumer explicitly accepts that behavior.
Verify under a realistic incident
Throttle the sink, let lag accumulate for a measured interval, then restore capacity. Plot backlog slope before and after, checkpoint completion time and event-time freshness. Compare the actual catch-up time with the calculation. If one input partition saturates, adding workers alone may leave the slope unchanged.
Implementation
def catch_up_seconds(backlog_records, incoming_per_second, processing_per_second):
spare = processing_per_second - incoming_per_second
if backlog_records < 0 or incoming_per_second < 0 or spare <= 0:
raise ValueError("catch-up needs positive spare throughput")
return backlog_records / spare
assert catch_up_seconds(720_000, 1_700, 2_100) == 1_800
assert catch_up_seconds(24_000, 500, 700) == 120Performance and operating cost
The estimate is O(1) arithmetic but assumes stable rates. In a real pipeline, processing capacity changes with key distribution, checkpoint load and sink throttling. Backlog consumes broker storage; open windows and deduplication state consume stream storage. Track both, and recalculate from observed net drain rather than treating a one-time capacity test as a guarantee.
Common Mistakes
- Do not interpret upstream busy time as proof the upstream stage is slow.
- Do not plan recovery with zero spare throughput.
- Do not disable checkpoints or drop records without a written loss contract.
Read next
- Keyed state, checkpoints and recovery
- Pipeline SLOs, freshness and error budgets
- Failure classification and retry budgets
- Project: release a recoverable streaming risk ledger
- Skewed keys and salted aggregation
Continue the workflow: Pipeline RPO, RTO and replication lag.
