如何将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
相关产品推荐
相关产品推荐

