Spring Controller中如何将Flux<DataBuffer>转换为StreamingResponseBody
如何将Flux流式写入StreamingResponseBody
需求:在Spring Controller中持有Flux<DataBuffer>对象,需要以流式方式将其写入StreamingResponseBody,避免将整个文件加载到内存中。
现有Controller代码:
@GetMapping(value="/attachment") public ResponseEntity<StreamingResponseBody> get(@PathVariable long attachmentId) { Flux<DataBuffer> dataBuffer = this.myService.getFile(attachmentId); StreamingResponseBody stream = out -> { // 此处应如何实现? } return ResponseEntity.ok() .contentType(MediaType.APPLICATION_OCTET_STREAM) .body(stream); } }
补充:myService.getFile()方法实现
public Flux<DataBuffer> getFile(long attachmentId) { return this.webClient.get() .uri("https://{host}/attachments/{attachmentId}", host, attachmentId) .attributes(clientRegistrationId("attachmentClient")) .accept(MediaType.ALL) .exchangeToFlux(clientResponse -> clientResponse.bodyToFlux(DataBuffer.class)); }
解决方案
核心思路是订阅Flux<DataBuffer>,将每个数据块实时写入输出流,同时确保异步流处理完成前保持同步阻塞,避免提前结束响应。关键要注意释放DataBuffer以防止内存泄漏,以及正确处理异常。
完整实现代码:
import org.springframework.core.io.buffer.DataBuffer; import org.springframework.core.io.buffer.DataBufferUtils; import org.springframework.http.MediaType; import org.springframework.http.ResponseEntity; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.PathVariable; import org.springframework.web.servlet.mvc.method.annotation.StreamingResponseBody; import reactor.core.publisher.Flux; import java.io.IOException; import java.io.OutputStream; import java.util.concurrent.CountDownLatch; import java.util.concurrent.atomic.AtomicReference; @GetMapping(value="/attachment") public ResponseEntity<StreamingResponseBody> get(@PathVariable long attachmentId) { Flux<DataBuffer> dataBufferFlux = this.myService.getFile(attachmentId); StreamingResponseBody stream = out -> { CountDownLatch latch = new CountDownLatch(1); AtomicReference<Throwable> errorRef = new AtomicReference<>(); dataBufferFlux.subscribe( buffer -> { try { // 读取DataBuffer内容并写入输出流 byte[] content = new byte[buffer.readableByteCount()]; buffer.read(content); out.write(content); out.flush(); // 按需刷新,确保数据及时发送到客户端 } catch (IOException e) { errorRef.set(e); } finally { // 必须释放DataBuffer,避免池化资源泄漏 DataBufferUtils.release(buffer); } }, // 处理流异常 error -> { errorRef.set(error); latch.countDown(); }, // 流处理完成时触发 latch::countDown ); try { // 同步等待流处理完成 latch.await(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new RuntimeException("Stream processing interrupted", e); } // 抛出订阅过程中捕获的异常 Throwable error = errorRef.get(); if (error != null) { if (error instanceof IOException) { throw (IOException) error; } else { throw new RuntimeException("Failed to process file stream", error); } } }; return ResponseEntity.ok() .contentType(MediaType.APPLICATION_OCTET_STREAM) .body(stream); }
关键要点说明
- DataBuffer释放:
DataBufferUtils.release(buffer)是必须的,因为Spring WebFlux使用池化的DataBuffer,不释放会导致内存泄漏。 - 同步阻塞:
CountDownLatch用于等待异步Flux处理完成,因为StreamingResponseBody.writeTo是同步执行的方法,必须阻塞直到所有数据写入完毕。 - 流式处理:每个DataBuffer到达时立即写入输出流,不会将整个文件加载到内存,完全符合流式传输的要求。
- 异常传递:通过
AtomicReference捕获订阅过程中的异常,在同步等待结束后统一抛出,确保错误能正确反馈给客户端。
内容的提问来源于stack exchange,提问作者user2480766
相关产品推荐
相关产品推荐

