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

使用ProcessPoolExecutor时aioboto3/aiobotocore异步调用挂起问题

问题分析与解决方案

核心问题

使用ProcessPoolExecutor并行运行包含异步S3调用的模型预测时,卡在key_exists函数的s3.head_object()调用处,单独运行正常,多进程环境下无限挂起。

原因拆解

  1. Fork进程的资源继承问题:用fork方式创建子进程时,父进程的异步IO资源(如socket、事件循环状态)会被复制到子进程,这些资源在子进程中处于无效状态,导致aioboto3客户端无法完成S3请求,进而挂起。
  2. 事件循环嵌套与初始化问题:主进程事件循环中嵌套调用子进程的asyncio.run,若子进程事件循环初始化不彻底,会引发异步任务阻塞。
  3. 异步客户端的跨进程状态污染:父进程中创建的异步客户端(如aioboto3 Session)被子进程继承,导致客户端状态异常。

修复方案

方案1:改用Spawn进程池+独立子进程事件循环

优先使用spawn方式创建进程(完全初始化新Python解释器,不继承父进程资源),并在每个子进程中独立初始化事件循环。

修改后的核心代码

import asyncio
import multiprocessing as mp

class ModelA(BaseModel):
    def predict(self, df):
        # 你的预测逻辑
        pass

    async def run_predictions(self):
        # 数据获取与预测执行逻辑
        df = self.get_data()  # 替换为你的数据获取方法
        self.predict(df)

class ModelB(BaseModel):
    async def predict(self, df):
        preds = []
        async with aiohttp.ClientSession() as session:
            # 并行处理各company的请求,提升异步效率
            tasks = [
                self._process_company(company, comp_df, session)
                for company, comp_df in df.groupby('company_id')
            ]
            preds.extend(await asyncio.gather(*tasks))
        # 后续预测逻辑处理
        return preds

    async def _process_company(self, company, comp_df, session):
        path = f'{some_path}/{company}.pqt'
        exists = await key_exists(WRITE_BUCKET, path)
        if exists:
            preds_df = await self.read_from_s3(path)
        else:
            # 用asyncio.to_thread包装同步OCR任务,避免阻塞事件循环
            preds_df = await asyncio.to_thread(run_ocr, comp_df)
        # 单company的预测逻辑
        return preds_df

    async def read_from_s3(self, path):
        # 子进程内独立创建S3客户端
        async with aioboto3.Session().client("s3") as s3_client:
            obj = await s3_client.get_object(Bucket=WRITE_BUCKET, Key=path)
            # 示例:读取Parquet文件
            import pyarrow.parquet as pq
            import io
            return pq.read_table(io.BytesIO(await obj['Body'].read())).to_pandas()

async def key_exists(bucket, key):
    async with aioboto3.Session().client("s3") as s3_client:
        try:
            LOGGER.info("entering get object in key_exists")
            await s3_client.head_object(Bucket=bucket, Key=key)
            LOGGER.info("exiting get object in key_exists")
            return True
        except Exception as e:
            LOGGER.debug(f"Key {key}不存在: {str(e)}")
            return False

def run_model_in_process(model):
    # 子进程独立初始化事件循环
    loop = asyncio.new_event_loop()
    asyncio.set_event_loop(loop)
    try:
        loop.run_until_complete(model.run_predictions())
    finally:
        loop.close()

def run_async_models(async_models):
    # 使用spawn方式创建进程池
    with mp.Pool(mp_context=mp.get_context('spawn')) as pool:
        pool.map(run_model_in_process, async_models)

# 调用入口
if __name__ == "__main__":
    async_models = [ModelA(), ModelB()]
    run_async_models(async_models)

关键修改点

  • 进程创建方式:用spawn替代fork,避免继承父进程异步资源。
  • 子进程事件循环:每个子进程独立创建并管理事件循环,彻底隔离异步状态。
  • 异步客户端隔离:所有aioboto3、aiohttp客户端均在子进程的异步上下文中创建,不依赖父进程资源。
  • 同步任务异步化:用asyncio.to_thread包装同步的run_ocr,避免阻塞事件循环。

方案2:调整ProcessPoolExecutor调用逻辑(若需保留原结构)

如果必须使用ProcessPoolExecutor,需确保子进程事件循环正确初始化:

async def run_async_models(async_models):
    with ProcessPoolExecutor(mp_context=mp.get_context('spawn')) as pool:
        event_loop = asyncio.get_running_loop()
        master_tasks = [
            event_loop.run_in_executor(
                pool,
                get_model_tasks,
                model,
            )
            for model in async_models
        ]
        await asyncio.gather(*master_tasks)

def get_model_tasks(model):
    # 子进程内重新初始化事件循环
    loop = asyncio.new_event_loop()
    asyncio.set_event_loop(loop)
    result = loop.run_until_complete(model.run_predictions())
    loop.close()
    return result

额外注意事项

  • 模型序列化:确保ModelA、ModelB对象可被pickle序列化,若包含不可序列化属性,需实现__getstate__和__setstate__方法处理。
  • 异步任务并行化:在ModelB中用asyncio.gather并行处理多个company的请求,充分发挥异步优势。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 23:48:19