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

MPI程序在非1或arraySize进程数下Send调用阻塞问题求助

MPI程序进程数非1或arraySize时阻塞问题分析

问题场景

编写了一个MPI测试程序,模拟不同进程对应不同工作负载的场景:每个进程负责数组中指定范围的元素,再将这些元素发送给其他所有进程,最终每个进程获取完整数组。但当使用的进程数不是1或arraySize(本案例为4)时,程序会在MPI_Send调用时阻塞,尤其是执行mpirun -np 2 MPItest时阻塞情况明显。其中1进程和4进程时程序正常,3进程会崩溃,2进程阻塞。尝试改用MPI_Isend后仍无法正常运行(flag始终为0),哪怕是-np 4也失效。

原始阻塞代码

#include <mpi.h>
#include <iostream>

int main(int argc, char** argv) {
    int rank, size;
    const int arraySize = 4;
    MPI_Init(&argc, &argv);
    MPI_Comm_rank(MPI_COMM_WORLD, &rank);
    MPI_Comm_size(MPI_COMM_WORLD, &size);

    // 每个进程负责数组中不同的工作负载(一个或多个元素),并将这些元素发送给其他所有进程
    int* sendbuf = new int[arraySize];
    int* recvbuf = new int[arraySize];

    int istart = arraySize/size * rank;
    int istop = (rank == size) ? arraySize : istart + arraySize/size;

    for (int i = istart; i < istop; i++) {
        sendbuf[i] = i;
    }

    std::cout << "Rank " << rank << " sendbuf :" << std::endl;
    // 接收前打印sendbuf
    for (int i = 0; i < arraySize; i++) {
        std::cout << sendbuf[i] << ", ";
    }
    std::cout << std::endl;

    // 将负责的元素发送给所有进程
    for(int i = istart; i < istop; i++){
        for(int j = 0; j < size; j++){
            MPI_Send(&sendbuf[i], 1, MPI_INT, j, i, MPI_COMM_WORLD);
        }
    }

    // 接收完整数组
    for(int i = 0; i < arraySize ; i++){
        int recvRank = i/(arraySize/size);
        MPI_Recv(&recvbuf[i], 1, MPI_INT, recvRank, i, MPI_COMM_WORLD, MPI_STATUS_IGNORE);
    }

    // 接收后打印recvbuf
    std::cout << "Rank " << rank << " recvbuf :" << std::endl;
    for (int i = 0; i < arraySize; i++) {
        std::cout << recvbuf[i] << ", ";
    }
    std::cout << std::endl;

    delete[] sendbuf;
    delete[] recvbuf;

    MPI_Finalize();
    return 0;
}

尝试的MPI_Isend修改版本代码

#include <mpi.h>
#include <iostream>

int main(int argc, char** argv) {
    int rank, size;
    const int arraySize = 4;
    MPI_Init(&argc, &argv);
    MPI_Comm_rank(MPI_COMM_WORLD, &rank);
    MPI_Comm_size(MPI_COMM_WORLD, &size);

    // 每个进程负责数组中不同的工作负载(一个或多个元素),并将这些元素发送给其他所有进程
    int* sendbuf = new int[arraySize];
    int* recvbuf = new int[arraySize];

    int istart = arraySize/size * rank;
    int istop = (rank == size) ? arraySize : istart + arraySize/size;

    for (int i = istart; i < istop; i++) {
        sendbuf[i] = i;
    }

    std::cout << "Rank " << rank << " sendbuf :" << std::endl;
    // 接收前打印sendbuf
    for (int i = 0; i < arraySize; i++) {
        std::cout << sendbuf[i] << ", ";
    }
    std::cout << std::endl;

    // 将负责的元素发送给所有进程
    for(int i = istart; i < istop; i++){
        for(int j = 0; j < size; j++){
            MPI_Request request;
            //MPI_Send(&sendbuf[i], 1, MPI_INT, j, i, MPI_COMM_WORLD);
            MPI_Isend(&sendbuf[i], 1, MPI_INT, j, i, MPI_COMM_WORLD, &request);
            // 检查发送是否完成
            int flag = 0;
            MPI_Test(&request, &flag, MPI_STATUS_IGNORE);
            const int numberOfRetries = 10;
            if(flag == 0){ // 操作未完成
                std::cerr << "Error in sending, waiting" << std::endl;
                for(int k = 0; k < numberOfRetries; k++){
                    MPI_Test(&request, &flag, MPI_STATUS_IGNORE);
                    if(flag == 1){
                        break;
                    }
                }
                if(flag == 0){
                    std::cerr << "Error in sending, aborting" << std::endl;
                    MPI_Abort(MPI_COMM_WORLD, 1);
                }
                
            }
        }
    }

    // 接收完整数组
    for(int i = 0; i < arraySize ; i++){
        int recvRank = i/(arraySize/size);
        MPI_Recv(&recvbuf[i], 1, MPI_INT, recvRank, i, MPI_COMM_WORLD, MPI_STATUS_IGNORE);
    }

    // 接收后打印recvbuf
    std::cout << "Rank " << rank << " recvbuf :" << std::endl;
    for (int i = 0; i < arraySize; i++) {
        std::cout << recvbuf[i] << ", ";
    }
    std::cout << std::endl;

  
    //MPI_Alltoall(sendbuf, 1, MPI_INT, recvbuf, 1, MPI_INT, MPI_COMM_WORLD);

    delete[] sendbuf;
    delete[] recvbuf;

    MPI_Finalize();
    return 0;
}

问题核心分析

1. MPI_Send版本死锁原因

  • 进程数为2时:每个进程会先执行完所有MPI_Send才进入接收阶段。比如Rank 0需要给Rank 1发送2个元素,Rank 1同样要先给Rank 0发送2个元素,两者都卡在发送对方的请求上——因为对方还没进入接收流程,直接导致双向死锁。
  • 进程数为4时:每个进程只负责1个元素,发送给包括自己在内的4个进程。其中给自己的MPI_Send可以直接完成(不需要等待外部接收),给其他进程的发送也能在所有进程发送完自身元素后,立刻进入接收流程,因此不会死锁。
  • 进程数为3时:arraySize/size是1(整数除法),Rank 0负责元素0,Rank1负责元素1,Rank2负责元素2,但数组有4个元素,Rank2的istop是2+1=3,漏掉了元素3。后续接收元素3时找不到对应的发送方,直接崩溃。

2. MPI_Isend版本失效原因

MPI_Isend是非阻塞发送,但需要接收方匹配MPI_Recv才能完成。代码中在发送后立刻用MPI_Test检查,尝试几次后就终止程序,而此时接收方还没开始接收,导致flag始终为0,触发终止逻辑。非阻塞操作需要配合等待或后续的接收操作,不能强行在发送阶段就要求完成。

解决方案

方法1:使用MPI_Alltoallv(推荐)

这是MPI标准的集体通信函数,专门用于所有进程间的数据分发,能自动处理不同进程负载不同的情况,彻底避免手动通信的死锁问题:

// 计算每个进程发送的元素数量和位移
int sendcounts[size];
int sdispls[size];
int recvcounts[size];
int rdispls[size];

for(int i=0; i<size; i++){
    sendcounts[i] = (i == size-1) ? arraySize - (arraySize/size)*i : arraySize/size;
    sdispls[i] = (arraySize/size)*i;
    recvcounts[i] = sendcounts[i];
    rdispls[i] = sdispls[i];
}

MPI_Alltoallv(sendbuf, sendcounts, sdispls, MPI_INT, recvbuf, recvcounts, rdispls, MPI_INT, MPI_COMM_WORLD);

方法2:调整手动通信顺序

如果必须手动实现,可以采用先接收后发送的逻辑,或者对通信排序(比如让Rank低的进程先接收Rank高的进程的数据,再发送),避免双向阻塞。也可以用MPI_Sendrecv替代单独的Send和Recv,直接匹配每个发送和接收操作。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 11:02:34