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

如何用WebFlux实现Spring MVC流式Zip文件下载的等效逻辑?

WebFlux 流式传输Zip文件的实现方案

完全可以实现和Spring MVC等效的流式Zip文件传输逻辑,核心思路是利用ServerHttpResponse直接写入ZipOutputStream,结合S3的异步响应式API逐段读取文件内容并写入Zip,全程不加载全量数据到内存。

实现步骤与代码示例

1. 核心选型说明

WebFlux中不适合返回Flux<SomeClass>(因为Zip是连续二进制流,无离散对象边界),正确方式是返回Mono<Void>,通过ServerHttpResponse获取输出流并包装为ZipOutputStream,逐个处理S3文档并写入Zip条目。

2. 依赖准备

确保引入AWS SDK v2的异步S3客户端(适配WebFlux响应式模型):

<dependency>
    <groupId>software.amazon.awssdk</groupId>
    <artifactId>s3</artifactId>
    <version>2.25.0</version>
</dependency>
<dependency>
    <groupId>software.amazon.awssdk</groupId>
    <artifactId>netty-nio-client</artifactId>
    <version>2.25.0</version>
</dependency>

3. 完整端点实现

import org.springframework.http.HttpHeaders;
import org.springframework.http.MediaType;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RestController;
import org.springframework.web.server.ServerWebExchange;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import software.amazon.awssdk.core.async.AsyncResponseTransformer;
import software.amazon.awssdk.services.s3.S3AsyncClient;
import software.amazon.awssdk.services.s3.model.GetObjectRequest;

import java.io.OutputStream;
import java.util.List;
import java.util.zip.ZipEntry;
import java.util.zip.ZipOutputStream;

@RestController
public class StreamingZipController {

    private final S3AsyncClient s3AsyncClient;

    public StreamingZipController(S3AsyncClient s3AsyncClient) {
        this.s3AsyncClient = s3AsyncClient;
    }

    @GetMapping("/download-zip")
    public Mono<Void> downloadStreamingZip(ServerWebExchange exchange) {
        // 设置响应头,让浏览器识别为下载文件
        exchange.getResponse().getHeaders().set(HttpHeaders.CONTENT_DISPOSITION,
            "attachment; filename=\"documents.zip\"");
        exchange.getResponse().getHeaders().setContentType(MediaType.APPLICATION_OCTET_STREAM);

        // 获取响应输出流并包装为ZipOutputStream
        return exchange.getResponse().writeWith(Mono.fromSupplier(() -> {
            OutputStream outputStream = exchange.getResponse().getOutputStream();
            return new ZipOutputStream(outputStream);
        }).flatMapMany(zipOutputStream -> {
            // 从业务逻辑获取S3文档列表(bucket + key + 显示文件名)
            List<S3Document> s3Documents = getS3DocumentList();

            // 遍历文档逐个写入Zip
            return Flux.fromIterable(s3Documents)
                .flatMap(s3Doc -> {
                    // 创建Zip条目
                    ZipEntry zipEntry = new ZipEntry(s3Doc.getFileName());
                    try {
                        zipOutputStream.putNextEntry(zipEntry);
                    } catch (Exception e) {
                        return Mono.error(e);
                    }

                    // 异步读取S3对象内容并写入Zip流
                    GetObjectRequest getObjectRequest = GetObjectRequest.builder()
                        .bucket(s3Doc.getBucket())
                        .key(s3Doc.getKey())
                        .build();

                    return Mono.fromCompletionStage(s3AsyncClient.getObject(getObjectRequest,
                            AsyncResponseTransformer.toOutputStream(zipOutputStream)))
                        .doOnSuccess(response -> {
                            try {
                                zipOutputStream.closeEntry();
                            } catch (Exception e) {
                                throw new RuntimeException(e);
                            }
                        });
                })
                .doOnComplete(() -> {
                    try {
                        zipOutputStream.finish();
                        zipOutputStream.close();
                    } catch (Exception e) {
                        // 处理流关闭异常
                    }
                });
        }));
    }

    // 模拟获取S3文档列表的业务方法
    private List<S3Document> getS3DocumentList() {
        return List.of(
            new S3Document("my-bucket", "docs/file1.pdf", "file1.pdf"),
            new S3Document("my-bucket", "docs/file2.txt", "file2.txt")
        );
    }

    // 封装S3文档信息的内部类
    private static class S3Document {
        private String bucket;
        private String key;
        private String fileName;

        public S3Document(String bucket, String key, String fileName) {
            this.bucket = bucket;
            this.key = key;
            this.fileName = fileName;
        }

        public String getBucket() { return bucket; }
        public String getKey() { return key; }
        public String getFileName() { return fileName; }
    }
}

4. 关键细节说明

  • 响应式适配:用AWS SDK v2的S3AsyncClient,返回的CompletionStage通过Mono.fromCompletionStage转为响应式类型,适配WebFlux线程模型。
  • 内存优化:每个S3文件的内容直接写入Zip流,不会把整个文件或Zip包加载到内存,和Spring MVC的StreamingResponseBody逻辑一致。
  • 资源管理:通过doOnComplete确保所有条目写入后完成并关闭Zip流,避免资源泄漏。

为什么不返回Flux?

Flux适用于返回离散的独立对象序列(如JSON数组),而Zip是连续二进制流,无明确对象边界,直接操作输出流的方式更符合场景需求,也能保证流式传输效率。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 22:43:11