mpi4py启动进程后无响应:并行Betweenness centrality故障排查
我正在尝试对**Betweenness centrality(介数中心性)**进行并行化实现,使用命令mpiexec -n 7 python3 ProjectMain.py启动测试程序。程序启动时无报错信息,但完全不执行任何内容,甚至连开头的测试打印语句都没有输出。经测试,函数外的代码可正常运行,判断问题出在自定义的betweenness_centrality函数中,怀疑是不同rank的进程未同步,执行顺序混乱。
相关代码如下:
from mpi4py import MPI import time import numpy as np import networkx as nx import time import betweenparalel as bpl from heapq import heappush, heappop from itertools import count import random print("test") name = MPI.Get_processor_name() comm = MPI.COMM_WORLD rank = comm.Get_rank() size = comm.Get_size() print(rank) def betweenness_centrality(G, k=None, normalized=True, weight=None, endpoints=False, seed=None): if rank == 0: betweenness = dict.fromkeys(G, 0.0) # b[v]=0 for v in G if k is None: nodes = G else: random.seed(seed) nodes = random.sample(G.nodes(), k) for s in nodes: comm.send(s, dest=s, tag=0) print("sent s = ", s) if rank != 0: betweenness = None betweenness = comm.bcast(betweenness, root=0) s = comm.recv(source=0, tag=0) if rank == 1: if weight is None: # use BFS S, P, sigma = bpl._single_source_shortest_path_basic(G, s) if endpoints: betweenness1 = bpl._accumulate_endpoints(betweenness, S, P, sigma, s) else: betweenness1 = bpl._accumulate_basic(betweenness, S, P, sigma, s) comm.send(betweenness1, dest=2, tag=1) print("1 done") if rank == 2: if weight is None: # use BFS S, P, sigma = bpl._single_source_shortest_path_basic(G, s) if endpoints: betweenness2 = bpl._accumulate_endpoints(betweenness, S, P, sigma, s) else: betweenness2 = bpl._accumulate_basic(betweenness, S, P, sigma, s) betweenness1 = comm.recv( source=1, tag=1) if endpoints: betweenness2 = bpl._accumulate_endpoints(betweenness1, S, P, sigma, s) else: betweenness2 = bpl._accumulate_basic(betweenness1, S, P, sigma, s) comm.send(betweenness2, dest=4, tag=2) print("2 done") if rank == 3: if weight is None: # use BFS S, P, sigma = bpl._single_source_shortest_path_basic(G, s) if endpoints: betweenness3 = bpl._accumulate_endpoints(betweenness, S, P, sigma, s) else: betweenness3 = bpl._accumulate_basic(betweenness, S, P, sigma, s) comm.send(betweenness3, dest=4, tag=3) print("3 done") if rank == 4: if weight is None: # use BFS S, P, sigma = bpl._single_source_shortest_path_basic(G, s) if endpoints: betweenness4 = bpl._accumulate_endpoints(betweenness, S, P, sigma, s) else: betweenness4 = bpl._accumulate_basic(betweenness, S, P, sigma, s) betweenness3 = comm.recv( source=3, tag=3) if endpoints: betweenness4 = bpl._accumulate_endpoints(betweenness3, S, P, sigma, s) else: betweenness4 = bpl._accumulate_basic(betweenness3, S, P, sigma, s) betweenness2 = comm.recv( source=2, tag=2) if endpoints: betweenness4 = bpl._accumulate_endpoints(betweenness2, S, P, sigma, s) else: betweenness4 = bpl._accumulate_basic(betweenness2, S, P, sigma, s) comm.send(betweenness4, dest=6, tag=4) print("4 done") if rank == 5: if weight is None: # use BFS S, P, sigma = bpl._single_source_shortest_path_basic(G, s) if endpoints: betweenness5 = bpl._accumulate_endpoints(betweenness, S, P, sigma, s) else: betweenness5 = bpl._accumulate_basic(betweenness, S, P, sigma, s) comm.send(betweenness5, dest=6, tag=5) print("5 done") if rank == 6: if weight is None: # use BFS S, P, sigma = bpl._single_source_shortest_path_basic(G, s) if endpoints: betweenness6 = bpl._accumulate_endpoints(betweenness, S, P, sigma, s) else: betweenness6 = bpl._accumulate_basic(betweenness, S, P, sigma, s) betweenness5 = comm.recv( source=5, tag=5) if endpoints: betweenness6 = bpl._accumulate_endpoints(betweenness5, S, P, sigma, s) else: betweenness6 = bpl._accumulate_basic(betweenness5, S, P, sigma, s) betweenness4 = comm.recv( source=4, tag=4) if endpoints: betweenness6 = bpl._accumulate_endpoints(betweenness4, S, P, sigma, s) else: betweenness6 = bpl._accumulate_basic(betweenness4, S, P, sigma, s) comm.send(betweenness6, dest=0, tag=6) print("6 done") if rank == 0: betweenness = comm.recv( source=6, tag=6) # rescaling betweenness = bpl._rescale(betweenness, len(G), normalized=normalized, directed=G.is_directed(), k=k) print("0 done") return betweenness if rank != 0: G = None if rank == 0: G = nx.Graph() G.add_node(1) G.add_node(2) G.add_node(3) G.add_node(4) G.add_node(5) G.add_node(6) print("nodes set") G.add_edge(1, 2) G.add_edge(2, 3) G.add_edge(1, 3) G.add_edge(3, 4) G.add_edge(4, 5) G.add_edge(4, 6) G.add_edge(5, 6) print("edges set") stimeA =time.time() betweeness_dictA = bpl.betweenness_centrality(G) print(betweeness_dictA) etimeA =time.time() print("Time of normal Betweenness centralization --- %s seconds ---" % (etimeA - stimeA)) G = comm.bcast(G, root=0) stimeB =time.time() betweeness_dict = betweenness_centrality(G) etimeB =time.time() if rank == 0: print(betweeness_dict) print("Time of paralel Betweenness centralization --- %s seconds ---" % (etimeB - stimeB))
核心问题分析
进程死锁:通信顺序完全颠倒
非0rank的代码逻辑是先执行comm.bcast(betweenness, root=0),再接收s,但rank0在函数内先循环发送s,完全没执行广播操作。这导致所有非0进程卡在广播等待,而rank0后续发送s时,目标进程无法接收,最终整个集群进入死锁状态,连开头的打印都无法正常输出(部分进程被阻塞,输出缓冲区无法刷新)。任务分配逻辑错误
rank0遍历图的节点(1-6),直接给对应编号的rank发送s,但每个非0进程只能接收一次s,当rank0发送s=2给rank2时,rank2还卡在广播等待,无法接收消息,rank0会阻塞在send操作,彻底卡死流程。依赖关系错误
比如rank2需要先接收rank1的消息才能继续,但rank1因为卡在广播等待,根本无法执行到发送步骤,进一步加剧死锁。
修复方案
调整通信顺序:先广播再分发任务
让rank0先完成betweenness的广播,再给其他进程发送任务节点s,确保所有进程同步进入后续逻辑:def betweenness_centrality(G, k=None, normalized=True, weight=None, endpoints=False, seed=None): # 先初始化并广播betweenness,所有进程同步完成这一步 betweenness = dict.fromkeys(G, 0.0) if rank == 0 else None betweenness = comm.bcast(betweenness, root=0) if rank == 0: if k is None: nodes = list(G.nodes()) else: random.seed(seed) nodes = random.sample(G.nodes(), k) # 按rank分配节点,跳过rank0 for idx, s in enumerate(nodes): target_rank = (idx + 1) % size comm.send(s, dest=target_rank, tag=0) else: s = comm.recv(source=0, tag=0) # 后续各rank的处理逻辑保持不变优化任务分配逻辑
不要让节点编号直接绑定rank,而是按进程数均匀拆分节点列表,避免节点编号与rank不匹配的问题(比如图节点数超过进程数、节点编号非连续等场景)。添加调试打印
在关键通信步骤前添加打印,比如print(f"Rank {rank} waiting for broadcast"),帮助快速定位阻塞点。
内容的提问来源于stack exchange,提问作者boris zic

