For an infinite Flux served as SSE, how does Spring WebFlux flush items incrementally, and how does backpressure prevent a fast producer from overwhelming a slow client?
answer
- demand-driven: request(n) downstream->upstream
- Netty channel writability gates emission
- cold source pauses; hot firehose overflows
- onBackpressureLatest/Drop/Buffer, sample, limitRate
- doOnCancel/doFinally for disconnect cleanup
basics
~20 sThe reactive stack writes and flushes each emitted item as bytes become sendable, keeping the connection open. Backpressure flows from the TCP write buffer up through Reactor: the encoder only requests the next item when the socket can accept more, so a slow client naturally slows the producer.
solid answer
~50 sWebFlux runs on a non-blocking server (Netty by default). The `ServerSentEventHttpMessageWriter` subscribes to your `Flux`, and for each emitted item serializes it and writes it to the response, flushing per element so the client sees data immediately. Crucially, the whole pipeline is **demand-driven**: Reactor's `Subscription.request(n)` signals propagate from the network layer upward. Netty only pulls the next item when its socket write buffer has room (the OS TCP send buffer plus the client's receive window). If the client reads slowly, TCP flow-control fills the buffers, backpressure requests stop, and your producing `Flux` is paused — no unbounded in-memory queue builds up. For a source that ignores demand (a hot/infinite firehose), you must add an explicit strategy — `onBackpressureBuffer`, `onBackpressureDrop`, `onBackpressureLatest`, or `sample`/`limitRate` — otherwise you risk `OverflowException` or memory growth. This is the core reason reactive streaming scales to many slow clients.
code
java · 14 lines// A hot firehose (interval keeps ticking regardless of demand).
// Without a backpressure strategy a slow client would overflow the buffer.
@GetMapping(path = "/ticks", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<ServerSentEvent<Tick>> ticks() {
return marketFeed.hotTicks() // emits on its own schedule
.onBackpressureLatest() // keep only newest if client lags
.sample(Duration.ofMillis(250)) // cap UI update rate
.map(t -> ServerSentEvent.<Tick>builder()
.id(String.valueOf(t.seq()))
.data(t)
.build())
.doOnCancel(() -> log.info("client disconnected, releasing feed"))
.doFinally(signal -> marketFeed.release());
}go deeper
Know items are flushed as they arrive and the connection stays open.
Explain request(n) demand signaling and that a slow client throttles the producer.
Distinguish cold demand-respecting sources from hot firehoses and pick the right onBackpressure*/sample operator; handle disconnect cleanup.
Reason about Netty channel writability, event-loop starvation from blocking calls, memory-bounding strategy per SLA, and capacity planning for many concurrent slow clients.
**The problem.** An infinite/fast `Flux` (say, market ticks every millisecond) feeding a slow client (a phone on 3G) could, in a naive push model, queue unbounded data in server memory. Reactive Streams exists to prevent exactly this via **backpressure**. **Backpressure = demand signaling.** In Reactive Streams / Project Reactor, a `Subscriber` calls `subscription.request(n)` to ask for at most `n` more items; the `Publisher` must not emit more than requested. Demand flows **downstream-to-upstream**: the ultimate consumer's readiness governs how fast the source runs. **How it works end-to-end in WebFlux SSE:** 1. The controller returns `Flux<T>`; Spring picks `ServerSentEventHttpMessageWriter` because of `produces=text/event-stream`. 2. The writer's reactive bridge to **Reactor Netty** subscribes to the Flux. 3. For each item: serialize -> wrap in an SSE frame -> write to the Netty channel -> **flush** so it goes out immediately (SSE flushes per item; that's what makes it a live stream rather than a buffered response). 4. Netty tracks channel **writability**. When the OS TCP send buffer fills (because the client's TCP receive window is closing — i.e., the client isn't reading fast enough), the channel becomes non-writable. 5. While non-writable, the adapter stops issuing `request(n)`. Reactor propagates this upstream, so your source Flux stops emitting (if it respects demand). 6. When the client drains data and the window reopens, the channel becomes writable, `request(n)` resumes, and emission continues. Thus a slow client transparently throttles the producer — no thread is blocked and memory stays bounded. **The catch: cold vs hot / demand-respecting sources.** - A **cold, demand-respecting** source (e.g., `Flux.range`, a reactive DB cursor via R2DBC, `Flux.generate`) naturally pauses — it only produces when requested. Backpressure just works. - A **hot / external firehose** that emits on its own schedule (a `Sinks.Many`, a websocket feed, `Flux.interval`, sensor callbacks) does **not** slow down just because you stopped requesting. If downstream demand is zero and the source keeps pushing, Reactor's default fixed-size buffer overflows -> `reactor.core.Exceptions$OverflowException` and the stream errors/terminates. **Explicit backpressure operators for hot sources:** - `onBackpressureBuffer([maxSize])` — queue overflow items (bounded, with an overflow strategy/callback). Risks memory if unbounded. - `onBackpressureDrop()` — discard items that can't be delivered (fine for "latest value only" telemetry). - `onBackpressureLatest()` — keep only the most recent item, drop older undelivered ones (great for live prices/dashboards). - `limitRate(n)` — cap how many items are prefetched/requested at a time. - `sample(Duration)` / `sampleTimeout(...)` — emit at most one item per window, naturally rate-limiting a firehose to what a UI can render. **Gotchas.** - Forgetting `flush` semantics: SSE writer flushes per event, but if you accidentally serve as a buffered type, nothing streams. - Assuming backpressure protects you from a hot source — it doesn't unless the source respects demand or you add an operator; otherwise you get overflow. - Blocking calls inside the map/flatMap on a Netty event-loop thread stall ALL connections on that loop — keep the pipeline non-blocking (offload blocking work with `subscribeOn(Schedulers.boundedElastic())`). - Detecting client disconnect: use `doOnCancel`/`doFinally` to release resources (subscriptions, DB cursors) when the client goes away — a cancel signal propagates up when the connection closes.
- Why can a cold source like Flux.range stream to a slow client safely, while Flux.interval can overflow?Flux.range is demand-respecting: it only produces the next value when request(n) arrives, so a slow client simply slows it down. Flux.interval is time-driven and emits regardless of downstream demand; if demand is zero its items pile into Reactor's bounded buffer and eventually throw OverflowException. Hot/time-driven sources need onBackpressure* or sample/limitRate.
- How do you release server resources when an SSE client disconnects?A closed connection sends a cancel signal up the reactive chain. Hook doOnCancel (client-initiated) and/or doFinally (any termination: complete, error, cancel) to unsubscribe from feeds, close R2DBC cursors, or release pooled resources. Otherwise you leak subscriptions per abandoned connection.
saying these in an interview costs you the question
- Believing backpressure automatically protects against any source, including hot ones
- Saying WebFlux buffers the whole infinite stream in memory
- Doing blocking I/O on the event-loop thread inside the stream
- Not cleaning up on client disconnect (no doOnCancel/doFinally)
- Confusing flush-per-item with request/response buffering