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以上大文件时因全量加载内存超出缓冲区限制失败,且不想通过提高缓冲区上限做临时修复。
错误尝试及问题
- 合并所有
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))
本质仍是将整个文件全量加载到内存导致的问题。
- 其他无效尝试:
- 订阅
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
相关产品推荐
相关产品推荐

