Python多进程能否作用于非自定义函数?NetworkX聚类多进程调用报错求解
问题原因与修复方案
报错根因
- 第一个错误:你直接调用了
ProcessPoolExecutor类的map方法,没有先创建执行器实例,类方法和实例方法的调用逻辑不匹配,这是触发TypeError的直接原因 - 第二个错误:
nx.clustering需要传入图对象和待计算的节点/节点列表两个核心参数,你当前的传参逻辑是把图对象迭代后传入,完全不符合函数参数要求
修复思路
参考并行介数中心性的分块计算逻辑,将全量节点拆分为多个子集,每个进程处理一个子集的聚类系数计算,最后合并所有进程的结果即可,不需要修改NetworkX内置函数的实现。
修复后代码
import concurrent.futures import csv import networkx as nx from functools import partial # 读入节点和边的逻辑保持不变 with open('nodes.csv', 'r') as nodecsv: nodereader = csv.reader(nodecsv) nodes = [n for n in nodereader][1:] node_names = [n[0] for n in nodes] with open('edges.csv', 'r') as edgecsv: edgereader = csv.reader(edgecsv) edges = [tuple(e) for e in edgereader][1:] DG = nx.DiGraph() DG.add_nodes_from(node_names) DG.add_edges_from(edges) # 定义单进程计算逻辑,固定图参数 def calc_clustering_chunk(nodes_chunk, G): return {node: nx.clustering(G, node) for node in nodes_chunk} # 将节点拆分为多个块,分块数量建议和CPU物理核心数保持一致 chunk_num = 4 chunk_size = len(node_names) // chunk_num node_chunks = [node_names[i*chunk_size : (i+1)*chunk_size] for i in range(chunk_num)] # 处理拆分后剩余的节点 if len(node_names) % chunk_num != 0: node_chunks.append(node_names[chunk_num*chunk_size:]) # 多进程计算入口 if __name__ == '__main__': clustering_res = {} # 实例化执行器后再调用map with concurrent.futures.ProcessPoolExecutor() as executor: # 绑定固定参数DG partial_calc = partial(calc_clustering_chunk, G=DG) for chunk_res in executor.map(partial_calc, node_chunks): clustering_res.update(chunk_res) # 最终clustering_res就是所有节点的聚类系数结果 print(clustering_res)
额外说明
- 必须加
if __name__ == '__main__'保护多进程入口,Windows系统下不加会反复启动进程报错 - 如果你不需要每个节点的聚类系数,只需要全局平均聚类系数,可以在每个子进程里计算完块的平均后再合并,能进一步减少进程间通信开销
内容的提问来源于stack exchange,提问作者eeno
相关产品推荐
相关产品推荐

