使用multiprocessing pool.imap_unordered时遇无法Pickle函数问题求助
问题:multiprocessing中pickle np.vectorize包装的函数失败
复现代码
import numpy as np import pandas as pd import multiprocessing as mp # Define test function def test_func(idx, a, b): return [idx, a**1 + b**2] test_func = np.vectorize(test_func) test_df = pd.DataFrame(np.random.normal(size = (10, 2)), columns = ['a', 'b']) test_df = test_df.reset_index() func_dict = {col:test_df[col] for col in test_df.columns} # Function used to evaluate test_func def expandCall(kargs): # Get the function argument func = kargs['func'] # Delete it from the fictionary del kargs['func'] # Evaluate function with other arguments out = func(**kargs) return out parts = np.linspace(0, test_df.shape[0], 2).astype(int) jobs = [] for i, j in zip(parts[:-1], parts[1:]): job = {key:func_dict[key][i:j] for key in func_dict} job.update({'func':test_func}) jobs.append(job) pool = mp.Pool(processes = 2) outputs = pool.imap_unordered(expandCall, jobs) out = [] # Process asynchronous output, report progress for out_ in outputs: out.append(out_) pool.close() pool.join()
报错信息
PicklingError:无法pickle函数<function test_func at 0x0000027FBF960F40>:它与__main__.test_func不是同一个对象
原方案无效原因
Lopez de Prado提供的代码是针对**类方法(types.MethodType)**的pickle适配,而当前问题是np.vectorize返回的numpy.vectorize实例对象,两者类型完全不同,因此旧方案无法解决问题。
解决方案
方案1:移除np.vectorize,改用原生批量处理
np.vectorize本质是语法糖,内部仍是循环,且返回对象pickle兼容性差。直接用map处理批量数据更高效且易序列化:
import numpy as np import pandas as pd import multiprocessing as mp # 定义处理单个元素的函数,无需vectorize def test_func(idx, a, b): return [idx, a + b**2] test_df = pd.DataFrame(np.random.normal(size=(10, 2)), columns=['a', 'b']) test_df = test_df.reset_index() func_dict = {col: test_df[col] for col in test_df.columns} def expandCall(kargs): func = kargs['func'] del kargs['func'] # 对传入的Series逐行应用函数 out = list(map(func, *kargs.values())) return out parts = np.linspace(0, test_df.shape[0], 2).astype(int) jobs = [] for i, j in zip(parts[:-1], parts[1:]): job = {key: func_dict[key][i:j] for key in func_dict} job.update({'func': test_func}) jobs.append(job) if __name__ == '__main__': # Windows系统必须添加该保护 pool = mp.Pool(processes=2) outputs = pool.imap_unordered(expandCall, jobs) out = [] for out_ in outputs: out.append(out_) pool.close() pool.join() # 合并结果 final_out = [item for sublist in out for item in sublist] print(final_out)
方案2:为np.vectorize对象自定义pickle逻辑
如果必须保留np.vectorize,可以通过copyreg注册其序列化/反序列化逻辑:
import numpy as np import pandas as pd import multiprocessing as mp import copyreg from numpy.lib.function_base import vectorize # 定义vectorize对象的pickle方法 def pickle_vectorize(vfunc): # 保存原始函数及vectorize配置参数 return unpickle_vectorize, (vfunc.pyfunc, vfunc.signature, vfunc.excluded) def unpickle_vectorize(pyfunc, signature, excluded): # 反序列化时重建vectorize对象 return vectorize(pyfunc, signature=signature, excluded=excluded) # 注册vectorize类的pickle处理 copyreg.pickle(vectorize, pickle_vectorize, unpickle_vectorize) # 定义函数并包装为vectorize def test_func(idx, a, b): return [idx, a + b**2] test_func = np.vectorize(test_func) test_df = pd.DataFrame(np.random.normal(size=(10, 2)), columns=['a', 'b']) test_df = test_df.reset_index() func_dict = {col: test_df[col] for col in test_df.columns} def expandCall(kargs): func = kargs['func'] del kargs['func'] out = func(**kargs) return out parts = np.linspace(0, test_df.shape[0], 2).astype(int) jobs = [] for i, j in zip(parts[:-1], parts[1:]): job = {key: func_dict[key][i:j] for key in func_dict} job.update({'func': test_func}) jobs.append(job) if __name__ == '__main__': # Windows系统必须添加该保护 pool = mp.Pool(processes=2) outputs = pool.imap_unordered(expandCall, jobs) out = [] for out_ in outputs: out.append(out_) pool.close() pool.join() print(out)
内容的提问来源于stack exchange,提问作者Charles0349
相关产品推荐
相关产品推荐

