You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Python多进程技术问询:向ProcessPoolExecutor提交对象列表时传递的是拷贝还是引用?

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.30 20:14:05