Skip to content

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
  end
flowchart 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"]
  end

To 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 pp processes, the ranks are numbered consecutively:

0,1,2,…,p−10, 1, 2, \dots, p - 1

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 */

Most MPI distributions (such as MPICH or OpenMPI) provide a compiler wrapper called mpicc:

Terminal window
$ mpicc -g -Wall -o mpi_hello mpi_hello.c

A 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:

Terminal window
$ mpiexec -n <number_of_processes> ./mpi_hello

For example, running with a single process:

Terminal window
$ mpiexec -n 1 ./mpi_hello
Greetings from process 0 of 1!

Running with 4 processes:

Terminal window
$ mpiexec -n 4 ./mpi_hello
Greetings 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:

  1. The operating system starts 4 separate instances (processes) of the mpi_hello executable.
  2. The runtime assigns each instance a distinct rank (0,1,2,30, 1, 2, 3).
  3. The MPI implementation maps processes to CPU cores and establishes communication channels across the interconnect.

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).

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, NULL can be passed for both parameters.
  • MPI_Finalize: Cleans up internal structures and frees resources. No MPI calls should be made after MPI_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_sz receives the total number of processes in comm.
  • my_rank receives the integer ID of the calling process (0≤my_rank<comm_sz0 \le \text{my\_rank} < \text{comm\_sz}).

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.

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 */
);
  1. 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.
MPI DatatypeC Datatype
MPI_CHARsigned char
MPI_SHORTsigned short int
MPI_INTsigned int
MPI_LONGsigned long int
MPI_LONG_LONGsigned long long int
MPI_UNSIGNED_CHARunsigned char
MPI_UNSIGNED_SHORTunsigned short int
MPI_UNSIGNEDunsigned int
MPI_UNSIGNED_LONGunsigned long int
MPI_FLOATfloat
MPI_DOUBLEdouble
MPI_LONG_DOUBLElong double
MPI_BYTERaw uninterpreted byte (8 bits)
MPI_PACKEDSerialized contiguous buffer
  1. 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 assign tag = 0 for logs and tag = 1 for 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.

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 (or MPI_STATUS_IGNORE if details are not needed).

For a message sent by process qq to be accepted by a receive called by process rr:

  1. recv_comm == send_comm
  2. dest == r and source == q
  3. recv_tag == send_tag
  4. The datatypes must be compatible and recv_buf_size >= send_msg_size.

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_Send returns immediately, even if the receiver has not yet called MPI_Recv.
  • Blocking: If the message exceeds the threshold, MPI_Send blocks until the transmission has begun and the user’s send buffer can be safely reused.

MPI guarantees that messages sent between the same pair of processes with matching tags are non-overtaking:

  • If process qq sends message 1 and then message 2 to process rr, message 1 is guaranteed to be available to rr before message 2.
  • However, there is no arrival ordering guarantee between different sending processes. If process qq and process tt both send messages to process rr, 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:

  1. Process 0 blocks inside MPI_Send, waiting for Process 1 to call MPI_Recv.
  2. Process 1 blocks inside MPI_Send, waiting for Process 0 to call MPI_Recv.
  3. Neither process ever reaches its MPI_Recv call. The program deadlocks (hangs indefinitely).

In later sections, we will explore safe communication patterns and combined operations like MPI_Sendrecv to systematically prevent deadlocks.