使用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
相关产品推荐
相关产品推荐

