如何实现适配asyncio的Python CPU密集型任务多进程装饰器
报错原因
你遇到的pickle报错是Python multiprocessing 模块的固有特性导致的:默认的pickle序列化器在序列化函数对象时,只会保存函数的模块名和函数限定名,不会序列化函数的实际代码,反序列化时会到对应模块的命名空间下查找同名对象。
你用装饰器修饰fact后,全局命名空间里的fact已经变成了装饰器返回的异步wrapper函数,闭包内保留的原fact函数没有绑定到__main__.fact这个名字上,pickle序列化原函数后,子进程反序列化时找不到对应对象,就触发了这个错误。
可行实现方案
可以通过「传递函数标识而非直接传递函数对象」的方式规避序列化限制,完全兼容你要求的使用方式,代码实现如下:
import functools import asyncio import importlib from typing import Callable, Awaitable, TypeVar, ParamSpec P = ParamSpec('P') R = TypeVar('R') # 注意这个辅助函数必须定义在模块顶层,保证可被pickle def _run_cpu_bound_func(func_qualname: str, module_name: str, args, kwargs): module = importlib.import_module(module_name) func = getattr(module, func_qualname) return func(*args, **kwargs) def cpu_bound(func: Callable[P, R]) -> Callable[P, Awaitable[R]]: func_qualname = func.__qualname__ module_name = func.__module__ @functools.wraps(func) async def wrapper(*args: P.args, **kwargs: P.kwargs) -> R: loop = asyncio.get_running_loop() executor = get_executor() # 你的自定义获取执行器逻辑 return await loop.run_in_executor( executor, _run_cpu_bound_func, func_qualname, module_name, args, kwargs ) return wrapper
注意事项
- 该方案要求被
cpu_bound装饰的函数必须是模块顶层定义的函数,不能是嵌套在其他函数/类内部的函数,这符合绝大多数CPU密集型工具函数的使用场景 - 传入被装饰函数的参数和返回值都必须是可被pickle序列化的类型,这是多进程通信的固有要求
- 如果允许引入第三方依赖,也可以使用
loky作为进程池实现,它自带的序列化器支持更多类型的函数序列化,不需要修改装饰器逻辑。
内容的提问来源于stack exchange,提问作者BlueGlassBlock
相关产品推荐
相关产品推荐

