Python多进程中共享自定义对象无法更新的问题求助
问题
尝试使用多进程从多个DataFrame中提取特定行,并将其存储到类字典的自定义对象中,相关代码如下:
import multiprocessing as mp class DataOperationMultiProcess: @staticmethod def get_row_single(new_symbol_df_shared, symbol_df, s, index, columns_dict, columns): df = symbol_df.symbols[s] df_columns = df[columns] new_symbol_df_shared.value.symbols[s] = df_columns.iloc[index] print(new_symbol_df_shared.value) @staticmethod def get_row_multi(symbol_df, symbol_data_frame_class_name, index, symbols, columns): new_symbol_df = symbol_data_frame_class_name() columns_dict = symbol_df.column_names() manager = mp.Manager() new_symbol_df_shared = manager.Value(type(new_symbol_df), new_symbol_df) pool = mp.Pool(mp.cpu_count()) for s in symbols: pool.apply_async(DataOperationMultiProcess.get_row_single, args=(new_symbol_df_shared, symbol_df, s, index, columns_dict, columns)) pool.close() pool.join() return new_symbol_df_shared.value
其中symbol_data_frame_class_name是类字典的自定义对象。
运行后发现:get_row_single函数中打印new_symbol_df_shared的值时显示为空对象,最终get_row_multi返回的结果也为空对象。请问问题原因是什么?如何解决?
原因分析
manager.Value的特性限制:mp.Manager().Value设计用于共享单个基本数据类型(如int、str),无法正确追踪自定义类字典对象的内部属性修改。子进程中对new_symbol_df_shared.value.symbols[s]的赋值,属于修改共享对象的内部属性,Value不会自动同步这类变化,导致主进程感知不到子进程的修改。- 进程内存隔离机制:每个子进程拥有独立内存空间,即便通过
Manager创建了共享对象,自定义对象的内部属性修改未通过Manager的同步机制完成,子进程的修改实际无法同步回主进程的共享对象。
解决方案
方案一:用manager.dict()中转存储结果
放弃直接共享自定义类对象,改用Manager提供的dict存储提取的行数据,最后将字典转换为目标自定义对象:
import multiprocessing as mp class DataOperationMultiProcess: @staticmethod def get_row_single(result_dict, symbol_df, s, index, columns): df = symbol_df.symbols[s] df_columns = df[columns] result_dict[s] = df_columns.iloc[index] @staticmethod def get_row_multi(symbol_df, symbol_data_frame_class_name, index, symbols, columns): manager = mp.Manager() result_dict = manager.dict() pool = mp.Pool(mp.cpu_count()) for s in symbols: pool.apply_async(DataOperationMultiProcess.get_row_single, args=(result_dict, symbol_df, s, index, columns)) pool.close() pool.join() # 将共享字典转换为自定义对象 new_symbol_df = symbol_data_frame_class_name() new_symbol_df.symbols = dict(result_dict) return new_symbol_df
方案二:让自定义类支持Manager代理
如果必须直接使用自定义对象,可让自定义类继承BaseManager的可代理类型,确保内部属性修改能被同步:
- 注册自定义类到
Manager:
from multiprocessing.managers import BaseManager # 假设你的自定义类名为SymbolDataFrame class SymbolDataFrame: def __init__(self): self.symbols = {} # 注册自定义类到Manager BaseManager.register('SymbolDataFrame', SymbolDataFrame)
- 修改多进程代码:
class DataOperationMultiProcess: @staticmethod def get_row_single(new_symbol_df_shared, symbol_df, s, index, columns): df = symbol_df.symbols[s] df_columns = df[columns] new_symbol_df_shared.symbols[s] = df_columns.iloc[index] print(new_symbol_df_shared.symbols) @staticmethod def get_row_multi(symbol_df, symbol_data_frame_class_name, index, symbols, columns): manager = BaseManager() manager.start() # 通过Manager创建自定义对象实例 new_symbol_df_shared = manager.SymbolDataFrame() pool = mp.Pool(mp.cpu_count()) for s in symbols: pool.apply_async(DataOperationMultiProcess.get_row_single, args=(new_symbol_df_shared, symbol_df, s, index, columns)) pool.close() pool.join() manager.shutdown() # 将共享对象内容复制到本地对象 new_symbol_df = symbol_data_frame_class_name() new_symbol_df.symbols = dict(new_symbol_df_shared.symbols) return new_symbol_df
内容的提问来源于stack exchange,提问作者Anthraxff
相关产品推荐
相关产品推荐

