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

aiomultiprocessing Pool冻结及OSError: Too many open files问题求助

问题:多进程异步任务中出现"Too many open files"错误

在32核机器上执行约700万次API请求并将数据导入PostgreSQL,程序运行1-2小时后无预警冻结,手动中断时触发OSError: [Errno 24] Too many open files错误。

核心代码结构

from aiomultiprocess import Pool
CPUS = 30 # 机器共32核
batch_size = 5000

for i in range(0, len(args_list), batch_size):
        log.info(f"<log message>")
        async with Pool(
            processes=CPUS,
            maxtasksperchild=100,
            childconcurrency=3,
            queuecount=int(CPUS / 3),
        ) as pool:
            await pool.starmap(fetch_api_results, args_list[i : i + batch_size])

补充:fetch_api_results及分页逻辑代码

fetch_api_results负责构造API请求,通过aiohttp递归请求直到无next_url:

from aiohttp import request

async def fetch_api_results(*args):
    try:
        result_objects = APIPaginator(*args)
        await result_objects.fetch()
        log.info("uploading data")
        # 数据导入PostgreSQL的函数调用

    except planned_exceptions as e:
        log.warning(e, exc_info=False)


class APIPaginator(object):
    async def query_data(self):
        url = self.api_base + "<str from arg>"
        payload = {"limit": 1000}
        await self.query_all(url, payload)

    async def query_all(self, url, payload):
        try:
            async with request(method="GET", url=url, params=payload) as response:
                log.info(f"status code: {response.status}")
                if response.status == 200:
                    results = await response.json()
                    self.results.append(results)
                    next_url = results.get("next_url")
                    if next_url:
                        await self.query_all(next_url)
                else:
                    response.raise_for_status()
        except:
            # 异常处理逻辑省略
            pass

    async def fetch(self):
        await self.query_data()

完整错误回溯

File "<path-to-file>.py", line 127, in import_data
    await pool.starmap(fetch_api_results, args_list[i : i + batch_size])
  File "/<path-to-env>/lib/python3.11/site-packages/aiomultiprocess/pool.py", line 136, in results
    return await self.pool.results(self.task_ids)
           ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/<path-to-env>/lib/python3.11/site-packages/aiomultiprocess/pool.py", line 312, in results
    await asyncio.sleep(0.005)
  File "/<path-to-env>/lib/python3.11/asyncio/tasks.py", line 639, in sleep
    return await future

File "/<path-to-env>/3.11.1/lib/python3.11/asyncio/runners.py", line 118, in run
    return self._loop.run_until_complete(task)
           ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/<path-to-env>/lib/python3.11/asyncio/base_events.py", line 653, in run_until_complete
    return future.result()
           ^^^^^^^^^^^^^^^
  File "/<path-to-file>/main.py", line 39, in add_all_data
    await import_data(args)
  File "/<path-to-file>/orchestrator.py", line 120, in import_data
    async with Pool(
  File "/<path-to-env>/lib/python3.11/site-packages/aiomultiprocess/pool.py", line 196, in __aexit__
    await self.join()
  File "/<path-to-env>/lib/python3.11/site-packages/aiomultiprocess/pool.py", line 379, in join
    await self._loop
  File "/<path-to-env>/lib/python3.11/site-packages/aiomultiprocess/pool.py", line 229, in loop
    self.processes[self.create_worker(qid)] = qid
                   ^^^^^^^^^^^^^^^^^^^^^^^
  File "/<path-to-env>/lib/python3.11/site-packages/aiomultiprocess/pool.py", line 261, in create_worker
    process.start()
  File "/<path-to-env>/lib/python3.11/site-packages/aiomultiprocess/core.py", line 153, in start
    return self.aio_process.start()
           ^^^^^^^^^^^^^^^^^^^^^^^^
  File "/<path-to-env>/lib/python3.11/multiprocessing/process.py", line 121, in start
    self._popen = self._Popen(self)
                  ^^^^^^^^^^^^^^^^^
  File "/<path-to-env>/lib/python3.11/multiprocessing/context.py", line 288, in _Popen
    return Popen(process_obj)
           ^^^^^^^^^^^^^^^^^^
  File "/<path-to-env>/lib/python3.11/multiprocessing/popen_spawn_posix.py", line 32, in __init__
    super().__init__(process_obj)
  File "/<path-to-env>/lib/python3.11/multiprocessing/popen_fork.py", line 19, in __init__
    self._launch(process_obj)
  File "/home/<path-to-env>/lib/python3.11/multiprocessing/popen_spawn_posix.py", line 58, in _launch
    self.pid = util.spawnv_passfds(spawn.get_executable(),
               ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/<path-to-env>/lib/python3.11/multiprocessing/util.py", line 451, in spawnv_passfds
    errpipe_read, errpipe_write = os.pipe()
                                  ^^^^^^^^^
OSError: [Errno 24] Too many open files

已尝试的方案与困惑

  • 调整maxtasksperchild参数,按文档说明该参数会在任务数达标后销毁旧进程并创建新进程,理论上可避免句柄泄漏,但实际未生效
  • 实现批处理逻辑:每处理5000个任务后通过async with Pool销毁并重建进程池,未解决问题
  • 临时方案:通过ulimit -n提高系统最大打开文件数,但担心长期运行仍会耗尽上限

请求解决方案或建议

寻求针对该文件句柄泄漏问题的根本解决方案或优化建议。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 13:24:56