Getting Started with MPI
Parallel computers in the MIMD (Multiple Instruction, Multiple Data) category are primarily divided into two major architectural paradigms: shared-memory systems and distributed-memory systems.
- In a distributed-memory system, the hardware consists of a collection of independent core-memory pairs (often separate compute nodes) connected by a communication network (interconnect). The memory associated with a given core is directly accessible only to that core.
- In contrast, in a shared-memory system, multiple processing cores connect directly to a single globally accessible memory space where every core can read and write to any physical memory location.
flowchart TD
subgraph Distributed["Figure 3.1: Distributed-Memory System"]
direction TB
subgraph Node0["Node 0"]
CPU0["CPU / Core"] <--> MEM0["Private Memory"]
end
subgraph Node1["Node 1"]
CPU1["CPU / Core"] <--> MEM1["Private Memory"]
end
subgraph Node2["Node 2"]
CPU2["CPU / Core"] <--> MEM2["Private Memory"]
end
subgraph NodeN["Node p-1"]
CPUn["CPU / Core"] <--> MEMn["Private Memory"]
end
MEM0 <==> NET["Interconnect Network"]
MEM1 <==> NET
MEM2 <==> NET
MEMn <==> NET
endflowchart TD
subgraph Shared["Figure 3.2: Shared-Memory System"]
direction TB
CPU0["CPU / Core 0"] <==> BUS["Interconnect / Bus"]
CPU1["CPU / Core 1"] <==> BUS
CPU2["CPU / Core 2"] <==> BUS
CPUn["CPU / Core p-1"] <==> BUS
BUS <==> GMEM["Globally Shared Memory"]
endTo program distributed-memory computers, software developers rely on message-passing. In message-passing programs, an autonomous instance of a program running on a core-memory pair is called a process. Since processes cannot access each other’s memory directly, they communicate explicitly across the network by calling library functions: one process executes a send function, while the target process executes a matching receive function.
The standard specification for message passing in high-performance computing is MPI (Message-Passing Interface). MPI is not a standalone programming language; rather, it is a standardized, portable library of functions callable from languages such as C, C++, and Fortran.
3.1 A First MPI Program: Parallel Greetings
Section titled “3.1 A First MPI Program: Parallel Greetings”In standard C, introductory programs begin with printing "hello, world". In an MPI environment, rather than having every process write uncoordinated text to the screen, a standard pattern is to designate a coordinator process (typically process rank 0) to collect and print messages sent to it by all other processes.
In parallel programming, processes in a group are identified by unique, non-negative integer identifiers called ranks. If an application runs with processes, the ranks are numbered consecutively:
Below is the complete C source code for Program 3.1 (mpi_hello.c):
#include <stdio.h>#include <string.h> /* For strlen */#include <mpi.h> /* For MPI functions, etc */
const int MAX_STRING = 100;
int main(void) { char greeting[MAX_STRING]; int comm_sz; /* Number of processes */ int my_rank; /* Calling process rank */
/* Initialize MPI execution environment */ MPI_Init(NULL, NULL);
/* Get the total number of processes */ MPI_Comm_size(MPI_COMM_WORLD, &comm_sz);
/* Get the individual rank of this process */ MPI_Comm_rank(MPI_COMM_WORLD, &my_rank);
if (my_rank != 0) { /* Worker processes construct and send their greetings */ sprintf(greeting, "Greetings from process %d of %d!", my_rank, comm_sz); MPI_Send(greeting, strlen(greeting) + 1, MPI_CHAR, 0, 0, MPI_COMM_WORLD); } else { /* Process 0 prints its own greeting first */ printf("Greetings from process %d of %d!\n", my_rank, comm_sz);
/* Then receives and prints greetings from processes 1, 2, ..., comm_sz-1 */ for (int q = 1; q < comm_sz; q++) { MPI_Recv(greeting, MAX_STRING, MPI_CHAR, q, 0, MPI_COMM_WORLD, MPI_STATUS_IGNORE); printf("%s\n", greeting); } }
/* Terminate MPI environment and release resources */ MPI_Finalize(); return 0;} /* main */3.1.1 Compilation and Execution
Section titled “3.1.1 Compilation and Execution”Most MPI distributions (such as MPICH or OpenMPI) provide a compiler wrapper called mpicc:
$ mpicc -g -Wall -o mpi_hello mpi_hello.cA wrapper script acts as an intermediary for the underlying C compiler (such as gcc or clang). It automatically passes the required include directories where mpi.h resides, sets compiler flags, and links against the MPI runtime libraries.
To execute an MPI binary across multiple processes, the launcher utility mpiexec (or mpirun) is used:
$ mpiexec -n <number_of_processes> ./mpi_helloFor example, running with a single process:
$ mpiexec -n 1 ./mpi_helloGreetings from process 0 of 1!Running with 4 processes:
$ mpiexec -n 4 ./mpi_helloGreetings from process 0 of 4!Greetings from process 1 of 4!Greetings from process 2 of 4!Greetings from process 3 of 4!When mpiexec -n 4 ./mpi_hello is invoked:
- The operating system starts 4 separate instances (processes) of the
mpi_helloexecutable. - The runtime assigns each instance a distinct rank ().
- The MPI implementation maps processes to CPU cores and establishes communication channels across the interconnect.
3.1.2 Anatomy of an MPI Program
Section titled “3.1.2 Anatomy of an MPI Program”Key structural elements of MPI C programs include:
- Header inclusion:
#include <mpi.h>contains all MPI function prototypes, constants, and type definitions. - Naming conventions: All MPI identifiers begin with the prefix
MPI_.- For function names and types, the first letter after the underscore is capitalized (e.g.,
MPI_Init,MPI_Comm,MPI_Datatype). - Defined constants and macros are written in all capital letters (e.g.,
MPI_COMM_WORLD,MPI_CHAR,MPI_STATUS_IGNORE).
- For function names and types, the first letter after the underscore is capitalized (e.g.,
3.1.3 MPI_Init and MPI_Finalize
Section titled “3.1.3 MPI_Init and MPI_Finalize”Every MPI program must initialize the runtime environment before invoking any other MPI functions, and must cleanly shut it down before exiting:
int MPI_Init( int* argc_p /* in/out: pointer to main's argc */, char*** argv_p /* in/out: pointer to main's argv */);
int MPI_Finalize(void);MPI_Init: Allocates internal communications buffers, initializes hardware interfaces, and assigns ranks to processes. If command-line arguments are not inspected by MPI,NULLcan be passed for both parameters.MPI_Finalize: Cleans up internal structures and frees resources. No MPI calls should be made afterMPI_Finalize().
The standard structure of an MPI program is:
#include <mpi.h>
int main(int argc, char* argv[]) { /* No MPI calls before this */ MPI_Init(&argc, &argv);
/* Parallel computation and message passing */
MPI_Finalize(); /* No MPI calls after this */ return 0;}3.1.4 Communicators: MPI_Comm_size and MPI_Comm_rank
Section titled “3.1.4 Communicators: MPI_Comm_size and MPI_Comm_rank”In MPI, a communicator is an opaque object representing a collection of processes that can send messages to one another. During initialization, MPI creates a default communicator named MPI_COMM_WORLD containing all processes spawned by mpiexec.
To inspect the communicator properties, processes call:
int MPI_Comm_size( MPI_Comm comm /* in: communicator */, int* comm_sz_p /* out: number of processes */);
int MPI_Comm_rank( MPI_Comm comm /* in: communicator */, int* my_rank_p /* out: rank of calling process */);comm_szreceives the total number of processes incomm.my_rankreceives the integer ID of the calling process ().
3.1.5 The SPMD Paradigm
Section titled “3.1.5 The SPMD Paradigm”MPI programs typically follow the Single Program, Multiple Data (SPMD) model:
- A single executable file is compiled and launched across all cores.
- Different processes perform different tasks by branching based on their rank:
if (my_rank == 0) { /* Code executed exclusively by coordinator process */} else { /* Code executed by worker processes */}A critical advantage of SPMD in MPI is elasticity: the exact same binary can execute on 1 core, 4 cores, or 10,000 cores without recompilation.
3.2 Point-to-Point Communication: MPI_Send and MPI_Recv
Section titled “3.2 Point-to-Point Communication: MPI_Send and MPI_Recv”Point-to-point communication involves an exchange of data between exactly two processes: a sender and a receiver.
3.2.1 MPI_Send
Section titled “3.2.1 MPI_Send”int MPI_Send( void* msg_buf_p /* in: pointer to data buffer */, int msg_size /* in: number of elements to send */, MPI_Datatype msg_type /* in: datatype of elements */, int dest /* in: rank of destination process*/, int tag /* in: message identifier tag */, MPI_Comm communicator /* in: communication context */);- Message Content (
msg_buf_p,msg_size,msg_type):msg_buf_p: Pointer to the start of the data in memory.msg_size: The number of items (not necessarily bytes) being transmitted. For a C string, this includes the null terminator\0(strlen(greeting) + 1).msg_type: An MPI datatype corresponding to the C data type.
Table 3.1: Some Predefined MPI Datatypes
Section titled “Table 3.1: Some Predefined MPI Datatypes”| MPI Datatype | C Datatype |
|---|---|
MPI_CHAR | signed char |
MPI_SHORT | signed short int |
MPI_INT | signed int |
MPI_LONG | signed long int |
MPI_LONG_LONG | signed long long int |
MPI_UNSIGNED_CHAR | unsigned char |
MPI_UNSIGNED_SHORT | unsigned short int |
MPI_UNSIGNED | unsigned int |
MPI_UNSIGNED_LONG | unsigned long int |
MPI_FLOAT | float |
MPI_DOUBLE | double |
MPI_LONG_DOUBLE | long double |
MPI_BYTE | Raw uninterpreted byte (8 bits) |
MPI_PACKED | Serialized contiguous buffer |
- Message Envelope (
dest,tag,communicator):dest: Target process rank.tag: Non-negative integer used to categorize messages. For example, if a worker sends both diagnostic measurements and final calculation results, it can assigntag = 0for logs andtag = 1for calculations.communicator: The communication universe. Messages sent in one communicator cannot be received in another, preventing accidental message cross-talk between independent libraries or modules.
3.2.2 MPI_Recv
Section titled “3.2.2 MPI_Recv”int MPI_Recv( void* msg_buf_p /* out: pointer to receive buffer */, int buf_size /* in: maximum number of elements */, MPI_Datatype buf_type /* in: datatype of elements */, int source /* in: rank of sending process */, int tag /* in: message tag to match */, MPI_Comm communicator /* in: communication context */, MPI_Status* status_p /* out: status object or IGNORE */);buf_size: The capacity of the buffer (in number of elements). It must be large enough to hold the incoming message (buf_size >= msg_size), otherwise a buffer overflow error occurs.status_p: Returns information about the received message (orMPI_STATUS_IGNOREif details are not needed).
3.2.3 Message Matching Rules
Section titled “3.2.3 Message Matching Rules”For a message sent by process to be accepted by a receive called by process :
recv_comm == send_commdest == randsource == qrecv_tag == send_tag- The datatypes must be compatible and
recv_buf_size >= send_msg_size.
Wildcard Arguments
Section titled “Wildcard Arguments”A receiving process may not always know in advance which process will finish first, or which tag it will send. MPI provides wildcard constants for MPI_Recv:
MPI_ANY_SOURCE: Accepts incoming messages from any sender rank within the communicator.MPI_ANY_TAG: Accepts messages regardless of the tag.
/* Receive results in whichever order worker processes finish */for (int i = 1; i < comm_sz; i++) { MPI_Recv(result, result_sz, result_type, MPI_ANY_SOURCE, result_tag, comm, MPI_STATUS_IGNORE); Process_result(result);}3.2.4 Inspecting Messages with MPI_Status and MPI_Get_count
Section titled “3.2.4 Inspecting Messages with MPI_Status and MPI_Get_count”When wildcards are used, the receiver can determine who sent the message, what tag was used, and how much data was received by passing a pointer to an MPI_Status struct:
MPI_Status status;MPI_Recv(recv_buf, max_count, MPI_DOUBLE, MPI_ANY_SOURCE, MPI_ANY_TAG, comm, &status);
int sender = status.MPI_SOURCE;int tag = status.MPI_TAG;To determine the exact number of elements actually transmitted:
int count;MPI_Get_count(&status, MPI_DOUBLE, &count);printf("Received %d doubles from process %d\n", count, sender);3.3 Semantics and Protocols: Buffering vs. Blocking
Section titled “3.3 Semantics and Protocols: Buffering vs. Blocking”What occurs under the hood when MPI_Send is executed?
flowchart LR
A["Process q: MPI_Send"] --> B{"Message Size <= Cutoff?"}
B -- Yes --> C["Copy to MPI Internal Buffer"] --> D["MPI_Send Returns Immediately"]
B -- No --> E["Block until Network / Recv Ready"] --> F["Data Transferred"] --> G["MPI_Send Returns"]
- Buffering: If the message is smaller than an internal system cutoff threshold, MPI copies the message into a private system buffer.
MPI_Sendreturns immediately, even if the receiver has not yet calledMPI_Recv. - Blocking: If the message exceeds the threshold,
MPI_Sendblocks until the transmission has begun and the user’s send buffer can be safely reused.
The Non-Overtaking Rule
Section titled “The Non-Overtaking Rule”MPI guarantees that messages sent between the same pair of processes with matching tags are non-overtaking:
- If process sends message 1 and then message 2 to process , message 1 is guaranteed to be available to before message 2.
- However, there is no arrival ordering guarantee between different sending processes. If process and process both send messages to process , their arrival order depends on network latency and scheduling.
3.4 Potential Pitfalls: Deadlocks and Unsafe Programs
Section titled “3.4 Potential Pitfalls: Deadlocks and Unsafe Programs”A program is unsafe if its correct execution relies on the system automatically buffering messages.
Consider two processes attempting to exchange data:
/* DEADLOCK SCENARIO *//* Process 0 */MPI_Send(send_data, count, MPI_INT, 1, 0, comm);MPI_Recv(recv_data, count, MPI_INT, 1, 0, comm, MPI_STATUS_IGNORE);
/* Process 1 */MPI_Send(send_data, count, MPI_INT, 0, 0, comm);MPI_Recv(recv_data, count, MPI_INT, 0, 0, comm, MPI_STATUS_IGNORE);If count exceeds the system buffer threshold:
- Process 0 blocks inside
MPI_Send, waiting for Process 1 to callMPI_Recv. - Process 1 blocks inside
MPI_Send, waiting for Process 0 to callMPI_Recv. - Neither process ever reaches its
MPI_Recvcall. The program deadlocks (hangs indefinitely).
In later sections, we will explore safe communication patterns and combined operations like MPI_Sendrecv to systematically prevent deadlocks.