如何在第二个multiprocessing.Pool中使用第一个Pool生成的变量?
问题:多进程中访问全局变量触发NameError
在Python中使用multiprocessing.Pool时,主进程生成的dfs字典无法在第二个进程池的函数中访问,触发NameError: name 'dfs' is not defined错误。
示例代码
from multiprocessing import Pool import pandas as pd lst = [1, 2, 3] def csv(code): df = pd.DataFrame({code: [code, code**2, code**3]}, index=lst) return {code: df} def mp1(): with Pool(8) as pool: rs = pool.map(csv, lst) dfs = dict((key, val) for k in rs for key, val in k.items()) return dfs def dosomthing(code): dfs[code] = dfs[code] * code return {code: dfs[code]} def mp_dosomething(): with Pool(8) as pool: rs = pool.map(dosomthing, lst) dfc = dict((key, val) for k in rs for key, val in k.items()) return dfc if __name__ == '__main__': dfs = mp1() dfc = mp_dosomething()
报错信息
multiprocessing.pool.RemoteTraceback: """ Traceback (most recent call last): File "C:\Users\NeNe\AppData\Local\Programs\Python\Python310\lib\multiprocessing\pool.py", line 125, in worker result = (True, func(*args, **kwds)) File "C:\Users\NeNe\AppData\Local\Programs\Python\Python310\lib\multiprocessing\pool.py", line 48, in mapstar return list(map(*args)) File "c:\Users\NeNe\OneDrive\Python\test.py", line 17, in dosomthing dfs[code] = dfs[code] * code NameError: name 'dfs' is not defined """ The above exception was the direct cause of the following exception: Traceback (most recent call last): File "c:\Users\NeNe\OneDrive\Python\test.py", line 28, in <module> dfc = mp_dosomething() File "c:\Users\NeNe\OneDrive\Python\test.py", line 22, in mp_dosomething rs = pool.map(dosomthing, lst) File "C:\Users\NeNe\AppData\Local\Programs\Python\Python310\lib\multiprocessing\pool.py", line 367, in map return self._map_async(func, iterable, mapstar, chunksize).get() File "C:\Users\NeNe\AppData\Local\Programs\Python\Python310\lib\multiprocessing\pool.py", line 774, in get raise self._value NameError: name 'dfs' is not defined
错误原因
Python多进程采用复制内存空间的方式创建子进程,主进程中的变量不会自动共享给子进程。dosomthing函数在子进程中执行时,无法找到主进程中定义的dfs变量,因此触发NameError。
解决方案
将dfs作为参数传递给子进程的函数,以下是两种可行的实现方式:
方式1:使用functools.partial绑定参数
通过partial工具将dfs绑定到dosomthing函数,让子进程能获取到该变量:
from multiprocessing import Pool import pandas as pd from functools import partial lst = [1, 2, 3] def csv(code): df = pd.DataFrame({code: [code, code**2, code**3]}, index=lst) return {code: df} def mp1(): with Pool(8) as pool: rs = pool.map(csv, lst) dfs = dict((key, val) for k in rs for key, val in k.items()) return dfs def dosomthing(dfs, code): # 避免修改原dfs,直接返回新的DataFrame updated_df = dfs[code] * code return {code: updated_df} def mp_dosomething(dfs): with Pool(8) as pool: # 绑定dfs到dosomthing函数 bound_func = partial(dosomthing, dfs) rs = pool.map(bound_func, lst) dfc = dict((key, val) for k in rs for key, val in k.items()) return dfc if __name__ == '__main__': dfs = mp1() dfc = mp_dosomething(dfs) print(dfc)
方式2:使用pool.starmap传递多参数
将dfs和code打包成元组,通过starmap传递给函数:
from multiprocessing import Pool import pandas as pd lst = [1, 2, 3] def csv(code): df = pd.DataFrame({code: [code, code**2, code**3]}, index=lst) return {code: df} def mp1(): with Pool(8) as pool: rs = pool.map(csv, lst) dfs = dict((key, val) for k in rs for key, val in k.items()) return dfs def dosomthing(dfs, code): updated_df = dfs[code] * code return {code: updated_df} def mp_dosomething(dfs): with Pool(8) as pool: # 打包参数为元组列表 args_list = [(dfs, code) for code in lst] rs = pool.starmap(dosomthing, args_list) dfc = dict((key, val) for k in rs for key, val in k.items()) return dfc if __name__ == '__main__': dfs = mp1() dfc = mp_dosomething(dfs) print(dfc)
说明
- 上述两种方式都通过序列化传递数据,DataFrame是可被
pickle序列化的类型,因此可以安全地在进程间传递。 - 不建议修改原
dfs变量(子进程中的修改不会同步到主进程),直接返回处理后的结果更可靠。
内容的提问来源于stack exchange,提问作者Ryan
相关产品推荐
相关产品推荐

