skip to content

In a Spring WebFlux app on Reactor Netty, how does backpressure propagate end-to-end from the HTTP client, across the network, to your reactive pipeline?

level: principalimportance: should knowfreq 40%

answer

  1. demand <-> Channel.isWritable + write watermarks
  2. slow client -> Netty stops requesting -> upstream slows
  3. crosses network via TCP receive window (sliding window)
  4. inbound: Netty auto-read off + release DataBuffers
  5. weakest link: block()/unbounded buffer/collectList breaks it

basics

~20 s

Reactor Netty maps reactive demand onto TCP flow control. When your response Publisher is slow (or a slow client stops reading), Netty stops requesting/writing, its socket buffers fill, the TCP receive window shrinks, and the sender is throttled — backpressure crosses the network via TCP.

solid answer

~50 s

WebFlux runs on Reactor Netty, which bridges Reactive Streams demand to TCP flow control. On the response path, your controller returns a Publisher; Netty subscribes and requests items as the socket can accept writes. If the client reads slowly, Netty's write buffers fill and it stops requesting more from your Publisher — backpressure flows upstream into your pipeline. Under the hood this rides TCP's sliding-window: a slow reader shrinks the advertised receive window, the OS stops ACKing new data, and the sender's congestion/flow window blocks further sends. On the request path, DataBuffers from the inbound body are released as you consume them; if you consume slowly, Netty pauses reading (channel auto-read off), the receive buffer fills, and the client is throttled the same way. The key principal insight: end-to-end backpressure is only as strong as its weakest link — an unbounded operator, a block(), or buffering the whole body breaks the chain and reintroduces OOM risk.

code

java · 14 lines
java
// GOOD: streaming preserves backpressure end-to-end (DB -> Netty -> TCP -> client)
@GetMapping(value = "/events", produces = MediaType.APPLICATION_NDJSON_VALUE)
public Flux<Event> stream() {
    return repository.findAllStreaming()   // R2DBC honors demand
        .limitRate(100);                   // bound per-round-trip fetch feeding the response
    // Netty requests from this Flux only while the socket is writable;
    // a slow client -> channel not writable -> DB fetch slows. No unbounded heap.
}

// BAD: collectList() materializes everything, decoupling heap from TCP flow control
@GetMapping("/events-bad")
public Mono<List<Event>> broken() {
    return repository.findAllStreaming().collectList(); // OOM risk on large sets
}

go deeper

for a junior

Understand at a high level that a slow client eventually slows your data source because demand propagates back.

for a middle

Explain that Reactor Netty maps demand to socket writability and that TCP flow control carries it over the network.

for a senior

Detail inbound (auto-read + DataBuffer release) and outbound (write watermarks) paths and the patterns that break the chain.

for a principal

Reason about the whole chain as layered flow control, tune watermarks/codec limits/prefetch, and design pipelines that never decouple heap from TCP backpressure.

## The layers involved A WebFlux request traverses: **TCP socket ↔ Reactor Netty channel ↔ HttpHandler/DispatcherHandler ↔ your reactive pipeline (controller + operators) ↔ downstream (DB/WebClient)**. Backpressure must propagate through all of them to be meaningful. Reactor Netty is the adapter that connects Reactive Streams demand to the network transport. ## Response (outbound) path 1. Your controller returns a `Flux<T>`/`Mono<T>` (or a `Publisher` of `DataBuffer`s after encoding). 2. Reactor Netty **subscribes** to that Publisher and issues `request(n)` based on whether the underlying **Netty channel is writable** (Netty exposes `Channel.isWritable()` and a configurable high/low write watermark). 3. If the client reads slowly, the OS socket send buffer fills, Netty's channel becomes **not writable**, and Netty **stops requesting** more elements from your Publisher. 4. That lack of demand propagates upstream through your operators to the data source (e.g. R2DBC stops fetching rows). Production naturally slows to match the client. ## How it crosses the network: TCP flow control TCP has a built-in **sliding-window / receive-window** mechanism. Each side advertises how much buffer space it has. A **slow reader** drains its receive buffer slowly, so it advertises a **smaller (eventually zero) window**; the sender's OS then **stops sending** until the window reopens. This is distinct from congestion control but achieves the same throttling. So application-level Reactive Streams demand is ultimately **backed by TCP's own backpressure** — the two are layered, not competing. ## Request (inbound) path 1. The request body arrives as a `Flux<DataBuffer>` (pooled, reference-counted buffers). 2. Reactor Netty uses Netty's **auto-read**: when your consumer is keeping up, it reads more from the socket; when you consume slowly (or stop requesting), Netty **disables auto-read**, the socket receive buffer fills, and TCP flow control throttles the client's upload. 3. You must **release** `DataBuffer`s (`DataBufferUtils.release`) as you consume them; the codecs handle this for typed `@RequestBody` binding. Failing to consume the body breaks the chain. ## Where the chain breaks (the principal-level failure modes) End-to-end backpressure is only as strong as its weakest link. Common breaks: - **`block()` / `toFuture().get()`** in the pipeline — requests unbounded and blocks a thread; destroys demand-based flow and can starve the event loop. - **Unbounded `onBackpressureBuffer()`** or `Flux.create(..., BUFFER)` — absorbs the mismatch into heap, decoupling from TCP flow control → OOM under load. - **Buffering the whole body/response** (`collectList()`, `bodyToMono(String.class)` on a huge payload) — materializes everything in memory regardless of demand. - **A blocking driver** (JDBC) wrapped naively — blocks event-loop threads; use `Schedulers.boundedElastic` or a reactive driver. - **Fixed-rate sources** (`Flux.interval`) inside the pipeline — ignore demand and overflow as covered earlier. ## WebClient symmetry `WebClient` (also Reactor Netty by default) applies the same model as a client: if you consume the response `Flux<DataBuffer>` slowly, it stops reading and TCP throttles the server. Streaming with `bodyToFlux` + bounded consumption preserves backpressure; `bodyToMono(byte[].class)` buffers it all. ## Operational tuning - Netty **write buffer watermarks** (`WRITE_BUFFER_WATER_MARK`) control when the channel flips writable/unwritable — the knob behind response backpressure sensitivity. - Codec **max-in-memory-size** (`spring.codec.max-in-memory-size`) bounds how much a buffering codec will hold, a safety net against unbounded body buffering. - `prefetch`/`limitRate` on DB streams bound per-round-trip fetches feeding the response. ## The takeaway Reactive Streams demand and TCP flow control are two layers of the same idea; Reactor Netty stitches them together so a slow client automatically throttles your database. Your job is to avoid the operators/patterns that silently buffer and thereby decouple your heap from the network's natural backpressure.

  • Concretely, what network mechanism carries backpressure across the wire?
    TCP flow control via the receive-window / sliding-window. A slow reader advertises a shrinking (eventually zero) window, so the sender's OS stops transmitting until the window reopens — throttling the sender without any application-level messages.
  • Name three code patterns that silently break end-to-end backpressure.
    block()/toFuture().get() (unbounded blocking), unbounded onBackpressureBuffer() or Flux.create(BUFFER), and fully materializing payloads with collectList()/bodyToMono(byte[]) — all decouple heap from demand/TCP and risk OOM.
  • On the inbound path, how does Netty throttle a fast uploading client?
    It disables channel auto-read when your consumer lags, so it stops reading from the socket; the receive buffer fills, TCP advertises a smaller window, and the client's OS stops sending until you consume (and release) more DataBuffers.

saying these in an interview costs you the question

  • Claiming WebFlux has some custom network protocol for backpressure instead of leveraging TCP flow control
  • Thinking backpressure survives a block() or collectList() in the pipeline
  • Believing an unbounded buffer preserves end-to-end backpressure
  • Not connecting Netty channel writability/watermarks to Reactive Streams demand
  • Forgetting DataBuffer release / that inbound uses auto-read toggling

context