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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 13:35:24