如何用Spring WebClient流式压缩响应并转发至其他请求?
问题描述
需要从远程位置下载多个大文件,流式打包为ZIP后发送至另一目标地址,必须采用流式方案避免大文件导致的缓冲区问题。
已实现基于线程的方案:通过PipedOutputStream和PipedInputStream将ZIP内容写入后转发,代码如下:
void test() throws IOException { var outputLocation = "http://localhost:5000/anything"; var pipedInputStream = new PipedInputStream(); var pipedOutputStream = new PipedOutputStream(pipedInputStream); var thread = new Thread(() -> { try { var zip = new ZipOutputStream(pipedOutputStream); for (var i = 0; i < 10; i++) { // 实际应为远程文件下载请求,返回可流式写入ZIP的InputStream zip.putNextEntry(new ZipEntry("test" + i + ".txt")); new ByteArrayInputStream(("Hello " + i).getBytes()) .transferTo(zip); zip.closeEntry(); } zip.close(); } catch (Exception e) { e.printStackTrace(); } }); thread.start(); var response = webClient.post() .uri(outputLocation) .body(BodyInserters.fromResource(new InputStreamResource(pipedInputStream))) .retrieve() .bodyToMono(JsonNode.class) .block(); System.out.println("Response: " + response); }
但希望采用更贴合Reactor的响应式实现。尝试的test2方法中,先订阅上传请求后写入ZIP,程序会在上传完成前提前结束:
void test2() throws IOException { var outputLocation = "http://localhost:5000/anything"; var pipedInputStream = new PipedInputStream(); var pipedOutputStream = new PipedOutputStream(pipedInputStream); var response = webClient.post() .uri(outputLocation) .body(BodyInserters.fromResource(new InputStreamResource(pipedInputStream))) .retrieve() .bodyToMono(JsonNode.class) .subscribe((a) -> { System.out.println("Completed"); }); try { var zip = new ZipOutputStream(pipedOutputStream); for (var i = 0; i < 10; i++) { // 实际应为远程文件下载请求,返回可流式写入ZIP的InputStream zip.putNextEntry(new ZipEntry("test" + i + ".txt")); // 实际应为:downloadFile().transferTo(zip); new ByteArrayInputStream(("Hello " + i).getBytes()) .transferTo(zip); zip.closeEntry(); } zip.close(); } catch (Exception e) { e.printStackTrace(); } // 方法在"Completed"输出前就返回,如何阻塞直到上传完成? }
推测应创建Flux作为请求体生产者,但不知具体实现方式。
响应式流式实现方案
可以通过Reactor的Flux结合DataBuffer实现全流式的ZIP打包与上传,无需手动管理线程或管道流,完全贴合响应式编程模型:
核心思路
- 用
Flux生成ZIP格式的数据流,包含每个文件的条目头、流式下载的文件内容、条目尾,以及ZIP结束标记 - 远程文件下载采用WebClient的响应式流式获取(
bodyToFlux(DataBuffer.class)),避免加载整个文件到内存 - 将生成的ZIP数据流作为WebClient请求体,通过
BodyInserters.fromPublisher实现流式上传 - 用Reactor的调度机制管理整个流程的生命周期,避免异步导致的提前结束问题
完整代码示例
import org.springframework.core.io.buffer.DataBuffer; import org.springframework.core.io.buffer.DataBufferFactory; import org.springframework.core.io.buffer.DefaultDataBufferFactory; import org.springframework.web.reactive.function.BodyInserters; import org.springframework.web.reactive.function.client.WebClient; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import com.fasterxml.jackson.databind.JsonNode; import java.io.ByteArrayOutputStream; import java.io.IOException; import java.util.zip.ZipEntry; import java.util.zip.ZipOutputStream; void reactiveZipUpload() { String outputLocation = "http://localhost:5000/anything"; DataBufferFactory bufferFactory = new DefaultDataBufferFactory(); WebClient webClient = WebClient.create(); // 生成ZIP格式的Flux数据流 Flux<DataBuffer> zipFlux = Flux.range(0, 10) .flatMap(i -> { String entryName = "test" + i + ".txt"; // 1. 生成ZIP条目头的数据流 Mono<DataBuffer> entryHeaderMono = generateZipEntryHeader(bufferFactory, entryName); // 2. 流式下载远程文件的数据流(替换为实际远程地址) Flux<DataBuffer> fileContentFlux = webClient.get() .uri("http://remote-server/files/" + i) .retrieve() .bodyToFlux(DataBuffer.class); // 3. 生成ZIP条目尾的数据流 Mono<DataBuffer> entryFooterMono = generateZipEntryFooter(bufferFactory); // 拼接条目头、文件内容、条目尾 return Flux.concat(entryHeaderMono, fileContentFlux, entryFooterMono); }) // 最后添加ZIP结束标记 .concatWith(generateZipEndMarker(bufferFactory)) // 处理异常,确保资源释放 .onErrorResume(e -> { e.printStackTrace(); return Flux.empty(); }); // 流式上传ZIP数据并等待完成 JsonNode response = webClient.post() .uri(outputLocation) .contentType(org.springframework.http.MediaType.APPLICATION_OCTET_STREAM) .body(BodyInserters.fromPublisher(zipFlux, DataBuffer.class)) .retrieve() .bodyToMono(JsonNode.class) .block(); // 若为非阻塞场景,可改用subscribe并配合doOnComplete等操作 System.out.println("Response: " + response); } // 生成ZIP条目头的DataBuffer private Mono<DataBuffer> generateZipEntryHeader(DataBufferFactory bufferFactory, String entryName) { return Mono.fromCallable(() -> { ByteArrayOutputStream baos = new ByteArrayOutputStream(); try (ZipOutputStream zipOut = new ZipOutputStream(baos)) { zipOut.putNextEntry(new ZipEntry(entryName)); // 仅写入条目头,不关闭流,避免写入结束标记 zipOut.flush(); } catch (IOException e) { throw new RuntimeException(e); } return bufferFactory.wrap(baos.toByteArray()); }); } // 生成ZIP条目尾的DataBuffer private Mono<DataBuffer> generateZipEntryFooter(DataBufferFactory bufferFactory) { return Mono.fromCallable(() -> { ByteArrayOutputStream baos = new ByteArrayOutputStream(); try (ZipOutputStream zipOut = new ZipOutputStream(baos)) { zipOut.closeEntry(); zipOut.flush(); } catch (IOException e) { throw new RuntimeException(e); } return bufferFactory.wrap(baos.toByteArray()); }); } // 生成ZIP结束标记的DataBuffer private Mono<DataBuffer> generateZipEndMarker(DataBufferFactory bufferFactory) { return Mono.fromCallable(() -> { ByteArrayOutputStream baos = new ByteArrayOutputStream(); try (ZipOutputStream zipOut = new ZipOutputStream(baos)) { // 关闭ZipOutputStream会自动写入结束标记 zipOut.close(); } catch (IOException e) { throw new RuntimeException(e); } return bufferFactory.wrap(baos.toByteArray()); }); }
关键说明
- 全流式处理:从远程文件下载到ZIP打包,再到上传,所有环节都是流式操作,不会将大文件加载到内存,避免缓冲区溢出
- 响应式生命周期管理:整个流程由Reactor调度,
block()会等待上传完成后再返回,解决了test2中程序提前结束的问题;若为非阻塞场景,可替换为subscribe()并添加doOnComplete、doOnError等回调 - 资源安全:通过
DataBuffer的自动管理(WebClient会负责释放)和try-with-resources确保流资源正确关闭,避免内存泄漏 - 灵活性:可轻松扩展为并行下载多个文件(调整
flatMap的并发参数),或添加进度监控等逻辑
内容的提问来源于stack exchange,提问作者Eric
相关产品推荐
相关产品推荐

