.NET 8实现S3多文件压缩流式回传至S3(内存限100MB)
.NET 8 S3 文件流式压缩上传(内存控制≤100MB)
要实现边下载S3文件、边压缩、边流式上传到目标桶,同时严格控制内存占用,核心是利用**管道(Pipe)**实现带背压的缓冲机制,结合ZipArchive流式写入和S3流式上传能力。以下是完整实现方案:
核心思路
- 使用
System.IO.Pipelines.Pipe创建带内存阈值限制的缓冲区,当缓冲区达到100MB时自动暂停写入(即暂停下载压缩),直到数据被上传端读取释放空间。 - 启动两个异步任务:
- 任务1:遍历源S3桶文件,逐个下载并写入
ZipArchive,ZipArchive直接输出到管道写入端。 - 任务2:从管道读取端获取压缩数据,流式上传到目标S3桶。
- 任务1:遍历源S3桶文件,逐个下载并写入
- 利用管道的背压机制自动平衡下载压缩速度和上传速度,确保内存占用不超限。
完整代码实现
必要命名空间
using Amazon.S3; using Amazon.S3.Model; using System.IO.Compression; using System.IO.Pipelines;
主逻辑
var sourceBucket = "your-source-bucket"; var targetBucket = "your-target-bucket"; var archiveKey = "compressed-archive.zip"; var maxMemoryLimit = 100 * 1024 * 1024; // 100MB // 初始化S3客户端(建议使用依赖注入管理生命周期) using var s3Client = new AmazonS3Client(); // 配置管道:设置背压阈值,缓冲区满100MB时暂停写入,降到80MB时恢复 var pipeOptions = new PipeOptions( pauseWriterThreshold: maxMemoryLimit, resumeWriterThreshold: (int)(maxMemoryLimit * 0.8), minimumSegmentSize: 4096 ); var pipe = new Pipe(pipeOptions); // 并行执行压缩写入和上传任务 var compressionTask = CompressS3FilesToPipe(s3Client, sourceBucket, pipe.Writer); var uploadTask = UploadPipeToS3(s3Client, targetBucket, archiveKey, pipe.Reader); await Task.WhenAll(compressionTask, uploadTask); Console.WriteLine("归档上传完成");
压缩写入管道方法
async Task CompressS3FilesToPipe(IAmazonS3 s3Client, string sourceBucket, PipeWriter writer) { try { // 分页遍历源桶所有对象 var listRequest = new ListObjectsV2Request { BucketName = sourceBucket }; ListObjectsV2Response listResponse; do { listResponse = await s3Client.ListObjectsV2Async(listRequest); foreach (var s3Obj in listResponse.S3Objects) { // 下载当前S3文件 using var getResponse = await s3Client.GetObjectAsync(sourceBucket, s3Obj.Key); using var fileStream = getResponse.ResponseStream; // 将PipeWriter包装为流,供ZipArchive写入 using var zipArchive = new ZipArchive(writer.AsStream(), ZipArchiveMode.Create, leaveOpen: true); var zipEntry = zipArchive.CreateEntry(s3Obj.Key, CompressionLevel.Optimal); // 将下载的文件流复制到Zip条目 using var entryStream = zipEntry.Open(); await fileStream.CopyToAsync(entryStream); // 刷新管道,确保数据推送到读取端 await writer.FlushAsync(); } listRequest.ContinuationToken = listResponse.NextContinuationToken; } while (listResponse.IsTruncated); // 完成写入,通知读取端无更多数据 await writer.CompleteAsync(); } catch (Exception ex) { // 通知读取端写入出错 await writer.CompleteAsync(ex); throw; } }
管道数据上传到S3方法
async Task UploadPipeToS3(IAmazonS3 s3Client, string targetBucket, string archiveKey, PipeReader reader) { try { var putRequest = new PutObjectRequest { BucketName = targetBucket, Key = archiveKey, InputStream = reader.AsStream(), // 将PipeReader包装为可读流 AutoCloseStream = false // 禁止自动关闭,由管道自身管理 }; // 流式上传到S3,AWS SDK会自动处理分块(默认8MB分块) await s3Client.PutObjectAsync(putRequest); // 完成读取 await reader.CompleteAsync(); } catch (Exception ex) { // 通知写入端读取出错 await reader.CompleteAsync(ex); throw; } }
关键细节说明
- 内存控制:管道的
pauseWriterThreshold设置为100MB,当缓冲区数据量达到该值时,后续的写入操作会阻塞,直到上传端读取数据使缓冲区降到resumeWriterThreshold(80MB),自动恢复写入。 - 流式压缩:
ZipArchive直接写入管道流,无需将整个归档文件存入内存,每个文件压缩后立即进入管道缓冲区。 - 流式上传:S3的
PutObjectAsync支持流式输入,管道的读取流会持续提供压缩数据,实现边读边上传,无需等待整个归档生成。 - 资源管理:所有流和客户端均使用
using语句确保自动释放,避免资源泄漏。
内容的提问来源于stack exchange,提问作者Maxime Rossini
相关产品推荐
相关产品推荐

