Python多进程技术问询:向ProcessPoolExecutor提交对象列表时传递的是拷贝还是引用?
我正尝试对一个大型项目的部分逻辑进行并行化改造。我有一个由点(Point)组成的晶格(Lattice),需要在晶格的每个点上执行计算。为了提升计算速度,我计划将晶格的点划分为多个子列表,通过ProcessPoolExecutor在独立进程中执行计算。但我发现向ProcessPoolExecutor传递列表时,似乎传递的是列表的拷贝而非引用。
以下是我的简化版代码:
from time import time from concurrent.futures import ProcessPoolExecutor class Lattice(): def __init__(self, rndl): self.Points = [] for val in rndl: P = Point(val) self.Points.append(P) class Point(): def __init__(self, val): self.val = val def calculate(self, it = 100): for i in range(it): some_tmp_value = self.val**(1.0/10) some_tmp_value = some_tmp_value**(1.0/10) self.val -= 1 #update values def do_calculation_ser_2(lattice, it = 100): workers = 2 splitted_list = [lattice.Points[i::workers] for i in range(workers)] for sublist in splitted_list: for P in sublist: P.calculate(it) def par_calc_helper(sublist, it): for i,P in enumerate(sublist): P.calculate(it) return None def do_calculation_par_partitioning(lattice, it): workers = 2 #define number of subprocesses #split list in chunks splitted_list = [lattice.Points[i::workers] for i in range(workers)] with ProcessPoolExecutor(max_workers = workers) as executor: for i, sublist in enumerate(splitted_list): future = executor.submit(par_calc_helper, sublist, it) def check_calc(lattice, rndl): for i,(P, val_rnd) in enumerate(zip(lattice.Points, rndl)): error = False P.val += 1 if P.val != val_rnd: print("ERROR - Calulation gone wrong") error = True break if not error: print("Calculation gone right") if __name__ == '__main__': max_iter = 2*8 rndl = [4,4,4,4] lat = Lattice(rndl) do_calculation_ser_2(lat, max_iter) check_calc(lat, rndl) do_calculation_par_partitioning(lat, max_iter) check_calc(lat, rndl)
代码运行输出:
Calculation gone right ERROR - Calulation gone wrong
测试时我初始化了4个val为4的Point对象,串行计算后check_calc验证通过,但并行计算后验证失败:calculate方法确实被调用(插入打印语句可确认),但Point的val值仍为4而非预期的3。我推测向ProcessPoolExecutor传递对象列表时,传递的是对象的拷贝而非引用(与串行场景的引用传递不同)。
请问我的推测是否正确?如果正确,如何避免每次传递时拷贝对象(这对大型计算场景影响很大)?传递列表拷贝并将计算结果替换主列表是否是最优方案?或者有更合适的Python并行计算实现思路?
你的推测完全正确!这是Python多进程编程中非常常见的一个坑,下面详细解释原因并给出可行的解决方案:
为什么传递的是拷贝而非引用?
Python的ProcessPoolExecutor基于multiprocessing模块实现,而多进程之间是完全内存隔离的——每个子进程都有自己独立的内存空间,和主进程互不干扰。当你把sublist(包含Point对象)传递给子进程时,Python会通过pickle模块对这些对象进行序列化(把对象转换成可传输的字节流),然后在子进程中反序列化生成全新的拷贝对象。子进程里对这些拷贝的修改,完全不会影响主进程中的原对象,这就是并行计算后val值没变化的根本原因。
可行的解决方案
1. 返回计算结果替换原列表(最推荐)
这是最直接、易维护且适合大多数场景的方案。核心思路是让子进程处理完数据后,把修改后的对象(或关键属性值)返回给主进程,由主进程替换原列表中的内容。
修改你的代码如下:
首先更新辅助函数,让它返回处理后的子列表:
def par_calc_helper(sublist, it): for P in sublist: P.calculate(it) return sublist # 返回修改后的点列表
然后在并行计算函数中收集结果,合并回原晶格的Points列表:
def do_calculation_par_partitioning(lattice, it): workers = 2 splitted_list = [lattice.Points[i::workers] for i in range(workers)] with ProcessPoolExecutor(max_workers=workers) as executor: # 提交所有任务并获取future对象 futures = [executor.submit(par_calc_helper, sublist, it) for sublist in splitted_list] # 收集所有子进程返回的结果 updated_sublists = [future.result() for future in futures] # 将更新后的点按原顺序放回lattice.Points # 因为我们是按步长拆分的(比如workers=2时,sublist0是[0,2], sublist1是[1,3]) lattice.Points = [] max_len = max(len(sublist) for sublist in updated_sublists) for i in range(max_len): for sublist in updated_sublists: if i < len(sublist): lattice.Points.append(sublist[i])
这个方案的优势在于:
- 逻辑简单,容易理解和调试
- 不需要复杂的共享内存管理
- 对于大型项目,只要合理划分chunk(每个子进程处理足够多的点),序列化/反序列化的开销可以忽略不计
2. 使用共享内存(适合属性简单的对象)
如果你的Point对象只有少数简单类型的属性(比如int、float),可以用multiprocessing提供的共享内存对象,让子进程直接修改共享内存中的值,主进程能实时看到变化。
修改Point类如下:
from multiprocessing import Value class Point(): def __init__(self, val): # 使用共享内存存储val,'d'表示双精度浮点数 self.val = Value('d', val) def calculate(self, it = 100): for i in range(it): some_tmp_value = self.val.value**(1.0/10) some_tmp_value = some_tmp_value**(1.0/10) # 修改共享内存中的值,主进程会同步看到变化 self.val.value -= 1
这样修改后,原有的并行计算函数不需要改动,子进程对self.val.value的修改会直接反映到主进程的原对象上。但要注意:
- 共享内存仅支持基本数据类型,复杂对象无法直接存储
- 如果多个子进程同时修改同一个共享变量,需要加锁(但这里每个Point只被一个子进程处理,所以无需锁)
- 代码复杂度会提升,维护成本更高
3. 避免多进程:改用线程池(仅适用于IO密集型计算)
如果你的计算任务是IO密集型(比如大量网络请求、文件读写),可以用ThreadPoolExecutor代替ProcessPoolExecutor。线程之间共享内存,不需要拷贝对象,修改会直接生效。但注意:你的calculate方法是CPU密集型的,Python的GIL(全局解释器锁)会导致线程无法真正并行利用多核,速度不会提升甚至变慢,所以这个方案不适合你的场景。
4. 优化序列化效率(减少拷贝开销)
默认的pickle序列化对于大型对象可能较慢,可以改用更高效的序列化库(如cloudpickle、dill),或者配置ProcessPoolExecutor使用更高效的序列化方式。这能减少对象拷贝的时间开销,但本质上还是传递拷贝,只是速度更快了。
5. 分布式计算框架(超大型项目)
如果你的项目数据量极大,手动管理进程太繁琐,可以考虑用Dask、Ray这类分布式计算框架。它们能自动处理数据分区、进程间通信和资源调度,比手动使用ProcessPoolExecutor更高效、更易扩展。
总结
对于你的CPU密集型晶格计算场景,返回计算结果替换原列表是最优方案,平衡了实现复杂度和性能。共享内存方案适合属性简单的小对象,但维护成本较高。如果未来项目规模持续增长,分布式框架会是更好的选择。
内容的提问来源于stack exchange,提问作者Felix Steinmetz

