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

Spring Boot 3迁移RestTemplate到WebClient:大文件转InputStream报错

Spring Boot 3 WebClient 流式处理大文件适配 StreamingResponseBody

问题背景

原使用RestTemplate的代码可直接将响应转为byte[]再包装成InputStream返回:

ResponseEntity<byte[]> responseEntity = restTemplate.getForEntity("myApi.com/file/12345", byte[].class);
retrieveResponse.setInputStream(new ByteArrayInputStream(responseEntity.getBody()));

升级到Spring Boot 3改用WebClient后,处理256KB以上大文件时因全量加载内存超出缓冲区限制失败,且不想通过提高缓冲区上限做临时修复。

错误尝试及问题

  1. 合并所有DataBuffer再转InputStream,触发内存溢出错误:
Flux<DataBuffer> responseBody = webClient
  .get("myApi.com/file/12345")
  .uri()
  .exchangeToFlux(responseEntity -> {
    return responseEntity.bodyToFlux(DataBuffer.class);
  });
InputStream inputStream = responseBody.reduce(DataBuffer::write).map(DataBuffer::asInputStream).block();
retrieveResponse.setInputStream(inputStream);

错误信息:

java.lang.IndexOutOfBoundsException: writerIndex(754) + minWritableBytes(2358) exceeds maxCapacity(754): PooledSlicedByteBuf(ridx: 0, widx: 754, cap: 754/754, unwrapped: PooledUnsafeDirectByteBuf(ridx: 1179, widx: 1179, cap: 1208))

本质仍是将整个文件全量加载到内存导致的问题。

  1. 其他无效尝试:
  • 订阅DataBuffer设置InputStream:无法保证执行顺序,无法正确流式输出
responseBody.subscribe(dataBuffer -> retrieveContentResponse.setInputStream(dataBuffer.asInputStream()));
  • 仅取第一个DataBuffer转InputStream:返回200但无法完整写入文件,出现无限加载
retrieveResponse.setInputStream(responseBody.blockFirst().asInputStream());

约束条件

响应类RetrieveResponse实现了StreamingResponseBody,大量遗留代码依赖InputStream字段无法替换:

public class RetrieveResponse implements StreamingResponseBody {
  private static final int BUFFER_SIZE = 8192;
  private InputStream inputStream;

  @Override
  public void writeTo(@NonNull OutputStream outputStream) throws IOException {
      int bytesRead;
      try (InputStream in = new BufferedInputStream(inputStream);
           OutputStream out = new BufferedOutputStream(outputStream)) {
          byte[] buf = new byte[BUFFER_SIZE];
          while ((bytesRead = in.read(buf)) != -1) {
              out.write(buf, 0, bytesRead);
          }
          out.flush();
      } catch (Exception e) {
          log.error("Exception when writing to output stream", e);
          throw new HttpException("something bad", e);
      }
  }
}

解决方案

方案1:直接对接流式DataBuffer与OutputStream(推荐)

核心思路是跳过InputStream转换,直接将WebClient返回的Flux<DataBuffer>流式写入StreamingResponseBody的目标OutputStream,彻底避免内存加载瓶颈。

修改RetrieveResponse类:

public class RetrieveResponse implements StreamingResponseBody {
    private static final int BUFFER_SIZE = 8192;
    private final Flux<DataBuffer> dataBufferFlux;

    // 构造方法传入WebClient返回的流式DataBuffer
    public RetrieveResponse(Flux<DataBuffer> dataBufferFlux) {
        this.dataBufferFlux = dataBufferFlux;
    }

    @Override
    public void writeTo(@NonNull OutputStream outputStream) throws IOException {
        // 利用DataBufferUtils将流式数据直接写入OutputStream,自动管理缓冲区
        DataBufferUtils.write(dataBufferFlux, outputStream)
                .doOnError(e -> {
                    log.error("Exception when writing data buffer to output stream", e);
                    DataBufferUtils.release(dataBufferFlux); // 异常时释放所有缓冲区资源
                })
                .block(); // 阻塞直到所有数据写入完成,适配StreamingResponseBody的同步调用逻辑
    }
}

调整WebClient调用代码:

RetrieveResponse retrieveResponse = webClient.get()
        .uri("myApi.com/file/12345")
        .exchangeToMono(response -> Mono.just(new RetrieveResponse(response.bodyToFlux(DataBuffer.class))))
        .block();

方案2:自定义流式InputStream适配(兼容遗留代码)

如果必须保留InputStream字段,可自定义一个适配类,将Flux<DataBuffer>的数据流转换为InputStream接口,实现真正的流式读取:

public class FluxDataBufferInputStream extends InputStream {
    private final Iterator<DataBuffer> bufferIterator;
    private byte[] currentBuffer;
    private int currentPosition;

    public FluxDataBufferInputStream(Flux<DataBuffer> dataBufferFlux) {
        this.bufferIterator = dataBufferFlux.toIterable().iterator();
    }

    @Override
    public int read() throws IOException {
        if (currentBuffer == null || currentPosition >= currentBuffer.length) {
            if (!bufferIterator.hasNext()) {
                return -1;
            }
            DataBuffer buffer = bufferIterator.next();
            try {
                currentBuffer = new byte[buffer.readableByteCount()];
                buffer.read(currentBuffer);
                currentPosition = 0;
            } finally {
                DataBufferUtils.release(buffer); // 读取完成后释放缓冲区资源
            }
        }
        return currentBuffer[currentPosition++] & 0xFF;
    }

    @Override
    public int read(byte[] b, int off, int len) throws IOException {
        if (currentBuffer == null || currentPosition >= currentBuffer.length) {
            if (!bufferIterator.hasNext()) {
                return -1;
            }
            DataBuffer buffer = bufferIterator.next();
            try {
                currentBuffer = new byte[buffer.readableByteCount()];
                buffer.read(currentBuffer);
                currentPosition = 0;
            } finally {
                DataBufferUtils.release(buffer);
            }
        }
        int bytesToRead = Math.min(len, currentBuffer.length - currentPosition);
        System.arraycopy(currentBuffer, currentPosition, b, off, bytesToRead);
        currentPosition += bytesToRead;
        return bytesToRead;
    }

    @Override
    public void close() throws IOException {
        super.close();
        // 关闭时清理所有未读取的缓冲区
        while (bufferIterator.hasNext()) {
            DataBufferUtils.release(bufferIterator.next());
        }
    }
}

使用方式:

Flux<DataBuffer> responseBody = webClient.get()
        .uri("myApi.com/file/12345")
        .exchangeToFlux(response -> response.bodyToFlux(DataBuffer.class));

RetrieveResponse retrieveResponse = new RetrieveResponse();
retrieveResponse.setInputStream(new FluxDataBufferInputStream(responseBody));

内容的提问来源于stack exchange,提问作者canpan14

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 13:58:36