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

如何将MPI阻塞式矩阵乘法代码转换为非阻塞通信版本

MPI非阻塞通信版矩阵乘法实现

核心改写逻辑

原阻塞版本使用MPI_Send/MPI_Recv实现通信,通信发起后进程会阻塞等待操作完成才能执行后续逻辑。非阻塞版本替换为MPI_Isend/MPI_Irecv接口,通信发起后进程可以继续执行其他无关操作,待需要依赖通信数据时再调用MPI_Wait/MPI_Waitall确认通信完成,可实现通信与计算的重叠,提升并行效率。
改动点如下:

  • 为每个非阻塞通信操作分配独立的MPI_Request句柄,用于跟踪通信状态
  • 主进程分发数据时,为每个从进程预存待发送的偏移量、行数值的副本,避免非阻塞发送未完成时变量被修改导致数据错误
  • 从进程先批量发起所有数据的接收请求,待所有数据接收完成后再执行矩阵乘法计算
  • 所有发送、接收操作完成后必须调用等待接口确认通信结束,再操作对应缓冲区,避免内存读写冲突
  • 保留原代码的矩阵分块规则、乘法计算逻辑、结果打印逻辑不变,仅替换通信层实现

注意:当前代码默认矩阵维度N可被从进程总数(总进程数-1)整除,如果N无法被整除,需要额外补充余数行的分配逻辑。

非阻塞版完整代码

#include <stdlib.h>
#include <stdio.h>
#include "mpi.h"
#include <time.h>
#include <sys/time.h>

// 矩阵行列数
#define N 4

// 矩阵存储
double matrix_a[N][N],matrix_b[N][N],matrix_c[N][N];

int main(int argc, char **argv)
{
  int processCount, processId, slaveTaskCount, source, dest, rows, offset;

  MPI_Init(&argc, &argv);
  MPI_Comm_rank(MPI_COMM_WORLD, &processId);
  MPI_Comm_size(MPI_COMM_WORLD, &processCount);
  slaveTaskCount = processCount - 1;

  // 主进程逻辑
 if (processId == 0) {
    // 随机生成矩阵A、B
    srand(time(NULL));
    for (int i = 0; i<N; i++) {
      for (int j = 0; j<N; j++) {
        matrix_a[i][j]= rand()%10;
        matrix_b[i][j]= rand()%10;
      }
    }
    
    printf("\n\t\tMatrix - Matrix Multiplication using MPI (Non-blocking)\n");
    // 打印矩阵A
    printf("\nMatrix A\n\n");
    for (int i = 0; i<N; i++) {
      for (int j = 0; j<N; j++) {
        printf("%.0f\t", matrix_a[i][j]);
      }
        printf("\n");
    }
    // 打印矩阵B
    printf("\nMatrix B\n\n");
    for (int i = 0; i<N; i++) {
      for (int j = 0; j<N; j++) {
        printf("%.0f\t", matrix_b[i][j]);
      }
        printf("\n");
    }

    rows = N/slaveTaskCount;
    offset = 0;
    // 预存每个从进程的发送参数,避免非阻塞发送时变量被修改
    int *send_offsets = (int*)malloc(sizeof(int)*slaveTaskCount);
    int *send_rows = (int*)malloc(sizeof(int)*slaveTaskCount);
    // 发送请求数组:每个从进程4个发送操作(offset、rows、A子块、B矩阵)
    MPI_Request *send_reqs = (MPI_Request*)malloc(sizeof(MPI_Request)*slaveTaskCount*4);

    // 批量发起非阻塞发送
    for (dest=1; dest <= slaveTaskCount; dest++)
    {
      int idx = dest - 1;
      send_offsets[idx] = offset;
      send_rows[idx] = rows;
      // 发送偏移量
      MPI_Isend(&send_offsets[idx], 1, MPI_INT, dest, 1, MPI_COMM_WORLD, &send_reqs[idx*4]);
      // 发送分配行数
      MPI_Isend(&send_rows[idx], 1, MPI_INT, dest, 1, MPI_COMM_WORLD, &send_reqs[idx*4+1]);
      // 发送A矩阵对应分块
      MPI_Isend(&matrix_a[offset][0], rows*N, MPI_DOUBLE,dest,1, MPI_COMM_WORLD, &send_reqs[idx*4+2]);
      // 发送完整B矩阵
      MPI_Isend(&matrix_b, N*N, MPI_DOUBLE, dest, 1, MPI_COMM_WORLD, &send_reqs[idx*4+3]);
      
      offset = offset + rows;
    }

    // 接收请求数组:每个从进程3个接收操作(offset、rows、C子块)
    MPI_Request *recv_reqs = (MPI_Request*)malloc(sizeof(MPI_Request)*slaveTaskCount*3);
    int *recv_offsets = (int*)malloc(sizeof(int)*slaveTaskCount);
    int *recv_rows = (int*)malloc(sizeof(int)*slaveTaskCount);
    // 批量发起非阻塞接收
    for (int i = 1; i <= slaveTaskCount; i++)
    {
      int idx = i - 1;
      source = i;
      MPI_Irecv(&recv_offsets[idx], 1, MPI_INT, source, 2, MPI_COMM_WORLD, &recv_reqs[idx*3]);
      MPI_Irecv(&recv_rows[idx], 1, MPI_INT, source, 2, MPI_COMM_WORLD, &recv_reqs[idx*3+1]);
      MPI_Irecv(&matrix_c[idx*rows][0], rows*N, MPI_DOUBLE, source, 2, MPI_COMM_WORLD, &recv_reqs[idx*3+2]);
    }

    // 等待所有发送、接收操作完成
    MPI_Waitall(slaveTaskCount*4, send_reqs, MPI_STATUSES_IGNORE);
    MPI_Waitall(slaveTaskCount*3, recv_reqs, MPI_STATUSES_IGNORE);

    // 校正结果矩阵块位置,适配消息乱序场景
    for(int i=0; i<slaveTaskCount; i++){
      if(recv_offsets[i] != i*rows){
        for(int r=0; r<recv_rows[i]; r++){
          for(int c=0; c<N; c++){
            double tmp = matrix_c[i*rows + r][c];
            matrix_c[i*rows + r][c] = matrix_c[recv_offsets[i]+r][c];
            matrix_c[recv_offsets[i]+r][c] = tmp;
          }
        }
      }
    }

    // 打印结果矩阵
    printf("\nResult Matrix C = Matrix A * Matrix B:\n\n");
    for (int i = 0; i<N; i++) {
      for (int j = 0; j<N; j++)
        printf("%.0f\t", matrix_c[i][j]);
      printf ("\n");
    }
    printf ("\n");

    // 释放动态分配内存
    free(send_offsets);
    free(send_rows);
    free(send_reqs);
    free(recv_offsets);
    free(recv_rows);
    free(recv_reqs);
  }

  // 从进程逻辑
  if (processId > 0) {
    source = 0;
    MPI_Status recv_status[4], send_status[3];
    MPI_Request recv_reqs[4], send_reqs[3];
    int recv_offset, recv_rows;

    // 批量发起非阻塞接收
    MPI_Irecv(&recv_offset, 1, MPI_INT, source, 1, MPI_COMM_WORLD, &recv_reqs[0]);
    MPI_Irecv(&recv_rows, 1, MPI_INT, source, 1, MPI_COMM_WORLD, &recv_reqs[1]);
    MPI_Irecv(&matrix_a, N*N, MPI_DOUBLE, source, 1, MPI_COMM_WORLD, &recv_reqs[2]);
    MPI_Irecv(&matrix_b, N*N, MPI_DOUBLE, source, 1, MPI_COMM_WORLD, &recv_reqs[3]);

    // 等待所有数据接收完成
    MPI_Waitall(4, recv_reqs, recv_status);

    // 执行矩阵乘法计算
    for (int k = 0; k<N; k++) {
      for (int i = 0; i<recv_rows; i++) {
        matrix_c[i][k] = 0.0;
        for (int j = 0; j<N; j++)
          matrix_c[i][k] = matrix_c[i][k] + matrix_a[i][j] * matrix_b[j][k];
      }
    }

    // 批量发起非阻塞发送
    MPI_Isend(&recv_offset, 1, MPI_INT, 0, 2, MPI_COMM_WORLD, &send_reqs[0]);
    MPI_Isend(&recv_rows, 1, MPI_INT, 0, 2, MPI_COMM_WORLD, &send_reqs[1]);
    MPI_Isend(&matrix_c, recv_rows*N, MPI_DOUBLE, 0, 2, MPI_COMM_WORLD, &send_reqs[2]);

    // 等待所有发送完成
    MPI_Waitall(3, send_reqs, send_status);
  }

  MPI_Finalize();
  return 0;
}

编译运行方法

  • 编译命令:mpicc nonblock_matmul.c -o nonblock_matmul -O2
  • 运行示例(N=4时启动4个进程,1主3从刚好均分4行):mpirun -np 4 ./nonblock_matmul

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.01 01:15:39