You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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))

核心问题分析

  1. 进程死锁:通信顺序完全颠倒
    非0rank的代码逻辑是先执行comm.bcast(betweenness, root=0),再接收s,但rank0在函数内先循环发送s,完全没执行广播操作。这导致所有非0进程卡在广播等待,而rank0后续发送s时,目标进程无法接收,最终整个集群进入死锁状态,连开头的打印都无法正常输出(部分进程被阻塞,输出缓冲区无法刷新)。

  2. 任务分配逻辑错误
    rank0遍历图的节点(1-6),直接给对应编号的rank发送s,但每个非0进程只能接收一次s,当rank0发送s=2给rank2时,rank2还卡在广播等待,无法接收消息,rank0会阻塞在send操作,彻底卡死流程。

  3. 依赖关系错误
    比如rank2需要先接收rank1的消息才能继续,但rank1因为卡在广播等待,根本无法执行到发送步骤,进一步加剧死锁。


修复方案

  1. 调整通信顺序:先广播再分发任务
    让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的处理逻辑保持不变
    
  2. 优化任务分配逻辑
    不要让节点编号直接绑定rank,而是按进程数均匀拆分节点列表,避免节点编号与rank不匹配的问题(比如图节点数超过进程数、节点编号非连续等场景)。

  3. 添加调试打印
    在关键通信步骤前添加打印,比如print(f"Rank {rank} waiting for broadcast"),帮助快速定位阻塞点。

内容的提问来源于stack exchange,提问作者boris zic

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.30 16:39:21