Backpressure

Backpressure is the propagation of downstream capacity limits toward upstream producers or operators. When a consumer cannot accept more work, upstream progress slows or blocks. The bottleneck may be a sink, a hot key, computation, or another resource; high backpressure identifies a symptom rather than a unique cause.

If arrivals remain at 1,000 records per second and processing handles 800, pending work grows by 200 per second, or 120,000 in ten minutes. At 1,500 per second of processing with the same arrivals, net drain is 500 per second and that backlog needs 240 seconds to clear in this constant-rate model. Matching arrival rate only stops growth; it does not clear existing work.

A buffer absorbs a temporary mismatch but adds no processing capacity. If the real-world source cannot slow down, backlog can move to a broker or eventually exceed retention and lose input. Observe throughput, queue age and capacity across the path, then address the bottleneck. Adding workers may not help an indivisible hot key or a rate-limited external sink.

Reference: Flink backpressure monitoring.


Discover more from Insightful Data Lab

Subscribe to get the latest posts sent to your email.

Similar Posts

Questions, corrections, or additional insights?

This site uses Akismet to reduce spam. Learn how your comment data is processed.