How do you consume a streaming HTTP response with WebClient as a Flux<T>, and how does bodyToFlux differ from bodyToMono in backpressure and memory behavior?
answer
- retrieve().bodyToFlux(T.class)
- incremental decode + backpressure
- bodyToMono buffers one value
- consume body → connection release
- 256KB maxInMemorySize on Mono
basics
~20 sCall retrieve() (or exchangeToFlux) then bodyToFlux(T.class). It decodes the response body incrementally into a stream of T with backpressure, so you process elements as they arrive without buffering the whole response. bodyToMono decodes into a single value/object, buffering that one payload.
solid answer
~40 sUse webClient.get()...retrieve().bodyToFlux(T.class) to decode a streaming response element-by-element. WebClient reads network DataBuffers as they arrive and the decoder (e.g. Jackson) emits T's incrementally; Reactor propagates backpressure to Netty so a slow consumer throttles reads instead of buffering unboundedly. This suits large collections, NDJSON, or SSE (text/event-stream) — often you'll see per-element flushes. bodyToMono(T.class), by contrast, aggregates the body into one value: for a POJO it buffers the full payload before emitting, for a large list it holds it all in memory. Choose bodyToFlux to keep memory flat and start processing early; bodyToMono when you genuinely need the whole thing as one object. For status/header control use exchangeToFlux/exchangeToMono. Remember: nothing happens until you subscribe.
code
java · 15 lines// Streaming decode, flat memory, backpressured
Flux<Quote> quotes = webClient.get()
.uri("/quotes")
.accept(MediaType.APPLICATION_NDJSON) // or TEXT_EVENT_STREAM for SSE
.retrieve()
.bodyToFlux(Quote.class);
quotes.take(100) // cancels the response after 100
.doOnNext(this::process) // acts on each as it arrives
.subscribe();
// Single aggregated object (buffers whole payload)
Mono<Report> report = webClient.get().uri("/report")
.retrieve()
.bodyToMono(Report.class);go deeper
Knows bodyToFlux gives a stream of T, bodyToMono a single object.
Explains incremental decoding and that bodyToMono buffers the full payload.
Adds backpressure-to-Netty, connection release/leak, maxInMemorySize, cancellation semantics.
Discusses pool exhaustion failure modes, SSE lifecycle, codec tuning, and choosing streaming vs aggregation by SLA/memory budget.
**Consuming a response reactively.** After building the request you call either `retrieve()` (convenient, auto-throws `WebClientResponseException` on 4xx/5xx) or `exchangeToFlux(...)`/`exchangeToMono(...)` (full access to `ClientResponse` for status/headers/conditional decoding). Then you extract the body: - **`bodyToFlux(Class<T>)`** — decode the body into a **`Flux<T>`**, a stream of elements. WebClient reads inbound network bytes as `DataBuffer`s and feeds a **`Decoder`** (e.g. `Jackson2JsonDecoder`). As elements are parsed they are emitted downstream immediately. **Backpressure**: if your subscriber processes slowly, Reactor requests fewer items, which translates into Reactor Netty reading fewer bytes off the socket — so memory stays bounded even for a huge or infinite stream. - **`bodyToMono(Class<T>)`** — decode into a **single `Mono<T>`**. For a POJO it collects the entire body and emits once. If `T` is a `List<Foo>`, the *whole list* is buffered in memory before emission — no incremental processing. **Streaming formats where bodyToFlux shines.** - **SSE** (`text/event-stream`): `bodyToFlux(MyEvent.class)` yields a potentially infinite stream of events; the connection stays open. - **NDJSON** (`application/x-ndjson`): each line becomes one `T`. - **A big JSON array** (`application/json`): Jackson's non-blocking parser can still emit array elements one at a time via `bodyToFlux`, so you don't hold the whole array — a real advantage over blocking clients. **Memory / latency.** `bodyToFlux` gives **flat memory** and **early first-element latency** (you act on element 1 before element N arrives). `bodyToMono` gives you a single object but pays full-buffer memory and waits for completion. **Gotchas.** - **Nothing runs until subscribe.** The returned publisher is cold; forgetting to subscribe/return it means no request is sent. - **Connection release.** With `retrieve().bodyToFlux(...)` the response body must be fully consumed (or cancelled) so the connection returns to the pool. If you `exchangeTo...` / use `ClientResponse`, you must consume or release the body (`bodyToFlux`, or `releaseBody()`), otherwise you **leak connections** and eventually stall (pool exhaustion). This is a classic `exchange()`-era bug the newer `exchangeToFlux`/`retrieve()` APIs mitigate. - **Backpressure isn't infinite buffering avoidance for SSE** — but Netty won't read faster than requested. - **Cancellation**: cancelling the `Flux` subscription aborts the HTTP response; fine for 'take first N' but the server may still be sending. - **Error decoding**: `retrieve()` maps error statuses to `WebClientResponseException`; customize with `onStatus(...)`. - **In-memory codec limit**: `bodyToMono` of a large payload can hit the default 256 KB `maxInMemorySize` and throw `DataBufferLimitException` — configure via `ExchangeStrategies`/`codecs()`. `bodyToFlux` of individually-small elements typically avoids this because each element is small.
- Why can consuming the response body be mandatory even if you don't care about it?Because until the body is fully read, cancelled, or released, the underlying connection isn't returned to the pool. Not consuming it leaks connections and eventually exhausts the pool, stalling further requests.
- How is bodyToFlux able to avoid buffering a large JSON array entirely?Jackson's non-blocking (async) parser tokenizes the incoming DataBuffers and emits each array element as it's fully parsed, so only one element is in flight at a time rather than the whole array.
- What is DataBufferLimitException and when do you hit it?It's thrown when a buffered body exceeds the codec's maxInMemorySize (default 256 KB). Typical with bodyToMono of a big payload; raise the limit via ExchangeStrategies codecs configuration.
saying these in an interview costs you the question
- Thinking bodyToFlux still buffers the entire response
- Not consuming/releasing the body and leaking connections
- Believing the request fires without subscribing
- Assuming bodyToMono streams a large list incrementally
- Ignoring the 256KB maxInMemorySize default