MPI: Distributed-Memory Parallelism
The rank model
One executable, N copies, zero shared memory. Each copy knows its identity (rank), the group size (size), and communicates only by send/recv or collectives. The launch command is mpiexec -n 4 ./app (a.k.a. mpirun). Every program starts with MPI_INIT and ends with MPI_FINALIZE:
program hello_mpi
use mpi_f08 ! modern Fortran 2008 bindings, all lowercase
implicit none
integer :: rank, size, ierr
call MPI_Init(ierr)
call MPI_Comm_rank(MPI_COMM_WORLD, rank, ierr)
call MPI_Comm_size(MPI_COMM_WORLD, size, ierr)
print '(a,i0,a,i0,a)', 'hello from rank ', rank, ' of ', size
call MPI_Finalize(ierr)
end program hello_mpi
mpifort -O2 hello_mpi.f90 -o hello_mpi
mpiexec -n 4 ./hello_mpi
# hello from rank 3 of 4
# hello from rank 0 of 4
# ...
Order of the printed lines varies — ranks run concurrently, and that is exactly the point: the program must work for any interleaving. Modern projects use the mpi_f08 module (type-safe, lowercase names); legacy code uses include 'mpif.h' — recognize it, do not adopt it.
Point-to-point: send and receive
The atomic unit of MPI programming is a message: MPI_Send(data, destination_rank, tag) and MPI_Recv(buffer, source_rank, status). The sender and receiver must match on communicator, ranks, tag and buffer type — the compiler cannot check this, only the two ends can agree. The classic pattern is a reduction by hand: rank 0 distributes work, every rank returns a chunk result, rank 0 sums them:
program ring_reduce
use mpi_f08
implicit none
integer :: rank, size, ierr, tag
real :: my_value, received
call MPI_Init(ierr)
call MPI_Comm_rank(MPI_COMM_WORLD, rank, ierr)
call MPI_Comm_size(MPI_COMM_WORLD, size, ierr)
my_value = real(rank) + 1.0 ! rank 0 → 1.0, rank 3 → 4.0
tag = 1
if (rank == 0) then
call MPI_Recv(received, 1, MPI_REAL, MPI_ANY_SOURCE, tag, &
MPI_COMM_WORLD, MPI_STATUS_IGNORE, ierr)
print '(a,f6.1)', 'got a partial from rank ', received
else
call MPI_Send(my_value, 1, MPI_REAL, 0, tag, MPI_COMM_WORLD, ierr)
end if
call MPI_Finalize(ierr)
end program ring_reduce
The MPI_ANY_SOURCE wildcard and MPI_STATUS_IGNORE placeholder keep this short; real code stores the status to learn which rank answered. Point-to-point is the substrate every higher pattern builds on — deadlocks here are the developer's classic first MPI rite of passage (MPI_Send blocks; nested MPI_Send before MPI_Recv on both sides is a symptom to recognize).
Collectives: one call, every rank
Almost all real reductions use the collective library instead of hand-rolled rings: MPI_Reduce (fold into one root), MPI_Allreduce (fold into all roots), MPI_Bcast (one root to everyone), MPI_Gather/MPI_Scatter (collect or distribute array slices), and MPI_Barrier. All ranks must call the same collective — a single rank calling MPI_Reduce alone hangs the communicator:
program allreduce_demo
use mpi_f08
implicit none
integer :: rank, ierr
real :: local, global
call MPI_Init(ierr)
call MPI_Comm_rank(MPI_COMM_WORLD, rank, ierr)
local = real(rank * rank)
call MPI_Allreduce(local, global, 1, MPI_REAL, MPI_SUM, &
MPI_COMM_WORLD, ierr)
print '(a,i0,a,f8.0)', 'rank ', rank, ' sees global sum = ', global
call MPI_Finalize(ierr)
end program allreduce_demo
Domain decomposition: the pattern that scales
Large grids are split into tiles, one per rank; each tile carries a halo of ghost rows/columns copied from neighbours before every step. The stencil from the performance lesson becomes an MPI solver in three pieces: MPI_Cart_create to name neighbours, MPI_Sendrecv to exchange halos without deadlock, and local stencil math after every exchange.
Fig. 1 — Each rank owns a tile and a halo; sendrecv swaps boundaries each step.
! Sketch of the halo exchange — the heart of every grid solver.
call MPI_Sendrecv(u_right, 1, MPI_DOUBLE_PRECISION, east_rank, 0, &
u_left, 1, MPI_DOUBLE_PRECISION, west_rank, 0, &
cart_comm, MPI_STATUS_IGNORE, ierr)
call MPI_Sendrecv(u_bottom, 1, MPI_DOUBLE_PRECISION, south_rank, 0, &
u_top, 1, MPI_DOUBLE_PRECISION, north_rank, 0, &
cart_comm, MPI_STATUS_IGNORE, ierr)
Hybrid MPI + OpenMP
The production recipe for multicore clusters: one MPI rank per node or per socket, OpenMP threads inside each rank exploiting the shared memory, one rank per GPU for offload regions. Message volume drops by the thread factor, and Amdahl's serial fractions shrink with each level. The pattern composes cleanly — the parallel lesson's OpenMP loops run unmodified inside an MPI rank.
# typical launch for a hybrid job (SLURM):
# srun --ntasks=4 --cpus-per-task=8 ./hybrid_app
export OMP_NUM_THREADS=8
mpiexec -n 4 ./hybrid_app
mingw-w64-ucrt-x86_64-openmpi), then compile with mpifort. The demos for this lesson include a complete sendrecv ring that runs on your laptop with mpiexec -n 4.