优化处理2亿行大文件构建NetworkX图的Python代码 提升运行效率
Python文件处理与NetworkX建图性能优化方案
核心瓶颈定位
你当前代码性能指数级下降的核心原因是建图阶段每处理一行就遍历整张图的所有节点更新属性,属于O(n²)时间复杂度,节点数越多遍历耗时越长,这也是前10万行快、后续越来越慢的根本原因。其他次要瓶颈包括不必要的重复计算、字典操作不规范、文件重复读取、逻辑错误导致的冗余计算等。
第一优先级:代码逻辑优化(改完即可获得几十到上百倍性能提升)
- 移除每行遍历节点的逻辑:节点属性不需要逐行更新,等当前批次10万行建图完成后,统一遍历一次节点设置属性即可,直接把O(n²)复杂度降为O(n)。
- 修复字典快照的浅拷贝问题:原来的
bl=add_balance是引用赋值,后续修改add_balance时bl会同步变化,逻辑完全错误,需要改成bl = add_balance.copy()、snd = de_send.copy()、rec = de_rec.copy()生成独立快照。 - 优化字典查找写法:把所有
if x in dict.keys()改为if x in dict,Python字典的in操作直接查哈希表,不需要额外生成keys列表,性能提升明显。 - 移除文件重复读取:不要提前计算
num_lines,2亿行文件读两次会浪费大量时间,直接按行迭代到文件末尾自动终止即可。 - 简化NetworkX操作:添加单个边直接用
G.add_edge(u, v, weight=amount),比传单元素列表给add_weighted_edges_from更快;同时不需要提前调用add_node,添加边时会自动创建不存在的节点。 - 用
collections.defaultdict代替普通字典:比如add_balance = defaultdict(int),就不需要判断key是否存在,直接做加减运算即可,代码更简洁性能也更好。
第二优先级:多进程并行改造
你的任务逻辑天然适合并行化:主进程负责IO密集的文件读取、三个全局字典的更新、生成每批次需要的快照和原始行数据,CPU密集的建图、CreateTies调用、文件保存操作扔给子进程并行处理。
具体实现步骤:
- 导入
multiprocessing.Pool,设置进程数为CPU核心数-1,避免资源占用过高。 - 主进程每读满20万行(2*gr),就把当前批次的三个字典快照、后10万行的原始行数据、批次编号作为参数,通过
apply_async提交给进程池处理。 - 子进程收到参数后单独完成建图、属性设置、调用
CreateTies、保存gexf文件的全流程,和主进程完全不冲突。
简化示例代码
import networkx as nx from collections import defaultdict from multiprocessing import Pool # 子进程执行的建图保存逻辑 def process_chunk(chunk_data): bl, snd, rec, lines, batch_id = chunk_data gr = 100000 G = nx.DiGraph() GsimA = nx.DiGraph() in_deg = defaultdict(int) out_deg = defaultdict(int) add_balance = defaultdict(int) de_send = defaultdict(int) de_rec = defaultdict(int) tr = None for idx, line in enumerate(lines): if idx % 2 == 0: inc = 0 fields = line.split(";") tr = int(fields[0].split(",")[1]) G.add_node(-tr) GsimA.add_node(-tr) for ff in fields[1:]: ffsplit = ff.split(",") address = int(ffsplit[0]) amount = int(ffsplit[1]) inc += 1 G.add_edge(address, -tr, weight=amount) add_balance[address] -= amount de_send[address] += 1 in_deg[inc] += 1 else: outc = 0 fields = line.split(";") for ff in fields: ffsplit = ff.split(",") address = int(ffsplit[0]) amount = int(ffsplit[1]) outc += 1 G.add_edge(-tr, address, weight=amount) add_balance[address] += amount de_rec[address] += 1 out_deg[outc] += 1 # 统一设置节点属性,只遍历一次 for i in G.nodes(): if i > 0: G.nodes[i]["balance"] = add_balance.get(i, 0) GsimA.nodes[i]["balance"] = bl.get(i, 0) G.nodes[i]["in-degree"] = de_send.get(i, 0) GsimA.nodes[i]["in-degree"] = snd.get(i, 0) G.nodes[i]["out-degree"] = de_rec.get(i, 0) GsimA.nodes[i]["out-degree"] = rec.get(i, 0) GsimA = CreateTies(GsimA, in_deg, out_deg) nx.write_gexf(G, f"C:/Users/uccio/Desktop/tesi2/realgraph{batch_id}.gexf") nx.write_gexf(GsimA, f"C:/Users/uccio/Desktop/tesi2/simAgraph{batch_id}.gexf") def Creation(file): gr = 100000 add_balance = defaultdict(int) de_send = defaultdict(int) de_rec = defaultdict(int) count = 0 current_chunk = [] batch = 0 # 进程池初始化,进程数按自身CPU核心数调整 pool = Pool(processes=7) with open(file, 'r') as f: for line in f: line = line.strip() if not line: continue # 前gr行只更新全局字典 if count < (batch * 2 + 1) * gr: if count % 2 == 0: fields = line.split(";") for ff in fields[1:]: ffsplit = ff.split(",") address = int(ffsplit[0]) amount = int(ffsplit[1]) add_balance[address] -= amount de_send[address] += 1 else: fields = line.split(";") for ff in fields: ffsplit = ff.split(",") address = int(ffsplit[0]) amount = int(ffsplit[1]) add_balance[address] += amount de_rec[address] += 1 # 后gr行收集到块中,交给子进程处理 else: current_chunk.append(line) # 满一个批次提交任务 if len(current_chunk) == gr: bl = add_balance.copy() snd = de_send.copy() rec = de_rec.copy() batch_id = f"{int(batch*gr)}-{int((batch+1)*gr)}" pool.apply_async(process_chunk, args=((bl, snd, rec, current_chunk, batch_id),)) current_chunk = [] batch += 1 count += 1 # 处理最后不足一个批次的剩余数据 if current_chunk: bl = add_balance.copy() snd = de_send.copy() rec = de_rec.copy() batch_id = f"{int(batch*gr)}-{int(count/2)}" process_chunk((bl, snd, rec, current_chunk, batch_id)) pool.close() pool.join() return ""
内容的提问来源于stack exchange,提问作者Duccio Borchi
相关产品推荐
相关产品推荐

