What memory, backpressure, and buffer-lifecycle pitfalls must you handle when consuming reactive multipart uploads, and how do you build a safe zero-buffer streaming pipeline?
answer
- concatMap not flatMap (sequential window)
- release on read/discard/cancel/error
- doOnDiscard(DataBuffer.class, release)
- limits: maxInMemorySize/maxParts/maxDiskUsagePerPart
- PartEvent -> WebClient = constant-memory relay
basics
~10 sConsume parts sequentially (concatMap), always release DataBuffers on every path including discard/cancel/error, set reader limits (maxInMemorySize/maxParts/maxDiskUsagePerPart), let backpressure throttle reads, and prefer PartEvent to relay uploads downstream without touching heap or disk.
solid answer
~40 sThree classes of pitfall. (1) Ordering: the parser is a sequential window over one connection, so parts and their content must be consumed in order — use concatMap, never concurrent flatMap, and always drain a part's body before advancing. (2) Buffer lifecycle: DataBuffers are pooled/reference-counted; you must release them on every terminal path — successful read, filter/discard, cancellation, and error — via DataBufferUtils.release and Flux#doOnDiscard(DataBuffer.class, ...). Leaks exhaust Netty's off-heap pool. (3) Resource limits: set maxInMemorySize, maxParts, maxDiskUsagePerPart, maxHeadersSize so a malicious client can't OOM or fill disk (zip-bomb-style). Backpressure means a slow sink (S3) naturally throttles socket reads, so you never over-read. For a true pass-through gateway, use Flux<PartEvent> and re-emit events into a WebClient body so bytes never hit heap or a temp file. Combine with a request-size/timeout guard at the edge.
code
java · 21 lines// Constant-memory upload relay: client -> this service -> upstream, zero disk/heap buffering
@PostMapping(value = "/proxy", consumes = MediaType.MULTIPART_FORM_DATA_VALUE)
Mono<ResponseEntity<Void>> proxy(@RequestBody Flux<PartEvent> events) {
return webClient.post()
.uri("/upstream/upload")
.contentType(MediaType.MULTIPART_FORM_DATA)
// PartEvents stream straight through; nothing is fully buffered
.body(BodyInserters.fromProducer(events, PartEvent.class))
.retrieve()
.toBodilessEntity();
}
// Safe manual consumption: release on EVERY path
Mono<Void> drainSafely(Flux<Part> parts) {
return parts.concatMap(part ->
part.content()
.map(buf -> { /* use bytes */ DataBufferUtils.release(buf); return buf; })
.doOnDiscard(DataBuffer.class, DataBufferUtils::release) // cancel/error/filter
.then())
.then();
}go deeper
Know uploads should stream and buffers must be released.
Use concatMap and prefer transferTo; know the reader limits exist.
Handle release on all terminal paths and reason about backpressure end-to-end.
Design a DoS-hardened, constant-memory pipeline (limits + edge guard + PartEvent relay) and verify with leak detection.
**Why this is a principal-level concern.** Reactive multipart sits at the intersection of off-heap memory management, backpressure, and untrusted input. Getting it wrong causes native-memory leaks, event-loop stalls, disk exhaustion, or DoS. **1. Sequential consumption.** `DefaultPartHttpMessageReader` reads one socket stream; each `Part`'s `content()` is a *window* over it. The next `Part` is only emitted after the current part's body is fully consumed. Therefore: - Use `concatMap`, not `flatMap` (which would interleave/parallelize and corrupt ordering). - Never obtain a `Part` and skip its body — always `transferTo`, `DataBufferUtils.write`, or explicitly drain+release. **2. DataBuffer reference counting.** Netty-backed `DataBuffer` (a `PooledDataBuffer`) is reference-counted and lives off-heap. Every buffer you receive must be released **exactly once on every terminal path**: - Normal read: release after use. - Dropped by an operator (`filter`, `take`, `takeUntil`): use `.doOnDiscard(DataBuffer.class, DataBufferUtils::release)`. - Cancellation/error: reactor may discard in-flight elements — `doOnDiscard` covers these; also `DataBufferUtils.releaseConsumer()` as a subscriber. Failure mode: Netty leak detector logs `LEAK: ByteBuf.release() was not called`, the pooled arena grows, and under load you get native OOM or degraded throughput. The high-level helpers (`transferTo`, `DataBufferUtils.write`) release for you — prefer them. **3. Resource-limit / DoS hardening.** Configure on the reader: - `maxInMemorySize` — per-part heap ceiling before temp-file spill. - `maxDiskUsagePerPart` — cap temp-file growth (defends against a single huge part filling disk). - `maxParts` — cap part count (defends against millions of tiny parts). - `maxHeadersSize` — cap per-part header bytes. Also enforce an overall request size/time budget at the edge (gateway/filter) because these per-part limits don't bound the whole request. Clean up temp files: they are removed when the part is consumed/released, so a half-consumed request that errors must still complete/release to avoid orphaned temp files. **4. Backpressure end-to-end.** Because `content()` is a `Flux`, a slow downstream sink (object store, DB) reduces demand and the reactive stack throttles reads from the client socket — the server never pulls a multi-GB body faster than it can persist it. This is only true if you keep the chain reactive; a `block()` or `collectList()` breaks it and reintroduces unbounded buffering. **5. Zero-buffer streaming gateway (the payoff).** For proxying uploads to another service, bind `@RequestBody Flux<PartEvent>` and stream those events straight into a `WebClient` multipart request (`BodyInserters.fromProducer(events, PartEvent.class)` / `.body(fromProducer(...))`). `FilePartEvent`s carry content slices; nothing is ever fully buffered to heap or disk. This gives constant memory regardless of file size and full end-to-end backpressure. Release any events you drop. **Checklist:** concatMap; release on every path (doOnDiscard); reader limits set; edge size/time guard; keep the chain non-blocking; PartEvent for relays; verify with Netty leak detection in tests.
- How do you make sure buffers are released when the stream is cancelled or errors mid-part?Reactor discards in-flight elements on cancel/error; attach Flux#doOnDiscard(DataBuffer.class, DataBufferUtils::release) (or subscribe via DataBufferUtils.releaseConsumer). High-level helpers like transferTo/DataBufferUtils.write already handle all terminal paths, so prefer them over manual consumption.
- Which limits defend against a malicious multipart upload and what does each stop?maxInMemorySize caps per-part heap before temp-file spill; maxDiskUsagePerPart caps temp-file size (stops one huge part filling disk); maxParts caps part count (stops millions of tiny parts); maxHeadersSize caps header bytes. Add an overall request size/time budget at the edge since these are per-part.
- Why does block() or collectList() in the pipeline undermine reactive multipart?It removes backpressure and forces the entire content into memory, so the server pulls the whole body at line rate regardless of sink speed — reintroducing the OOM and stall risks the streaming model was designed to avoid, and can block an event-loop thread.
saying these in an interview costs you the question
- Using flatMap and assuming order/consumption is preserved
- Only releasing buffers on the happy path, not on cancel/error/discard
- Relying on per-part limits to bound total request size
- Using PartEvent but forgetting to release dropped events
- Calling block()/collectList() and claiming it is still streaming