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

如何让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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 16:43:28