如何用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_
相关产品推荐
相关产品推荐

