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

如何无内存开销聚合多个WebClient的Flux<DataBuffer>为Zip响应

问题:流式聚合多文件为Zip(避免内存溢出)

我需要一种无需将所有文件完全加载到内存或文件系统,就能把多个文件聚合为Zip的方法。目前通过WebClient调用获取文件,返回Flux<DataBuffer>形式的流,代码如下:

public Mono<ResponseEntity<Flux<DataBuffer>>> getDocumentContentAsDataBuffer(
    String objectId) {
  return momentumApiClient
      .get()
      .uri("/objects/" + objectId + "/contents/file")
      .retrieve()
      .toEntityFlux(DataBuffer.class);
}

之后我为多个对象调用该接口,将结果存入BufferedFile对象(仅包含文件名和Flux流)。部分对象无内容时,Flux为空,会创建对应文件夹:

List<BufferedFile> zipDataBuffers =
    objects.stream()
        .map(
            o ->
                new BufferedFile(
                    (objectMap.get(o).getContentStreams() != null
                            && !objectMap.get(o).getContentStreams().isEmpty())
                        ? client
                            .getDocumentContentAsDataBuffer(
                                objectMap.get(o).getObjectId(), jwt)
                            .flatMapMany(
                                response ->
                                    response.getBody() != null
                                        ? response.getBody()
                                        : Flux.empty())
                        : Flux.empty(),
                    generatePath(objectMap, o)))
        .toList();

// create zip
return ZipDataBuffer.toFlux(zipDataBuffers);

ZipDataBuffer.toFlux需要把所有文件和文件夹聚合为Zip返回给客户端。目前找到的可行方案是下载完整文件后写入ZipOutputStream,处理大量小文件没问题,但大文件会被完全加载到内存中,引发内存溢出。

当前ZipDataBuffer的实现如下:

import java.io.InputStream;
import java.util.List;
import java.util.zip.ZipEntry;
import java.util.zip.ZipOutputStream;
import lombok.SneakyThrows;
import org.springframework.core.io.buffer.DataBuffer;
import org.springframework.core.io.buffer.DataBufferUtils;
import org.springframework.core.io.buffer.DefaultDataBuffer;
import org.springframework.core.io.buffer.DefaultDataBufferFactory;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;

public class ZipDataBuffer {
  final DefaultDataBufferFactory factory = new DefaultDataBufferFactory();
  final DefaultDataBuffer buffer = factory.allocateBuffer(1024);
  final ZipOutputStream zos = new ZipOutputStream(buffer.asOutputStream());

  public static Flux<DataBuffer> toFlux(List<BufferedFile> resources) {
    ZipDataBuffer zipDataBuffer = new ZipDataBuffer();
    return Flux.fromIterable(resources)
        .flatMap(zipDataBuffer::write)
        .concatWith(Mono.fromCallable(zipDataBuffer::finish));
  }

  public Mono<DataBuffer> write(BufferedFile bufferedFile) {
    return DataBufferUtils.join(bufferedFile.getFile())
        .defaultIfEmpty(factory.allocateBuffer(0))
        .map(
            dataBuffer ->
                {
                  write(bufferedFile.getFilename(), dataBuffer);
                  DataBufferUtils.release(dataBuffer);
                  return read();
                });
  }

  @SneakyThrows
  private synchronized void write(String filename, DataBuffer dataBuffer) {
    ZipEntry zipEntry = new ZipEntry(filename);
    zos.putNextEntry(zipEntry);
    try (InputStream is = dataBuffer.asInputStream()) {
      is.transferTo(zos);
    }
    zos.closeEntry();
  }

  private DataBuffer read() {
    return buffer.split(buffer.writePosition());
  }

  @SneakyThrows
  public DataBuffer finish() {
    zos.finish();
    zos.close();
    return read();
  }
}

该实现无法处理50GB级别的大文件,会因内存不足崩溃。我尝试过多种策略,但核心问题是客户端下载速度慢于应用的WebClient下载速度,导致应用缓冲区快速填满并崩溃。

有没有办法让WebClient仅在缓冲区有空闲时才请求数据?这是我至今没解决的问题,甚至怀疑是否存在这样的方案。

最小可复现示例
import java.io.IOException;
import java.io.InputStream;
import java.util.zip.ZipEntry;
import java.util.zip.ZipOutputStream;
import lombok.SneakyThrows;
import lombok.extern.slf4j.Slf4j;
import org.springframework.core.io.buffer.DataBuffer;
import org.springframework.core.io.buffer.DataBufferUtils;
import org.springframework.core.io.buffer.DefaultDataBuffer;
import org.springframework.core.io.buffer.DefaultDataBufferFactory;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RestController;
import org.springframework.web.reactive.function.client.WebClient;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.core.scheduler.Schedulers;

@RestController
public class MinimalWorkingExample {
  final WebClient webclient;

  public MinimalWorkingExample(WebClient.Builder webclientBuilder) {
    this.webclient = webclientBuilder.build();
  }

  // example one where the complete file is written to the zip file buffer and then read
  @GetMapping(value = "/exampleOutOfMemory", produces = "application/zip")
  public Flux<DataBuffer> exampleOutOfMemory() {
    final DefaultDataBufferFactory factory = new DefaultDataBufferFactory();
    final DefaultDataBuffer outputBuffer = factory.allocateBuffer(1024);
    final ZipOutputStream zos = new ZipOutputStream(outputBuffer.asOutputStream());
    // aggregate all Databuffers into one
    return DataBufferUtils.join(
            webclient
                .get()
                .uri("https://files.icyflamestudio.com/50MB.zip")
                //.uri("http://speedtest.tele2.net/50GB.zip")
                .retrieve()
                .bodyToFlux(DataBuffer.class))
        .map(
            // write the aggregated Databuffer to a zip file
            databuffer ->
                {
                  write("testfile", zos, databuffer);
                  DataBufferUtils.release(databuffer);
                  // Problem: will read after the complete file has been written to the zip file buffer
                  // I want a method which returns this buffer for example every 1024 bytes
                  return read(outputBuffer);
                })
        .concatWith(Mono.fromCallable(() -> finish(zos, outputBuffer)));
  }

  // example two where every Databuffer is written to the zip file buffer and then read
  // This will never finish for big files
  @GetMapping(value = "/exampleCorruptFile", produces = "application/zip")
  public Flux<DataBuffer> exampleCorruptFile() {
    final DefaultDataBufferFactory factory = new DefaultDataBufferFactory();
    final DefaultDataBuffer outputBuffer = factory.allocateBuffer(1024);
    final ZipOutputStream zos = new ZipOutputStream(outputBuffer.asOutputStream());
    // call a function where every Databuffer is written to the zip output Stream buffer and will be
    // read immediately after
    return write(
            "testfile",
            zos,
            outputBuffer,
        webclient
            .get()
            //.uri("http://speedtest.tele2.net/50GB.zip")
            .uri("https://files.icyflamestudio.com/50MB.zip")
            .retrieve()
            .bodyToFlux(DataBuffer.class))
        .concatWith(Mono.fromCallable(() -> finish(zos, outputBuffer)));
  }

  @SneakyThrows
  private synchronized void write(String filename, ZipOutputStream zos, DataBuffer dataBuffer) {
    ZipEntry zipEntry = new ZipEntry(filename);
    zos.putNextEntry(zipEntry);
    try (InputStream is = dataBuffer.asInputStream()) {
      is.transferTo(zos);
    }
    zos.closeEntry();
  }

  private Flux<DataBuffer> write(
      String filename,
      ZipOutputStream zos,
      DataBuffer outputDataBuffer,
      Flux<DataBuffer> inputDataBuffer) {
    return inputDataBuffer
        .publishOn(Schedulers.boundedElastic())
        .map(
            inputBuffer ->
                {
                  try (InputStream is = inputBuffer.asInputStream()) {
                    is.transferTo(zos);
                  } catch (IOException e) {
                    throw new RuntimeException(e);
                  }
                  return read(outputDataBuffer);
                })
        .doOnSubscribe(
            subscription ->
                {
                  ZipEntry zipEntry = new ZipEntry(filename);
                  try {
                    zos.putNextEntry(zipEntry);
                  } catch (Exception e) {
                    throw new RuntimeException(e);
                  }
                })
        .doOnTerminate(
            () ->
                {
                  try {
                    zos.closeEntry();
                  } catch (Exception e) {
                    throw new RuntimeException(e);
                  }
                });
  }

  private DataBuffer read(DataBuffer buffer) {
    return buffer.split(buffer.writePosition());
  }

  @SneakyThrows
  public DataBuffer finish(ZipOutputStream zos, DataBuffer buffer) {
    zos.finish();
    zos.close();
    return read(buffer);
  }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 20:07:32