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

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错误。

问题核心

  1. 连接完全未复用:每个文件上传都新建SSH/S3连接,单个连接本可处理多个上传任务,重复建连接纯粹浪费资源。
  2. 并发无限制:所谓的"批次"仅做文件分组,但所有批次任务同时执行,本质是1000个任务全量并发,瞬间耗尽系统资源。
  3. 高并发下资源创建速度远快于释放速度,进一步加剧描述符耗尽问题。

修复步骤

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 01:59:52