解决Python3中ProcessPoolExecutor.submit传递实例方法的Pickle错误
多进程调用实例方法的序列化问题
我对多进程完全是新手!
我的需求
- 将耗时的实例方法**wait_n_secs()**作为独立进程运行,让其他任务可并行执行。
- 实例方法执行完成后,通过multiprocessing模块的共享数组获取其输出并使用。
尝试运行的代码
import cv2 import time from multiprocessing import Array import concurrent.futures import copyreg as copy_reg import types def _pickle_method(m): if m.im_self is None: return getattr, (m.im_class, m.im_func.func_name) else: return getattr, (m.im_self, m.im_func.func_name) copy_reg.pickle(types.MethodType, _pickle_method) class Testing(): def __init__(self): self.executor = concurrent.futures.ProcessPoolExecutor() self.futures = None self.shared_array = Array('i', 4) def wait_n_secs(self,n): print(f"I wait for {n} sec") cv2.waitKey(n*1000) wait_array = (n,n,n,n) return wait_array def function(waittime): bbox = Testing().wait_n_secs(waittime) return bbox if __name__ =="__main__": testing = Testing() waittime = 5 # 无法运行! testing.futures = testing.executor.submit(testing.wait_n_secs,waittime) # 可以运行! #testing.futures = testing.executor.submit(function,waittime) stime = time.time() while 1: if not testing.futures.running(): print("Checking for results") testing.shared_array = testing.futures.result() print("Shared_array received = ",testing.shared_array) break time_elapsed = time.time()-stime if (( time_elapsed % 1 ) < 0.001): print(f"Time elapsed since some time = {time_elapsed:.2f} sec")
遇到的问题
1) Python 3.6中的错误:
Traceback (most recent call last): File "C:\Users\haide\AppData\Local\Programs\Python\Python36\lib\multiprocessing\queues.py", line 234, in _feed obj = _ForkingPickler.dumps(obj) File "C:\Users\haide\AppData\Local\Programs\Python\Python36\lib\multiprocessing\reduction.py", line 51, in dumps cls(buf, protocol).dump(obj) File "C:\Users\haide\AppData\Local\Programs\Python\Python36\lib\multiprocessing\queues.py", line 58, in __getstate__ context.assert_spawning(self) File "C:\Users\haide\AppData\Local\Programs\Python\Python36\lib\multiprocessing\context.py", line 356, in assert_spawning ' through inheritance' % type(obj).__name__ RuntimeError: Queue objects should only be shared between processes through inheritance
2) Python 3.8中的错误:
testing.shared_array = testing.futures.result() File "C:\Users\haide\AppData\Local\Programs\Python\Python38\lib\concurrent\futures\_base.py", line 437, in result return self.__get_result() File "C:\Users\haide\AppData\Local\Programs\Python\Python38\lib\concurrent\futures\_base.py", line 389, in __get_result raise self._exception File "C:\Users\haide\AppData\Local\Programs\Python\Python38\lib\multiprocessing\queues.py", line 239, in _feed obj = _ForkingPickler.dumps(obj) File "C:\Users\haide\AppData\Local\Programs\Python\Python38\lib\multiprocessing\reduction.py", line 51, in dumps cls(buf, protocol).dump(obj) TypeError: cannot pickle 'weakref' object
核心问题是多进程中实例方法无法被Pickle序列化,导致报错,和其他用户遇到的情况一致。
尝试的部分解决方案
最常提及的方案是使用copyreg来序列化实例方法:
- 不太理解copyreg的工作原理,尝试添加相关代码到文件顶部,但未成功。
- 重要说明: 使用Python 3,导入的是copyreg,而相关解决方案多基于Python 2的copy_reg。
未尝试的方案:
- 使用dill库,因为相关案例要么不涉及多进程,要么未使用concurrent.futures模块。
临时 workaround
- 将调用实例方法的普通函数传递给submit(),而非直接传递实例方法。
testing.futures = testing.executor.submit(function,waittime)
该方法可行,但不够优雅。
我的诉求
- 请指导如何正确使用copyreg,因为确实不理解它的工作方式。
或者 - 若这是Python 3的问题,请提供另一种可向
concurrent.futures.ProcessPoolExecutor.submit()传递实例方法的优雅解决方案。 :)
更新#1
请求分享“传递以实例为参数的模块级函数”的示例代码,或者纠正我的尝试错误:
这是我的尝试。:(
- 将实例与参数一起传递给模块级函数
inp_args = [waittime] testing.futures = testing.executor.submit(wrapper_func,testing,inp_args)
- 创建的模块级包装函数:
def wrapper_func(ins,*args): ins.wait_n_secs(args)
- 这又导致了错误:
TypeError: cannot pickle 'weakref' object
内容的提问来源于stack exchange,提问作者haider abbasi
相关产品推荐
相关产品推荐

