Python多进程共享内存:不相交数组操作赋值失效问题问询
你遇到的问题是Python多进程的核心特性导致的:每个子进程会创建主进程数据的独立副本。你定义的dump是普通的Python列表,当子进程启动时,它会复制这个列表到自己的内存空间里——你在子进程中修改的dump[i],只是修改了子进程本地的副本,完全不会影响主进程中的原始dump列表,所以最后主进程打印出来的还是全None的初始状态。
下面提供几种可行的解决方法,按推荐程度排序:
方法1:让子进程返回结果,主进程合并(最推荐)
这种方法避开了共享内存的复杂问题,是多进程编程的常规做法——让每个子进程处理自己的区间,然后返回该区间的结果,最后由主进程把所有结果合并到最终的dump列表中。
修改后的代码如下:
from multiprocessing import Pool import math bins = 1000 buckets = 4 # number of processes def iAmSlow(i): if i % 13 == 0: return [], [] return [k*i for k in range(5)], [k-i for k in range(5)] def iAmParallelizable(args): start, end = args # 存储当前进程处理的结果,而不是修改全局列表 result = {} for i in range(start, end): a, b = iAmSlow(i) result[i] = (a, b) print(result[i], i) return result incr = math.ceil(bins / buckets) starts = [min(incr*j, bins) for j in range(buckets+1)] arguments = [(starts[i], starts[i+1]) for i in range(buckets)] dump = [None for i in range(bins)] with Pool(buckets) as thrd: # 获取所有子进程返回的结果字典 process_results = thrd.map(iAmParallelizable, arguments) # 合并所有结果到dump列表 for res in process_results: for idx, value in res.items(): dump[idx] = value # 验证结果 for i in range(len(dump)): print(dump[i])
为什么这能行?:每个子进程处理完自己的区间后,把结果以字典形式返回给主进程,主进程再统一把这些结果填充到原始的dump列表中——所有修改都发生在主进程,完全没有共享内存的问题。
方法2:使用multiprocessing.Manager创建共享列表
如果你确实需要让子进程直接修改共享的数据结构,可以用multiprocessing.Manager,它会创建一个跨进程共享的列表对象,所有子进程操作的都是同一个列表。
代码示例:
from multiprocessing import Pool, Manager import math bins = 1000 buckets = 4 # number of processes def iAmSlow(i): if i % 13 == 0: return [], [] return [k*i for k in range(5)], [k-i for k in range(5)] def iAmParallelizable(args): dump, start, end = args for i in range(start, end): a, b = iAmSlow(i) dump[i] = (a, b) print(dump[i], i) incr = math.ceil(bins / buckets) starts = [min(incr*j, bins) for j in range(buckets+1)] # 把共享列表和区间参数一起传入 arguments = [(dump, starts[i], starts[i+1]) for i in range(buckets)] with Manager() as manager: # 创建跨进程共享的列表 dump = manager.list([None for i in range(bins)]) with Pool(buckets) as thrd: thrd.map(iAmParallelizable, arguments) # 把共享列表转成普通列表(可选) dump = list(dump) # 验证结果 for i in range(len(dump)): print(dump[i])
注意:Manager是通过进程间通信实现的共享,比方法1的性能稍差,但适合需要实时共享数据的场景。
方法3:使用multiprocessing.Array存储对象
如果你想使用底层的共享内存,可以用multiprocessing.Array,但因为要存储复杂的Python对象(元组、列表),需要指定ctypes.py_object类型:
from multiprocessing import Pool, Array import math import ctypes bins = 1000 buckets = 4 # number of processes def iAmSlow(i): if i % 13 == 0: return [], [] return [k*i for k in range(5)], [k-i for k in range(5)] def iAmParallelizable(args): dump, start, end = args for i in range(start, end): a, b = iAmSlow(i) dump[i] = (a, b) print(dump[i], i) incr = math.ceil(bins / buckets) starts = [min(incr*j, bins) for j in range(buckets+1)] arguments = [(dump, starts[i], starts[i+1]) for i in range(buckets)] # 创建共享内存数组,存储Python对象 dump = Array(ctypes.py_object, [None for _ in range(bins)]) with Pool(buckets) as thrd: thrd.map(iAmParallelizable, arguments) # 转成普通列表查看结果 dump_list = list(dump) for i in range(len(dump_list)): print(dump_list[i])
注意:这种方法直接操作共享内存,性能比Manager好,但要注意对象的序列化问题,且如果多个进程同时修改同一个位置(你这里是互不相交的区间,所以没问题),需要加锁。
内容的提问来源于stack exchange,提问作者lee

