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

使用aioboto3与aiofiles从S3下载大文件速度极慢,求排查与优化方案

问题排查与优化方案

核心慢因分析

从提供的代码和性能数据来看,1GB文件耗时300秒(约5MB/s)远低于S3正常下载速率,主要问题集中在以下几点:

1. 重复创建S3客户端

每次调用fetch_file都会新建S3客户端,客户端初始化涉及TCP连接建立、身份验证等开销,高负载下重复创建会大幅增加耗时。

2. 小分块串行下载+写入

默认情况下,response["Body"]的分块尺寸很小(通常几KB到几十KB),串行遍历分块并逐块写入磁盘,会导致大量IO系统调用和网络往返,拖慢整体速度。

3. 未利用并行分段下载

单连接下载大文件受限于单条TCP流的带宽上限,无法充分利用网络资源。

针对性优化方案

方案1:复用S3客户端

将客户端实例化逻辑移到初始化阶段,避免重复创建,同时开启连接复用和HTTP2提升并发:

class S3Storage:
    def __init__(self, *args, **kwargs):
        ...
        self._config = botocore.config.Config(
            read_timeout=read_timeout,
            connect_timeout=connect_timeout,
            retries={
                "total_max_attempts": ...,
                "max_attempts": ...,
            },
            signature_version="v4",
            tcp_keepalive=True,  # 开启TCP长连接复用
            http2=True  # 启用HTTP2提升并发效率
        )
        self._session = self._create_session()
        # 初始化客户端并复用
        self._client = self._create_client()

    def _create_session():
        return aioboto3.Session()

    def _create_client(self):
        return self._session.client(
            service_name="s3",
            endpoint_url=self._endpoint_url,
            region_name=self._region,
            aws_access_key_id=self._access_key_id,
            aws_secret_access_key=self._secret_access_key,
            config=self._config,
        )

    async def fetch_file(self, key: str, bucket: str | None = None) -> str:
        try:
            # 复用已创建的客户端
            response = await self._client.get_object(Bucket=bucket or self._bucket, Key=key)
            # ... 后续写入逻辑
        except Exception as e:
            # 异常处理:客户端失效时可重新创建
            self._client = self._create_client()
            ...

方案2:增大分块尺寸并批量写入

修改分块读取逻辑,使用更大的缓冲区,减少IO调用次数:

async def fetch_file(self, key: str, bucket: str | None = None) -> str:
    try:
        response = await self._client.get_object(Bucket=bucket or self._bucket, Key=key)
        async with aiofiles.tempfile.NamedTemporaryFile(
            "wb",
            suffix=_get_file_extension(key),
            delete=False,
            buffering=1024*128  # 设置128KB缓冲区
        ) as file:
            # 自定义1MB分块读取,减少网络IO次数
            chunk_size = 1024 * 1024
            while True:
                chunk = await response["Body"].read(chunk_size)
                if not chunk:
                    break
                await file.write(chunk)
            return str(file.name)
    # ... 异常处理

方案3:并行分段下载(大文件最优解)

利用S3的Range请求,将大文件分成多个片段并行下载,最后合并,充分利用带宽:

import asyncio

class S3Storage:
    # ... 其他初始化代码

    async def _download_chunk(self, bucket: str, key: str, start: int, end: int, temp_path: str):
        """下载单个文件片段"""
        response = await self._client.get_object(
            Bucket=bucket,
            Key=key,
            Range=f"bytes={start}-{end}"
        )
        async with aiofiles.open(temp_path, "rb+") as f:
            await f.seek(start)
            chunk_size = 1024 * 1024
            while True:
                chunk = await response["Body"].read(chunk_size)
                if not chunk:
                    break
                await f.write(chunk)

    async def fetch_file(self, key: str, bucket: str | None = None) -> str:
        bucket = bucket or self._bucket
        try:
            # 获取文件总大小
            head_response = await self._client.head_object(Bucket=bucket, Key=key)
            file_size = head_response["ContentLength"]
            
            # 分块大小设置为50MB(可根据网络调整,建议10-100MB)
            chunk_size = 50 * 1024 * 1024
            async with aiofiles.tempfile.NamedTemporaryFile(
                "wb",
                suffix=_get_file_extension(key),
                delete=False
            ) as temp_file:
                # 预分配文件空间,提升写入效率
                await temp_file.truncate(file_size)
                temp_file_path = temp_file.name

            # 生成所有分段任务
            chunk_tasks = []
            for i in range(0, file_size, chunk_size):
                start = i
                end = min(i + chunk_size - 1, file_size - 1)
                chunk_tasks.append(self._download_chunk(bucket, key, start, end, temp_file_path))
            
            # 并行下载所有分段
            await asyncio.gather(*chunk_tasks)
            return temp_file_path
        # ... 异常处理

额外优化建议

  • 调整超时配置:将read_timeout设置为300秒以上,避免大文件下载时触发超时;connect_timeout设置为5-10秒即可。
  • 简化重试策略:如果网络稳定,可适当降低重试次数,避免不必要的重试等待。
  • 检查存储位置:确保文件存储在与计算节点同区域的S3桶,跨区域会增加延迟;避免使用低频访问存储类(如Glacier),这类存储读取速度远低于标准存储。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 00:00:16