MPI C++如何高效传递大型二维数组实现程序并行加速
解决方案
先说明你之前方案的问题
- 小块分发方案慢的核心原因:单次MPI通信有固定的协议开销,1000×1000数组对应10000个10×10块,累计2万次点对点通信的开销远超过块计算本身的开销,自然比串行慢。
- 全量数组分发方案随机崩溃的原因:你写的
arrayAlloc函数是每行单独申请内存,整个二维数组的存储是不连续的,&data[0][0]仅指向第一行的首地址,直接接收总长度size*size的数据会越界访问未分配内存,触发随机崩溃。 - 内存地址传递方案确实不可行,MPI进程拥有独立的虚拟地址空间,跨进程传递地址没有任何意义。
推荐实现方案
你的场景中所有10×10块的计算完全独立,属于典型的易并行场景,用下面的方案可以做到极低通信开销:
第一步:统一使用连续存储的二维数组
不管是主进程还是子进程,都用连续内存分配二维数组,既保留[][]的访问方式,又满足MPI连续数据收发的要求:
// 连续二维数组分配 double** arrayAlloc(int size) { double** arr = new double*[size]; // 一次性申请所有元素的连续内存 arr[0] = new double[size * size]; for (int i = 1; i < size; i++) { // 每行指针指向对应位置 arr[i] = arr[0] + i * size; } return arr; } // 对应的释放逻辑 void arrayFree(double** arr) { delete[] arr[0]; delete[] arr; }
第二步:用MPI集体通信实现低开销并行
1000×1000的double数组仅8MB,哪怕是更大的10000×10000数组也才800MB,完全可以一次广播给所有进程,仅需要2次集体通信就能完成整个流程,比点对点通信效率高几个量级:
int main(int argc, char** argv) { MPI_Init(&argc, &argv); int rank, process_cnt; MPI_Comm_rank(MPI_COMM_WORLD, &rank); MPI_Comm_size(MPI_COMM_WORLD, &process_cnt); const int dataSize = 1000; const int blockSize = 10; const int totalBlockCnt = (dataSize / blockSize) * (dataSize / blockSize); double** data = nullptr; // 主进程初始化原始数组 if (rank == 0) { data = arrayAlloc(dataSize); // 这里写你原始数组的初始化逻辑 } else { // 子进程提前分配连续数组空间,用于接收广播 data = arrayAlloc(dataSize); } // 1. 一次广播把全量数组发给所有进程,开销极低 MPI_Bcast(&(data[0][0]), dataSize * dataSize, MPI_DOUBLE, 0, MPI_COMM_WORLD); // 2. 每个进程计算自己负责的块范围 int blockPerProcess = totalBlockCnt / process_cnt; int startBlock = rank * blockPerProcess; // 最后一个进程负责剩下的所有块 int endBlock = (rank == process_cnt - 1) ? totalBlockCnt : startBlock + blockPerProcess; int localResultCnt = endBlock - startBlock; double* localResults = new double[localResultCnt]; // 3. 本地计算,无任何通信 for (int i = startBlock; i < endBlock; i++) { // 计算块对应的二维坐标 int blockY = (i / (dataSize / blockSize)) * blockSize; int blockX = (i % (dataSize / blockSize)) * blockSize; // 拷贝10*10块到缓存,或者直接传指针给calc函数 double blockBuf[100]; for (int y = 0; y < blockSize; y++) { for (int x = 0; x < blockSize; x++) { blockBuf[y * blockSize + x] = data[blockY + y][blockX + x]; } } localResults[i - startBlock] = function->calc(blockBuf); } // 4. 一次聚合把所有结果收集回主进程 double* allResults = nullptr; if (rank == 0) { allResults = new double[totalBlockCnt]; } MPI_Gather(localResults, localResultCnt, MPI_DOUBLE, allResults, blockPerProcess, MPI_DOUBLE, 0, MPI_COMM_WORLD); // 最后一个进程多出来的块单独收一下(可选,也可以用MPI_Gatherv处理) if (rank == process_cnt - 1 && totalBlockCnt % process_cnt != 0) { int extraCnt = totalBlockCnt % process_cnt; MPI_Send(localResults + blockPerProcess, extraCnt, MPI_DOUBLE, 0, 0, MPI_COMM_WORLD); } if (rank == 0 && totalBlockCnt % process_cnt != 0) { int extraCnt = totalBlockCnt % process_cnt; MPI_Recv(allResults + process_cnt * blockPerProcess, extraCnt, MPI_DOUBLE, process_cnt - 1, 0, MPI_COMM_WORLD, MPI_STATUS_IGNORE); } // 后续主进程处理allResults,释放内存即可 // ... 省略释放内存逻辑 MPI_Finalize(); return 0; }
如果数组极大(超过单进程内存承受范围)
可以改用MPI_Scatterv给每个进程仅分发它需要的块对应的数据,不需要广播全量数组,逻辑类似,只是把广播替换为按需分发即可。
内容的提问来源于stack exchange,提问作者Raspberry
相关产品推荐
相关产品推荐

