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.

Advertisement

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

Kafka topicdurable bufferSourcepauses readingParsewaits for creditEnrichthe bottleneckSinkdatabase writesno free buffers: no creditupstream stalls in turnlag grows herebackpressuredbackpressuredbusy ~100%mostly idleThree responses to rate mismatchBlock: slow the producer (credits, demand)Buffer: absorb bursts in a bounded queueDrop or sample: shed excess deliberatelyPlus: add capacity at the bottleneckQueue arithmeticarrival 50k/s, capacity 40k/sbacklog grows 10k/s = 36M per hourdelay = backlog / drain ratedrain only with spare capacity
A slow enrichment stage stops granting credit to its upstream. Parse and source stall in turn, the source stops reading, and the excess accumulates as consumer lag in the durable Kafka topic, where it is safe and measurable.

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.

Advertisement

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 t

Flink'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

SignalWhyAlert on
Consumer lag in time, not only recordsRecords mean nothing without the rateLag time above your freshness target
Busy and backpressured ratios per operatorLocates the bottleneckAny operator busy above 90 percent for long periods
Queue depth for every bounded bufferShows where pressure accumulatesTime spent above the high watermark
End-to-end latency, event time to outputWhat users feelPercentiles above target
Dropped or sampled eventsShedding must be visibleAny drop outside a declared policy
StrategyUse whenPrice
Block (backpressure)Every event matters and delay is tolerableLatency grows; needs a durable buffer upstream
Bounded bufferShort bursts with spare capacity afterwardsMemory and recovery time
Drop or sampleFreshness matters more than completenessLost data, which must be counted
Scale outSustained overload at a parallelisable stageCost; limited by partitions and dependencies

What to do next

  1. Draw your pipeline and mark every queue; give each one an explicit bound and overflow behaviour.
  2. Make sure excess ends up in a durable log, and set retention longer than your worst plausible backlog.
  3. In Flink, read busy and backpressured ratios to find the real bottleneck before tuning anything; try buffer debloating if checkpoints slow under load.
  4. 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.
  5. Bound concurrency to every external dependency from its capacity, with timeouts and a shedding policy.
  6. Alert on lag in seconds and on end-to-end latency, and autoscale on lag rather than CPU.
  7. Decide with the product owner which streams may be sampled under overload, and count every dropped event.
Key takeaway: When arrival rate exceeds a stage's capacity the queue must grow somewhere, so design where. Propagate backpressure end to end, through Reactive Streams demand, Flink credits and Kafka pause and resume, until the excess sits as lag in a durable log, then bound every in-memory queue and external call, locate the true bottleneck from busy and backpressured time, and choose deliberately between blocking, buffering, dropping and scaling.