如何无内存开销聚合多个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
相关产品推荐
相关产品推荐

