You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.16 13:55:21