如何让MPI其他Rank等待Rank 0完成指定任务后再执行?
MPI同步其他进程等待Rank 0完成指定任务的实现方案
我是MPI新手,当前代码结构(hello_world.py)如下:
import time def A(): # some simple codes here to execute. Lets say print('Hello World from A') time.sleep(5) def B(): # some simple codes here to execute.Lets say print('Hello World from B') time.sleep(5) def C(): print('Hello World from C') print('>>>>>>>>>>>>>>>>> The place where I will start to make parallelization') scores = [] for i in range (1, 1000): score = score_calculator(i) # ignore this function scores.append(score) print('<<<<<<<<<<<<<<<<< The place where I will stop to make parallelization')
我将通过mpiexec -n 4 hello_world.py启动4个Rank进程,需求如下:
- A、B函数仅由Rank 0串行执行;
- C函数执行完
print('Hello World from C')后,其他Rank强制等待,直到Rank 0完成后续循环任务; - Rank 0完成后再进行后续并行配置。
以下是Python(基于mpi4py)和C语言的实现方案:
Python 实现方案
from mpi4py import MPI import time def A(): print('Hello World from A') time.sleep(5) def B(): print('Hello World from B') time.sleep(5) def C(comm, rank): print('Hello World from C (Rank {})'.format(rank)) # 可选:确保所有进程完成这步打印后再继续,可根据需求移除 comm.Barrier() if rank == 0: print('>>>>>>>>>>>>>>>>> The place where I will start to make parallelization') scores = [] for i in range(1, 1000): # 模拟score_calculator逻辑 score = i * 2 scores.append(score) print('<<<<<<<<<<<<<<<<< The place where I will stop to make parallelization') # 完成循环后触发同步,释放其他进程 comm.Barrier() else: # 其他进程等待Rank 0完成循环 comm.Barrier() # 后续并行配置逻辑,所有进程从这里开始执行 print('Post-parallelization step starting for Rank {}'.format(rank)) if __name__ == '__main__': comm = MPI.COMM_WORLD rank = comm.Get_rank() # 仅Rank 0执行A、B函数 if rank == 0: A() B() # 所有进程执行C函数 C(comm, rank) MPI.Finalize()
核心逻辑说明
- 利用
MPI.COMM_WORLD获取全局通信器和当前进程Rank; - 通过
comm.Barrier()实现同步:其他进程调用该方法后会阻塞,直到Rank 0完成循环并调用相同方法,所有进程才会继续执行后续代码; - 第一个
Barrier用于同步所有进程的打印输出,可根据实际需求移除。
C语言实现方案
#include <stdio.h> #include <stdlib.h> #include <mpi.h> #include <unistd.h> void A() { printf("Hello World from A\n"); sleep(5); } void B() { printf("Hello World from B\n"); sleep(5); } void C(int rank, MPI_Comm comm) { printf("Hello World from C (Rank %d)\n", rank); // 可选:同步所有进程的打印输出,可根据需求移除 MPI_Barrier(comm); if (rank == 0) { printf(">>>>>>>>>>>>>>>>> The place where I will start to make parallelization\n"); int* scores = malloc(999 * sizeof(int)); for (int i = 1; i < 1000; i++) { // 模拟score_calculator逻辑 scores[i-1] = i * 2; } printf("<<<<<<<<<<<<<<<<< The place where I will stop to make parallelization\n"); free(scores); // 完成循环后触发同步 MPI_Barrier(comm); } else { // 等待Rank 0完成循环 MPI_Barrier(comm); } // 后续并行配置逻辑 printf("Post-parallelization step starting for Rank %d\n", rank); } int main(int argc, char** argv) { MPI_Init(&argc, &argv); int rank; MPI_Comm_rank(MPI_COMM_WORLD, &rank); // 仅Rank 0执行A、B函数 if (rank == 0) { A(); B(); } // 所有进程执行C函数 C(rank, MPI_COMM_WORLD); MPI_Finalize(); return 0; }
核心逻辑说明
- 调用
MPI_Init初始化MPI环境,MPI_Comm_rank获取当前进程Rank; MPI_Barrier是同步核心:阻塞调用进程,直到通信器内所有进程都调用该函数,完美实现其他Rank等待Rank 0完成任务的需求;- 编译执行命令示例:
mpicc -o hello_world hello_world.c && mpiexec -n 4 ./hello_world
内容的提问来源于stack exchange,提问作者Ankur Lahiry
相关产品推荐
相关产品推荐

