如何在multiprocessing Pool中结合yield并行生成对称矩阵?
解决multiprocessing中yield生成器无法序列化的问题
你的报错原因很明确:multiprocessing的Pool在进程间传递结果时依赖pickle序列化,但生成器对象是没办法被pickle的——func用了yield,所以它返回的是一个生成器,子进程没法把这个生成器传回主进程,就抛出了can't pickle generator objects的错误。
给你两种可行的解决方案,适配你的场景:
方案一:将生成器改为返回结果列表
最直接的办法是把func里的生成器逻辑改成收集结果到列表后返回,因为列表是可序列化的。修改后的代码如下:
import multiprocessing as mp import random import numpy as np def func(param): i, js = param # 用列表收集所有该行的结果 row_results = [] for j in js: row_results.append((i, j, random.random())) return row_results N = 100 M = np.eye(N, dtype=np.float) params_parallel = [(i, range(i)) for i in range(N)] with mp.Pool(processes=4) as p: # 先遍历每个子进程返回的结果列表,再遍历列表里的元素 for result_list in p.imap_unordered(func, params_parallel): for i, j, value in result_list: M[i,j] = M[j,i] = value
这种方法简单高效,对于N=100的场景完全没问题,即使N更大,每行的元素数量是O(N)级,整体内存开销也在可控范围内。
方案二:用Queue手动实现流式结果传递(进阶)
如果你的矩阵特别大,不想一次性返回整行结果,可以用multiprocessing.Queue来让子进程逐个发送结果到主进程。这种方式需要手动管理进程,代码稍复杂,但能实现类似生成器的流式传递:
import multiprocessing as mp import random import numpy as np def worker(param, queue): i, js = param for j in js: # 把结果逐个放入队列 queue.put((i, j, random.random())) # 放一个结束标记表示该任务完成 queue.put(None) N = 100 M = np.eye(N, dtype=np.float) params_parallel = [(i, range(i)) for i in range(N)] queue = mp.Queue() processes = [] # 启动所有子进程 for param in params_parallel: p = mp.Process(target=worker, args=(param, queue)) processes.append(p) p.start() # 从队列读取结果 completed_tasks = 0 while completed_tasks < len(params_parallel): item = queue.get() if item is None: completed_tasks += 1 continue i, j, value = item M[i,j] = M[j,i] = value # 等待所有进程结束 for p in processes: p.join()
这种方式不需要依赖Pool的序列化机制,每个结果直接通过队列传递,适合超大规模矩阵的场景。
另外你提到无法使用if __name__ == '__main__',在Linux/Unix系统下,multiprocessing用fork方式创建进程,这个限制不会影响代码运行;如果是Windows系统,可能会有进程启动的问题,但你的原代码在Windows下本来也会有问题,所以如果是Linux环境,上面的两种方案都可以直接运行。
内容的提问来源于stack exchange,提问作者Vítor Mangaravite
相关产品推荐
相关产品推荐

