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

.NET 8实现S3多文件压缩流式回传至S3(内存限100MB)

.NET 8 S3 文件流式压缩上传(内存控制≤100MB)

要实现边下载S3文件、边压缩、边流式上传到目标桶,同时严格控制内存占用,核心是利用**管道(Pipe)**实现带背压的缓冲机制,结合ZipArchive流式写入和S3流式上传能力。以下是完整实现方案:

核心思路

  1. 使用System.IO.Pipelines.Pipe创建带内存阈值限制的缓冲区,当缓冲区达到100MB时自动暂停写入(即暂停下载压缩),直到数据被上传端读取释放空间。
  2. 启动两个异步任务:
    • 任务1:遍历源S3桶文件,逐个下载并写入ZipArchive,ZipArchive直接输出到管道写入端。
    • 任务2:从管道读取端获取压缩数据,流式上传到目标S3桶。
  3. 利用管道的背压机制自动平衡下载压缩速度和上传速度,确保内存占用不超限。

完整代码实现

必要命名空间

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 06:42:17