Python异步批量上传文件遇OSError: [Errno 24] Too many open files问题排查
Python异步批量上传文件报错"Too many open files"分析与解决
你的异步代码工作逻辑
你的代码执行流程如下:
- 从
event_loop启动异步任务,将所有文件按每100个分组为批次 - 同时启动所有批次的处理任务(
await asyncio.gather(*batch_tasks)) - 每个批次内又同时启动该批次所有文件的上传任务,每个文件并行执行SFTP和S3上传
当文件量达到1000个时,相当于同时创建1000个SSH连接+1000个S3连接,再加上打开的本地文件,每个连接和文件都会占用一个系统文件描述符。而系统默认的文件描述符上限一般为1024,直接超出限制触发Errno 24错误。
问题核心
- 连接完全未复用:每个文件上传都新建SSH/S3连接,单个连接本可处理多个上传任务,重复建连接纯粹浪费资源。
- 并发无限制:所谓的"批次"仅做文件分组,但所有批次任务同时执行,本质是1000个任务全量并发,瞬间耗尽系统资源。
- 高并发下资源创建速度远快于释放速度,进一步加剧描述符耗尽问题。
修复步骤
1. 复用SFTP连接
提前初始化一次SFTP连接,所有上传任务共用该连接:
class my_class(): def __init__(self): self.sftp_conn = None self.sftp_client = None async def init_sftp(self): # 提前建立一次SFTP连接 self.sftp_conn = await asyncssh.connect( host="hostname", username="uname", password="pwd", known_hosts=None, ) self.sftp_client = await self.sftp_conn.start_sftp_client() async def sftp_uploader(self, source_file, destination_file): # 复用已建立的客户端上传 await self.sftp_client.put(source_file, destination_file) def event_loop(self): tuple_of_records = () # 你的文件列表 loop = asyncio.get_event_loop() # 先初始化SFTP连接 loop.run_until_complete(self.init_sftp()) # 执行上传任务 loop.run_until_complete(self.iterate_asynchronously(tuple_of_records)) # 任务完成后关闭连接 loop.run_until_complete(self.sftp_client.close()) loop.run_until_complete(self.sftp_conn.close()) loop.close()
2. 复用S3客户端
aiobotocore客户端自带连接池,提前初始化一次即可:
class my_class(): def __init__(self): self.sftp_conn = None self.sftp_client = None self.s3_client = None async def init_s3(self): ses = session.get_session() # 提前初始化S3客户端 self.s3_client = await ses.create_client( 's3', aws_secret_access_key='key', aws_access_key_id='id' ) async def s3_uploader(self, source_file, destination_file): with open(source_file, 'rb') as file: await self.s3_client.put_object( Bucket='bucket_name', Key=destination_file, Body=file ) def event_loop(self): tuple_of_records = () loop = asyncio.get_event_loop() # 同时初始化SFTP和S3连接 loop.run_until_complete(asyncio.gather(self.init_sftp(), self.init_s3())) loop.run_until_complete(self.iterate_asynchronously(tuple_of_records)) # 关闭所有资源 loop.run_until_complete(self.sftp_client.close()) loop.run_until_complete(self.sftp_conn.close()) loop.run_until_complete(self.s3_client.close()) loop.close()
3. 限制并发量
用asyncio.Semaphore控制同时运行的上传任务数量,比如限制最多20个并发:
class my_class(): def __init__(self): self.semaphore = asyncio.Semaphore(20) # 并发数上限 # 其他初始化... async def my_method(self, files): # 超过并发上限时自动等待 async with self.semaphore: source_file, destination_file = files # 假设files包含源文件和目标路径 await asyncio.gather( self.sftp_uploader(source_file, destination_file), self.s3_uploader(source_file, destination_file) )
4. 可选:串行处理批次
如果不想用信号量,可改为串行处理批次(效率略低但实现简单):
async def iterate_asynchronously(self, tuple_of_records): batch_size = 100 batches = [tuple_of_records[i:i+batch_size] for i in range(0, len(tuple_of_records), batch_size)] # 处理完一批再启动下一批 for batch in batches: await self.batch_process(batch)
额外提示
- 用
ulimit -n查看当前系统文件描述符上限,临时调高(如ulimit -n 4096)可缓解,但核心还是优化代码。 - 确保所有异步资源通过
async with或手动关闭,避免资源泄漏。
内容的提问来源于stack exchange,提问作者Tim
相关产品推荐
相关产品推荐

