Every system has a slowest component. When work arrives faster than that component can finish it, the excess has to go somewhere: into a queue, into memory, into a retry loop, or onto the floor. Backpressure is the architecture that makes this decision deliberately. A slow stage tells the stage before it to send less, and that signal travels back towards the source until something with the authority to wait, reject or degrade acts on it.

Without backpressure, overload shows up as unbounded queues, rising latency, memory exhaustion and eventually a crash that takes healthy requests down with the bad ones. This article covers why unbounded queues fail, sizing bounded ones with Little's Law, the three ways to push back, flow control in TCP, HTTP/2, gRPC, Reactive Streams and Kafka, and carrying the signal across services without a retry storm.

Advertisement

What backpressure is, and what it is not

Backpressure is flow control: a feedback path from consumer to producer that limits the rate or amount of in-flight work to what the consumer can absorb. It is related to, but distinct from, three neighbouring techniques. Rate limiting caps what a client is allowed to send regardless of how busy the server is; it is a policy, not a feedback signal. Load shedding discards work that has already arrived when the system is past saturation; it is what you do when backpressure has nowhere further upstream to go. A system has working backpressure when every buffer in the chain is bounded and every stage reacts to a full buffer in a defined way.

Backpressure: demand flows upstream, data flows downstream, overflow is decided at the edgeClientsretry budget + backoffEdge / gatewayadmission controlService Abounded queue: 200Service Bconcurrency limit: 64requests503 / slowqueue full429 + Retry-AfterKafka topicdurable buffer, pullConsumerpause() / resume()async workpollstop pollingTransport windowsTCP receive window, HTTP/2 WINDOW_UPDATE, gRPC isReadyThree responses to excess loadblock (wait) | signal (credits) | reject or dropSignals to watchqueue depth, age of oldest item, rejects, utilizationEvery queue is bounded. Every boundary returns a signal. Retries are budgeted so the signal is not amplified.
A request path and an asynchronous path. Each hop has a bounded buffer and returns a signal upstream: 503 or slow responses between services, 429 at the edge, pause and resume on a Kafka consumer, and flow-control windows in the transport underneath.

Why an unbounded queue fails: a worked example

Consider a service that can complete 1,000 requests per second. Traffic rises to 1,200 requests per second for one minute. The queue grows by 200 items every second, so after 60 seconds it holds 12,000 items. A new request now waits behind 12,000 others, which at 1,000 per second is 12 seconds of queueing before any work starts.

Now add the client timeout of 2 seconds. Requests that waited longer have been abandoned, but the server still processes them. For most of that minute the service is running at full capacity and producing nothing anyone will use. This is goodput collapse: throughput stays high while useful throughput falls towards zero.

Little's Law gives the fix. For a stable system, the average number of items in the system L equals the arrival rate λ multiplied by the average time each spends there W. Turn it round to size a buffer: if the service processes 1,000 requests per second and the latency budget for queueing is 200 milliseconds, then no more than 1,000 × 0.2 = 200 items should ever be waiting. A queue bounded at 200 caps queueing delay at about 200 milliseconds; the 201st request gets a fast, honest rejection instead of a slow, useless success.

Advertisement

Three ways a stage can push back

A bounded stage must choose what happens when it is full. Most real systems use all three answers at different layers.

Block. The producer waits until space frees up. This is what a bounded blocking queue does in-process, and what a socket write does when the kernel send buffer is full. Blocking is simple and lossless, but it holds a thread or connection while it waits, and if two stages can each block waiting for the other, the pipeline deadlocks.

Signal. The consumer tells the producer how much it is willing to receive, and the producer never sends more. This credit-based flow control underlies Reactive Streams, HTTP/2, gRPC and TCP's receive window.

Reject or drop. The stage refuses new work immediately with an error, or discards some items. Rejection is the only option at a boundary where you cannot make the caller wait indefinitely, such as a public API. Dropping is correct for data whose value decays quickly, such as metrics samples or video frames.

ResponseMechanismGood forCost
Blockproducer waits on a full bounded bufferin-process pipelines, file and socket I/Oties up threads; can deadlock in cycles
Signalconsumer grants credits or demand; producer never sends morestreams, reactive libraries, HTTP/2 and gRPCevery stage must honour the protocol
Rejectfast error (429, 503, RESOURCE_EXHAUSTED) at admissionrequest/response servicesclients must back off, not hammer
Dropdiscard newest, oldest or lowest prioritytelemetry, video, sampled eventsdata loss must be acceptable

Credit-based flow control with Reactive Streams

The Reactive Streams specification, adopted into the JDK as java.util.concurrent.Flow, makes demand explicit. A subscriber calls subscription.request(n) to say it can accept n more items; the publisher must not call onNext more times than the total requested. Project Reactor, RxJava and Akka Streams implement the same contract.

The subscriber below asks for the next batch of 32 only when it has finished the current one.

import java.util.concurrent.Flow;

final class BatchingSubscriber<T> implements Flow.Subscriber<T> {
    private static final int BATCH = 32;
    private Flow.Subscription subscription;
    private int outstanding;

    @Override public void onSubscribe(Flow.Subscription s) {
        this.subscription = s;
        outstanding = BATCH;
        s.request(BATCH);                 // initial credit
    }

    @Override public void onNext(T item) {
        process(item);                    // slow work happens here
        if (--outstanding == 0) {         // batch finished: grant more credit
            outstanding = BATCH;
            subscription.request(BATCH);
        }
    }

    @Override public void onError(Throwable t) { log(t); }
    @Override public void onComplete() { flush(); }

    private void process(T item) { /* write to a slow sink */ }
    private void log(Throwable t) { }
    private void flush() { }
}

The common mistake is calling request(Long.MAX_VALUE), which the specification treats as unbounded demand, and then doing slow work asynchronously. Demand has been granted for everything, so the publisher floods an internal buffer and the flow control is gone. For sources that cannot slow down, Reactor's onBackpressureBuffer, onBackpressureDrop and onBackpressureLatest make the buffer-or-drop choice explicit.

Backpressure in the transport: TCP, HTTP/2 and gRPC

TCP has flow control built in. The receiver advertises a receive window, the number of bytes it can buffer; the sender may not have more unacknowledged bytes in flight than that window. When an application stops reading from a socket, the kernel buffer fills, the advertised window shrinks to zero and the sender's writes block or return would-block.

HTTP/2 multiplexes many streams over one TCP connection, so it adds its own per-stream and per-connection windows, replenished with WINDOW_UPDATE frames. gRPC streaming is built on those windows. The trap is that application code can bypass them by writing into an unbounded in-memory queue. In gRPC Java, a server streaming a large response should check ServerCallStreamObserver.isReady() and register setOnReadyHandler to resume sending when the transport has room, rather than calling onNext in a tight loop.

void streamRows(Request req, StreamObserver<Row> raw) {
    ServerCallStreamObserver<Row> out = (ServerCallStreamObserver<Row>) raw;
    Iterator<Row> rows = db.scan(req);
    boolean[] done = {false};             // the ready handler can fire again after the last row
    Runnable drain = () -> {
        // send only while the HTTP/2 window has room; resume on the next ready callback
        while (!done[0] && out.isReady() && rows.hasNext()) {
            out.onNext(rows.next());
        }
        if (!done[0] && !rows.hasNext()) {
            done[0] = true;
            out.onCompleted();
        }
    };
    out.setOnReadyHandler(drain);
    drain.run();
}

Pull-based consumers: Kafka and durable queues

Log-based brokers such as Kafka give backpressure almost for free because consumers pull. A consumer that falls behind simply has lag: the messages wait durably in the broker, not in the consumer's memory. The broker is the bounded buffer, sized by retention rather than RAM, and producers are not slowed at all. See message queues compared.

Backpressure still matters inside the consumer. A common design polls a batch and hands records to a worker pool through a local queue. If the workers are slower than polling, that local queue grows without limit. The fix is to bound the local queue and stop fetching when it is full, using consumer.pause() on the assigned partitions while continuing to call poll() so the consumer stays in the group, then resume() when the queue drains. Keep max.poll.records small enough to finish inside max.poll.interval.ms.

from queue import Queue
from confluent_kafka import Consumer

work = Queue(maxsize=500)          # bounded hand-off to worker threads
HIGH, LOW = 450, 100
paused = False
consumer = Consumer({"bootstrap.servers": BROKERS, "group.id": "orders",
                     "enable.auto.commit": False})
consumer.subscribe(["orders"])

while running:
    msg = consumer.poll(0.5)       # keep polling even while paused: liveness
    if msg is not None and not msg.error():
        work.put(msg)              # never blocks for long: we pause before full
    depth = work.qsize()
    if not paused and depth >= HIGH:
        consumer.pause(consumer.assignment()); paused = True
    elif paused and depth <= LOW:
        consumer.resume(consumer.assignment()); paused = False
    commit_completed_offsets(consumer)   # commit only what workers finished

The high and low water marks provide hysteresis so the consumer does not flap between paused and resumed on every message. Committing only completed offsets keeps delivery at-least-once, so downstream handlers must be idempotent; see idempotency architecture.

Carrying the signal across service boundaries

Between synchronous services the signal is an error or latency. A service whose queue is full should fail fast with HTTP 503 or 429, or gRPC RESOURCE_EXHAUSTED or UNAVAILABLE, and may include a Retry-After header. The caller must treat that as a request to send less, not as a transient glitch to retry immediately. This is where most backpressure designs break: a three-tier call chain where each tier makes up to three attempts multiplies the load on the bottom tier by up to 27 during an incident.

Two techniques keep the signal intact. Retry budgets allow retries only while they are a small fraction, say 10 percent, of recent requests, with exponential backoff and jitter. Adaptive concurrency limits grow the allowed in-flight count slowly while latency is stable and cut it sharply when latency rises, like TCP congestion control; Netflix's concurrency-limits library and Envoy's adaptive concurrency filter implement variants.

import asyncio, time

class AimdLimiter:
    """Adaptive in-flight limit: +1 on healthy completions, halve on overload."""
    def __init__(self, initial=20, floor=4, ceiling=500, target_ms=150):
        self.limit, self.floor, self.ceiling = initial, floor, ceiling
        self.target_ms, self.in_flight = target_ms, 0

    def try_acquire(self) -> bool:
        if self.in_flight >= self.limit:
            return False                  # caller returns 503 immediately
        self.in_flight += 1
        return True

    def release(self, latency_ms: float, overloaded: bool):
        self.in_flight -= 1
        if overloaded or latency_ms > self.target_ms:
            self.limit = max(self.floor, int(self.limit * 0.5))
        else:
            self.limit = min(self.ceiling, self.limit + 1)

async def handle(limiter, request, downstream):
    if not limiter.try_acquire():
        return 503, {"Retry-After": "1"}   # fast, honest rejection
    start, overloaded = time.monotonic(), False
    try:
        return 200, await asyncio.wait_for(downstream(request), timeout=0.5)
    except (asyncio.TimeoutError, OverloadedError):
        overloaded = True                   # downstream pushed back
        return 503, {"Retry-After": "1"}
    finally:
        limiter.release((time.monotonic() - start) * 1000, overloaded)

A production limiter smooths latency against a measured baseline; the shape is what matters: admit what the dependency can handle, reject the rest in microseconds, recover gradually.

Failure modes to design against

  • Hidden unbounded buffers. Java's Executors.newFixedThreadPool uses an unbounded LinkedBlockingQueue; many HTTP clients, log appenders and in-memory channels do the same. Audit every queue in the request path and give each an explicit capacity and a rejection policy.
  • Buffer bloat. Bounded buffers at every layer still add up to seconds of work; size them against the end-to-end latency budget.
  • Deadlock from blocking cycles. If stage A blocks waiting on B and B blocks waiting on A, for example through a shared thread pool used for both calling and callback work, the system stops. Use separate pools per dependency, as in the bulkhead pattern, or non-blocking signalling.
  • Work for departed callers. Propagate deadlines, for example the gRPC deadline, and drop queued items whose deadline has already passed before starting them. This alone often restores goodput during overload.

Operating backpressure

Track, for every bounded buffer: current depth against capacity, the age of the oldest item, the rejection or pause rate, and consumer utilization. Queue age is the best single alert, because it is directly comparable with the latency objective. For Kafka consumers, alert on lag expressed in time rather than message count.

Test overload on purpose. A load test that ramps past capacity should show latency rising to the bound and then flattening while rejections rise, with goodput staying near capacity. If latency keeps climbing, there is an unbounded buffer somewhere; if goodput collapses, work is being done for callers that have left. When rejections are sustained, hand off to shedding and prioritisation, and route poisoned or repeatedly failing messages to a dead-letter queue rather than letting them block a partition.

The trade-off: a system with backpressure fails visibly and early, with 429s and 503s where an unbounded design would look healthy for a few seconds more, seconds borrowed from the crash that follows.

What to do next

  1. Draw your request path and list every buffer on it: load balancer, server accept queue, thread pools, connection pools, in-memory channels and broker topics.
  2. Give each buffer an explicit capacity derived from Little's Law and your latency budget, and a named policy for when it is full: block, signal, reject or drop.
  3. Replace fixed unbounded executors with bounded queues and a rejection handler that returns 503 or RESOURCE_EXHAUSTED quickly.
  4. In streaming code, honour transport readiness (gRPC isReady, Reactive Streams request(n)) instead of writing into unbounded queues.
  5. For Kafka consumers, bound the hand-off queue and use pause() and resume() with high and low water marks.
  6. Add retry budgets, exponential backoff with jitter and deadline propagation to every client.
  7. Alert on age of oldest queued item and on rejection rate, and run an overload test that confirms goodput stays near capacity.
Key takeaway: Backpressure is a feedback loop from the slowest stage to the source: bound every buffer, size it with Little's Law against the latency you can afford, and choose deliberately between blocking, credit signals and fast rejection at each boundary. Use the flow control that TCP, HTTP/2, gRPC, Reactive Streams and pull-based brokers already provide instead of bypassing it with hidden queues, keep retries budgeted so the signal is not amplified, and hand off to load shedding only when there is nowhere further upstream to push.