Every streaming pipeline eventually meets a moment when data arrives faster than some stage can process it. A downstream database slows, an enrichment service starts timing out, a traffic spike hits, or one partition carries far more than its share. Something has to give. Backpressure is the mechanism that makes the slow stage's limit visible upstream, so that the system slows down in a controlled way instead of filling memory until a process dies.
This article is about backpressure inside stream-processing pipelines: chains of operators fed by a log such as Kafka, written with Reactive Streams libraries, Flink or plain consumer code. The relay case, where one service forwards a gRPC or WebSocket stream to another, is covered in backpressure in bidirectional streams. Here we start from the arithmetic of queues, look at how three widely used systems implement backpressure, work through a lag incident with numbers, and finish with what to measure and how to choose between blocking, buffering and dropping.
The arithmetic that makes backpressure unavoidable
Model any stage as a queue in front of a server. Items arrive at rate λ and the stage completes them at rate μ. If λ stays below μ, the queue empties between bursts. If λ exceeds μ for a duration t, the queue grows by (λ − μ) × t, and nothing inside the stage can change that. Little's law adds the latency consequence: the average time an item spends in the system equals the average number of items in it divided by the throughput. A queue of a million items in front of a stage completing 40,000 per second means 25 seconds of added delay, whatever the hardware.
So when λ exceeds μ there are exactly four options. Block: tell the producer to slow down, which moves the queue upstream. Buffer: hold the excess somewhere, which only works if the overload ends and there is spare capacity afterwards to drain it. Drop: discard or sample the excess on purpose. Scale: raise μ by adding capacity at the bottleneck. Real systems combine them. Backpressure is the first option implemented end to end, so that the queue ends up somewhere durable and cheap rather than in a JVM heap.
The architecture at a glance
The picture shows the property you want. The bottleneck stage is busy all the time. Every stage before it is blocked waiting for permission to send. The stage after it is idle. And the excess is not in any process's memory: it is in the topic, as unread offsets, which costs disk rather than heap and survives restarts. The rest of the article is about how each layer achieves this.
Demand signalling: Reactive Streams
Reactive Streams is a small specification for asynchronous stream processing with non-blocking backpressure, adopted in Java 9 as java.util.concurrent.Flow and implemented by libraries such as Project Reactor, RxJava and Akka Streams. It defines four interfaces: Publisher, Subscriber, Subscription and Processor. The rule that matters is that a publisher may send at most as many onNext signals as the subscriber has requested through Subscription.request(n). Demand flows upstream; data flows downstream; neither side blocks a thread.
import java.util.concurrent.Flow;
/** Requests in batches and replenishes at half, so the publisher never runs ahead. */
final class BatchingSubscriber<T> implements Flow.Subscriber<T> {
private static final int BATCH = 64;
private Flow.Subscription sub;
private int outstanding; // requested but not yet delivered
@Override public void onSubscribe(Flow.Subscription s) {
sub = s;
outstanding = BATCH;
s.request(BATCH); // initial demand: nothing flows before this
}
@Override public void onNext(T item) {
process(item); // the slow part; demand is not renewed until done
if (--outstanding <= BATCH / 2) {
sub.request(BATCH - outstanding); // top demand back up to BATCH
outstanding = BATCH;
}
}
@Override public void onError(Throwable t) { log(t); }
@Override public void onComplete() { flush(); }
}Two details make this work in practice. Requesting one item at a time is correct but slow, because every item costs a signal; batching demand amortises it. And operators between source and sink must preserve demand: an operator that buffers internally without a bound, or that requests unbounded demand (request(Long.MAX_VALUE)) from its upstream, silently disconnects the chain. Most library operators that do this are named for it, such as buffer or drop variants with a strategy argument, and they force you to choose what happens on overflow.
Credits between machines: Flink
Apache Flink applies the same idea across the network. Since version 1.5 its network stack uses credit-based flow control. Each receiving subtask has a pool of network buffers for each input channel. It announces free buffers to the sender as credits, and the sender transmits a buffer of records only when it holds a credit. When the receiver is slow, its buffers fill, it stops granting credits, the sender's output buffers fill, and the sender's own processing stops. That stall propagates hop by hop until it reaches the source, which stops pulling from Kafka. Because credits are per channel, one slow channel does not block the shared TCP connection for others.
Flink exposes where time goes. Each subtask reports busyTimeMsPerSecond, backPressuredTimeMsPerSecond and idleTimeMsPerSecond, which add up to roughly 1,000. The web UI shows the same as busy and backpressured ratios. To find the bottleneck, look for the furthest-upstream operator that is busy near 100 percent while the operators before it are backpressured. Fixing anything upstream of it changes nothing.
Backpressure interacts badly with aligned checkpoints, because checkpoint barriers travel with the data and wait behind full buffers, so checkpoints slow down exactly when the job is under stress. Flink offers two mitigations. Unaligned checkpoints let barriers overtake buffered data and store that in-flight data in the checkpoint. Buffer debloating, added in 1.14 and enabled with taskmanager.network.memory.buffer-debloat.enabled, shrinks the amount of in-flight data to what the subtask can consume within a target time, set by taskmanager.network.memory.buffer-debloat.target. Less buffered data means faster barriers and smaller unaligned checkpoints, at some cost in throughput for very bursty jobs.
Pull consumers: Kafka
A Kafka consumer pulls. It fetches records when it calls poll, so a slow consumer cannot be overwhelmed by the broker; it just falls behind, and the gap between the latest offset and its committed offset is consumer lag. That makes the topic the natural place for excess to wait. The trap is in what the consumer does with polled records.
If processing happens inline and takes too long, the gap between poll calls can exceed max.poll.interval.ms. The consumer is then considered failed, leaves the group and triggers a rebalance, and its partitions are reassigned, often to another consumer that is equally slow. The result is a rebalance storm with no progress. Lower max.poll.records so each batch finishes in time, or hand records to a worker pool and use pause and resume so that poll keeps being called without fetching more.
// Bounded hand-off from a Kafka poll loop to a worker pool, with pause/resume.
BlockingQueue<ConsumerRecord<String, byte[]>> work = new ArrayBlockingQueue<>(10_000);
int high = 8_000, low = 2_000;
boolean paused = false;
while (running) {
ConsumerRecords<String, byte[]> batch = consumer.poll(Duration.ofMillis(200));
for (ConsumerRecord<String, byte[]> r : batch) {
work.put(r); // fits while max.poll.records <= 2,000
}
int depth = work.size();
if (!paused && depth > high) {
consumer.pause(consumer.assignment()); // poll() keeps the member alive,
paused = true; // but returns no records
} else if (paused && depth < low) {
consumer.resume(consumer.paused());
paused = false;
}
commitProcessedOffsets(consumer); // only offsets whose records finished
}The high and low watermarks give hysteresis, so the consumer does not flap between paused and resumed on every poll. Commit only offsets whose records have finished processing; committing on hand-off turns a crash into data loss. If one partition is far hotter than the rest, backpressure will hold the whole consumer to its pace, and the real fix is in the key design, covered in Kafka partitioning strategies.
Bounding calls to external services
The most common bottleneck is not CPU in the pipeline but a call to something outside it: a profile service, a feature store, a database. Unbounded concurrency there is not throughput, it is an attack on your dependency. Bound in-flight requests explicitly, and let the bound block the reader.
import asyncio
async def enrich_stream(events, lookup, max_in_flight=200):
"""Call an external service per event without ever exceeding max_in_flight requests."""
slots = asyncio.Semaphore(max_in_flight)
async def one(ev):
try:
ev["profile"] = await asyncio.wait_for(lookup(ev["user_id"]), timeout=0.5)
except asyncio.TimeoutError:
ev["profile"] = None # shed the enrichment, keep the event
metrics.incr("enrich.timeout")
finally:
slots.release() # free the slot whether it worked or not
return ev
pending = set()
async for ev in events:
await slots.acquire() # no free slot: stop reading the source
pending.add(asyncio.create_task(one(ev)))
done = {t for t in pending if t.done()}
for t in done:
yield t.result()
pending -= done
for t in asyncio.as_completed(pending):
yield await tFlink's async I/O operator takes the same idea as a capacity parameter: once that many requests are pending, the operator stops accepting input and backpressure takes over. Choose the bound from the dependency's capacity, not your own: if the profile service handles 8,000 requests per second at 25 ms, Little's law says about 200 requests in flight saturate it, and more only adds queueing inside the service. Pair the bound with a timeout and a load shedding policy for when the dependency is down, so backpressure does not turn an outage there into an unbounded stall here.
Worked example: a two-hour peak
A clickstream job reads a topic receiving 50,000 events per second at peak and 30,000 off-peak. Events average 1 KB. The enrichment stage, bounded as above, completes 40,000 per second; everything else is faster.
At peak the backlog grows by 10,000 events per second: 36 million events, about 36 GB of lag, per hour. Flink's metrics show enrich busy near 100 percent with source and parse backpressured, which correctly identifies the bottleneck. No process runs out of memory, because the excess stays in Kafka. After a two-hour peak, lag is 72 million events. When input falls to 30,000 per second, the job has 10,000 per second of spare capacity, so draining takes another two hours. At the worst point an event waits about 72 million divided by 40,000, roughly 30 minutes, before it is processed.
Three things follow. Topic retention must comfortably exceed the worst backlog duration, or lagging data is deleted before it is read. Adding consumers helps only up to the partition count. And whether 30 minutes of delay is acceptable is a product decision: if dashboards need data within a minute, the answer is more enrichment capacity, a cache in front of the profile service, or sampling during peaks, not more buffering.
Failure modes
- An unbounded queue somewhere. One stage with an unbounded in-memory buffer absorbs the backpressure signal and then dies of memory exhaustion, losing everything it held. Every queue needs a bound and an overflow policy.
- A source that cannot be paused. UDP, webhooks and device telemetry push regardless. Land them in a durable log first so the pipeline behind it can apply backpressure.
- Rebalance storms. Slow inline processing exceeds max.poll.interval.ms and consumers churn without progress.
- Retries that multiply load. A slow dependency triggers retries, which add load to the slowest thing in the system. Bound retries and add jitter.
- Scaling on CPU. An I/O-bound bottleneck shows low CPU while lag climbs, so CPU-based autoscaling never fires. Scale on lag or backpressure.
- Cycles. A pipeline that feeds output back into its own input can deadlock when every stage waits for credit from the next. Give feedback edges their own bounded, droppable buffer.
- Checkpoints that never finish. Under sustained backpressure aligned checkpoints time out repeatedly; use unaligned checkpoints or buffer debloating.
Operating it: what to measure
| Signal | Why | Alert on |
|---|---|---|
| Consumer lag in time, not only records | Records mean nothing without the rate | Lag time above your freshness target |
| Busy and backpressured ratios per operator | Locates the bottleneck | Any operator busy above 90 percent for long periods |
| Queue depth for every bounded buffer | Shows where pressure accumulates | Time spent above the high watermark |
| End-to-end latency, event time to output | What users feel | Percentiles above target |
| Dropped or sampled events | Shedding must be visible | Any drop outside a declared policy |
| Strategy | Use when | Price |
|---|---|---|
| Block (backpressure) | Every event matters and delay is tolerable | Latency grows; needs a durable buffer upstream |
| Bounded buffer | Short bursts with spare capacity afterwards | Memory and recovery time |
| Drop or sample | Freshness matters more than completeness | Lost data, which must be counted |
| Scale out | Sustained overload at a parallelisable stage | Cost; limited by partitions and dependencies |
What to do next
- Draw your pipeline and mark every queue; give each one an explicit bound and overflow behaviour.
- Make sure excess ends up in a durable log, and set retention longer than your worst plausible backlog.
- In Flink, read busy and backpressured ratios to find the real bottleneck before tuning anything; try buffer debloating if checkpoints slow under load.
- In Kafka consumers, either keep each poll batch well inside max.poll.interval.ms or move work to a bounded pool with pause and resume.
- Bound concurrency to every external dependency from its capacity, with timeouts and a shedding policy.
- Alert on lag in seconds and on end-to-end latency, and autoscale on lag rather than CPU.
- Decide with the product owner which streams may be sampled under overload, and count every dropped event.