skip to content

How do you stream an uploaded part's content as DataBuffers to storage without buffering the whole file in memory?

level: seniorimportance: must knowfreq 40%

answer

  1. content() -> Flux<DataBuffer>
  2. transferTo / DataBufferUtils.write stream + release
  3. reference-counted: must release()
  4. doOnDiscard(DataBuffer, release)
  5. peak memory ~ one chunk + maxInMemorySize

basics

~10 s

Use part.content(), which is a Flux<DataBuffer> streaming the body in chunks. Write it with DataBufferUtils.write(content, path) or FilePart.transferTo(path). Both stream chunk-by-chunk with backpressure and release each buffer.

solid answer

~40 s

Every Part exposes content() returning Flux<DataBuffer> — the body as a stream of chunks, not one big array. To persist without full buffering you subscribe reactively: FilePart.transferTo(Path) is the shortcut; under the hood it uses DataBufferUtils.write(content(), channel) which writes each DataBuffer as it arrives, honoring backpressure so you only pull as fast as the sink accepts. Crucially, every DataBuffer is reference-counted (may be pooled/off-heap), so you must release it — transferTo and DataBufferUtils.write do this for you; if you consume content() manually you must call DataBufferUtils.release(buffer) or you leak native memory. Peak memory stays at roughly one chunk (plus the reader's maxInMemorySize threshold per part) rather than the whole file. This is what lets WebFlux handle multi-GB uploads on a small heap.

code

java · 17 lines
java
// Manual streaming that still releases buffers correctly
Mono<Void> streamToChannel(FilePart part, AsynchronousFileChannel channel) {
    // DataBufferUtils.write subscribes with backpressure and releases each buffer
    return DataBufferUtils.write(part.content(), channel)
            .map(DataBufferUtils::release) // release written buffers
            .then();
}

// If you consume content() yourself, release on every path:
Mono<Long> hashUpload(FilePart part, MessageDigest digest) {
    return part.content()
            .doOnNext(buf -> digest.update(buf.toByteBuffer()))
            .doOnDiscard(DataBuffer.class, DataBufferUtils::release) // dropped/cancelled
            .map(buf -> { long n = buf.readableByteCount();
                          DataBufferUtils.release(buf); return n; })
            .reduce(0L, Long::sum);
}

go deeper

for a junior

Know content() gives a Flux<DataBuffer> and transferTo streams it.

for a middle

Know buffers are chunks and prefer transferTo/DataBufferUtils.write over reading bytes.

for a senior

Explain reference counting, mandatory release, backpressure, and doOnDiscard.

for a principal

Reason about off-heap pool exhaustion, leak detection, and end-to-end streaming to object storage with bounded memory.

**DataBuffer.** `org.springframework.core.io.buffer.DataBuffer` is Spring's abstraction over a chunk of bytes, backed by Netty `ByteBuf`, `java.nio.ByteBuffer`, or a byte array. In WebFlux the body of a request/part is exposed as `Flux<DataBuffer>` — a *stream* of chunks. `Part.content()` returns exactly this for a multipart part. **Why streaming matters.** Loading a whole upload into a `byte[]` or String forces the entire file into heap — an OOM risk and a scalability killer. Reactive streaming keeps only a bounded window in memory: the reader emits a chunk, your sink writes it, the buffer is released, backpressure throttles the next read. **The idiomatic writes:** - `FilePart.transferTo(Path)` — highest level; returns `Mono<Void>`. - `DataBufferUtils.write(Publisher<DataBuffer>, Path, OpenOption...)` — write a flux to a file. - `DataBufferUtils.write(content, AsynchronousFileChannel, position)` — for an explicit channel. All of these subscribe with backpressure and **release** each buffer after writing. **Reference counting — the #1 gotcha.** Netty-backed `DataBuffer`s are pooled and reference-counted (`PooledDataBuffer`). If you obtain buffers from `content()` and do your own processing, you **must** release each one — `DataBufferUtils.release(buffer)` — typically in a `doOnDiscard`, `doFinally`, or after reading. Forgetting leaks off-heap memory (Netty logs `LEAK` warnings). Operators that drop elements (filter, take, cancellation, errors) must also release discarded buffers; `DataBufferUtils.releaseConsumer()` and `Flux#doOnDiscard(DataBuffer.class, DataBufferUtils::release)` help. **Backpressure.** Because it's a `Flux`, a slow disk/S3 sink naturally slows the read from the socket — the framework won't over-read. This is the core advantage over Servlet blocking IO. **maxInMemorySize interaction.** The reader (`DefaultPartHttpMessageReader`) buffers up to `maxInMemorySize` bytes per part in memory before spilling to a temp file; `content()` then streams from there. So even before your code runs, a large part isn't fully in heap. **Reading into memory intentionally.** If you *do* need bytes, `DataBufferUtils.join(content())` concatenates into a single buffer (`Mono<DataBuffer>`) — but that reintroduces full buffering, so bound it with `maxInMemorySize` and use only for small parts. **When to use.** Any large-file upload, streaming to disk/object storage, or transforming bytes on the fly (hashing, encryption) — process `content()` chunk-by-chunk and release.

  • Why must DataBuffers be released and what happens if you don't?
    Netty-backed buffers are pooled and reference-counted (PooledDataBuffer). Not releasing them leaks off-heap/native memory; Netty emits LEAK warnings and eventually the pool exhausts. transferTo and DataBufferUtils.write release for you; manual consumption needs DataBufferUtils.release, including on discard/error paths.
  • How does backpressure protect the server during a large upload?
    content() is a Flux, so a slow sink (disk/S3) reduces demand upstream; the reactive stack throttles socket reads to match, keeping only a bounded window of buffers in memory instead of pulling the whole file at line rate.
  • What does DataBufferUtils.join do and why is it risky?
    It concatenates the Flux<DataBuffer> into one Mono<DataBuffer>, i.e. fully buffers the content. Handy for small parts but reintroduces the OOM risk for large uploads, so bound it with maxInMemorySize.

saying these in an interview costs you the question

  • Calling DataBufferUtils.join on a large upload and defeating streaming
  • Consuming content() without releasing buffers
  • Forgetting to release on filter/error/cancel paths
  • Blocking on part.content().collectList().block()

context