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

主async应用中使用ProcessPoolExecutor运行多进程async循环的实现问题

报错原因

ProcessPoolExecutor 依赖进程间IPC通信传递数据,所有跨进程传递的函数、参数都必须支持pickle序列化。原代码触发错误的核心问题有两个:

  • 传递给进程池的是绑定类实例的方法 self.blocking_task,类实例中包含不可序列化的 event_loop、pool_executor 属性,序列化整个实例时直接失败
  • 类内的异步方法绑定了主进程的运行上下文,子进程无法读取主进程的异步对象上下文
适配ProcessPoolExecutor的完整代码
import asyncio
import time
from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor

async def subtask(letter: str, multiplier: int):
    await asyncio.sleep(1)
    return letter * multiplier

async def task_gatherer(subtasks: list):
    return await asyncio.gather(*subtasks)

def blocking_task(word: str, multiplier: int):
    time.sleep(1)
    subtasks = [subtask(letter, multiplier) for letter in word]
    result = asyncio.run(task_gatherer(subtasks))
    return result


class MyClass:
    def __init__(self) -> None:
        self.event_loop = None
        self.pool_executor = ProcessPoolExecutor(max_workers=8)
        self.words = ["one", "two", "three", "four", "five"]
        self.multiplier = int(2)

    async def master_method(self):
        self.event_loop = asyncio.get_running_loop()
        master_tasks = [
            self.event_loop.run_in_executor(
                self.pool_executor,
                blocking_task,
                word,
                self.multiplier
            )
            for word in self.words
        ]

        results = await asyncio.gather(*master_tasks)
        print(results)
        self.pool_executor.shutdown()


if __name__ == "__main__":
    my_class = MyClass()
    asyncio.run(my_class.master_method())
关键修改说明
  • 将子进程需要运行的逻辑提取为顶层函数,不再绑定类实例,避免跨进程传递整个类实例
  • 仅将子进程需要的基础参数(word、multiplier)传递到进程池,这类基础类型天然支持pickle序列化
  • 子进程内部的异步逻辑完全独立,仅依赖传入的参数运行,不读取主进程的任何上下文对象

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 22:45:07