Collective Communication in MPI
In the trapezoidal rule program, we used point-to-point communication where each worker sent its partial integral to Process 0 in a sequential loop. While functional, this design suffers from serious scalability bottlenecks.
In this chapter, we explore collective communications—specialized functions that involve all processes within a communicator simultaneously, optimizing data movement across the cluster interconnect.
3.4 Tree-Structured Communication
Section titled “3.4 Tree-Structured Communication”In our naive point-to-point reduction, Process 0 performed sequential receives and additions:
- Time complexity on Process 0 was .
- Workers became idle immediately after transmitting their values.
The Binary Tree Reduction
Section titled “The Binary Tree Reduction”Instead of funneling all messages to Process 0, processes can pair up in parallel rounds, halving the active process count at each phase:
flowchart TD
subgraph Tree["Figure 3.6: Tree-Structured Global Sum (8 processes)"]
direction TB
subgraph P0["Phase 0: Initial Values"]
C0_0["P0: 5"]
C1_0["P1: 2"]
C2_0["P2: -1"]
C3_0["P3: -3"]
C4_0["P4: 6"]
C5_0["P5: 5"]
C6_0["P6: -7"]
C7_0["P7: 2"]
end
subgraph P1["Phase 1: 4 Parallel Adds"]
C0_1["P0: 7"]
C2_1["P2: -4"]
C4_1["P4: 11"]
C6_1["P6: -5"]
end
subgraph P2["Phase 2: 2 Parallel Adds"]
C0_2["P0: 3"]
C4_2["P4: 6"]
end
subgraph P3["Phase 3: Final Sum"]
C0_3["P0: 9"]
end
C0_0 --> C0_1
C1_0 -.->|Sends| C0_1
C2_0 --> C2_1
C3_0 -.->|Sends| C2_1
C4_0 --> C4_1
C5_0 -.->|Sends| C4_1
C6_0 --> C6_1
C7_0 -.->|Sends| C6_1
C0_1 --> C0_2
C2_1 -.->|Sends| C0_2
C4_1 --> C4_2
C6_1 -.->|Sends| C4_2
C0_2 --> C0_3
C4_2 -.->|Sends| C0_3
endComplexity Comparison
Section titled “Complexity Comparison”The tree reduction requires only phases:
- For : 3 phases vs. 7 sequential steps.
- For : phases vs. sequential steps (an acceleration factor of over 100x).
Alternative tree topologies (such as pairing Process with as shown in Figure 3.7) offer identical asymptotic complexity. Rather than requiring developers to manually handcraft and tune tree structures for specific network switches and topologies, MPI encapsulates global reductions inside library primitives.
3.5 MPI_Reduce
Section titled “3.5 MPI_Reduce”A function that aggregates data across all processes in a communicator into a single result is called a reduction.
int MPI_Reduce( void* input_data_p /* in: pointer to data to be combined */, void* output_data_p /* out: pointer to store result (dest) */, int count /* in: number of elements in buffer */, MPI_Datatype datatype /* in: MPI datatype of elements */, MPI_Op operator /* in: reduction operation */, int dest_process /* in: rank of process to receive res */, MPI_Comm comm /* in: communicator */);Predefined MPI Reduction Operators
Section titled “Predefined MPI Reduction Operators”| Operator | Meaning |
|---|---|
MPI_MAX | Maximum |
MPI_MIN | Minimum |
MPI_SUM | Sum |
MPI_PROD | Product |
MPI_LAND | Logical AND (&&) |
MPI_BAND | Bitwise AND (&) |
MPI_LOR | Logical OR (||) |
MPI_BOR | Bitwise OR (|) |
MPI_LXOR | Logical XOR |
MPI_BXOR | Bitwise XOR (^) |
MPI_MAXLOC | Maximum value and its rank location |
MPI_MINLOC | Minimum value and its rank location |
Using MPI_Reduce, the loops of point-to-point sends and receives from Program 3.2 condense into a single call:
MPI_Reduce(&local_int, &total_int, 1, MPI_DOUBLE, MPI_SUM, 0, MPI_COMM_WORLD);If count > 1, MPI_Reduce operates element-wise on arrays of vectors:
double local_x[N], sum[N];MPI_Reduce(local_x, sum, N, MPI_DOUBLE, MPI_SUM, 0, MPI_COMM_WORLD);3.6 Rules for Collective Communications
Section titled “3.6 Rules for Collective Communications”Collective communications differ fundamentally from point-to-point primitives:
- Global Participation: Every process in the communicator must call the collective function. Matching an
MPI_Reducecall on rank 0 with anMPI_Recvon rank 1 will cause the program to hang or crash. - Argument Compatibility: Arguments must be consistent. For example, all processes must pass the exact same
dest_processtoMPI_Reduce. - Buffer Arguments on Non-Destination Processes:
output_data_pis only used ondest_process, but all processes must still pass a valid pointer argument (orNULL). - No Tags — Matching by Sequence: Collective calls do not use tags. Collective calls are matched strictly by the communicator and the order in which they are invoked on each process.
Table 3.3: Multiple Collective Calls Matching by Invocation Order
Section titled “Table 3.3: Multiple Collective Calls Matching by Invocation Order”| Time | Process 0 | Process 1 | Process 2 |
|---|---|---|---|
| 0 | a = 1; c = 2; | a = 1; c = 2; | a = 1; c = 2; |
| 1 | MPI_Reduce(&a, &b, ...) | MPI_Reduce(&c, &d, ...) | MPI_Reduce(&a, &b, ...) |
| 2 | MPI_Reduce(&c, &d, ...) | MPI_Reduce(&a, &b, ...) | MPI_Reduce(&c, &d, ...) |
- Result on Process 0: Call 1 stores into variable ; Call 2 stores into variable .
3.7 MPI_Allreduce and the Butterfly Network
Section titled “3.7 MPI_Allreduce and the Butterfly Network”When every process in the communicator needs the result of a reduction (not just Process 0), we use MPI_Allreduce:
int MPI_Allreduce( void* input_data_p /* in: pointer to input data */, void* output_data_p /* out: pointer to store result */, int count /* in: number of elements */, MPI_Datatype datatype /* in: MPI datatype */, MPI_Op operator /* in: reduction operator */, MPI_Comm comm /* in: communicator */);Instead of running an MPI_Reduce followed by a broadcast, high-performance MPI libraries implement MPI_Allreduce using a Butterfly communication pattern:
flowchart LR
subgraph Butterfly["Figure 3.9: Butterfly Exchange Pattern (4 processes)"]
direction TB
subgraph Stage0["Stage 0: Local Values"]
P0_0["P0: v0"]
P1_0["P1: v1"]
P2_0["P2: v2"]
P3_0["P3: v3"]
end
subgraph Stage1["Stage 1: Distance-1 Exchange"]
P0_1["P0: v0+v1"]
P1_1["P1: v0+v1"]
P2_1["P2: v2+v3"]
P3_1["P3: v2+v3"]
end
subgraph Stage2["Stage 2: Distance-2 Exchange (Final Sum)"]
P0_2["P0: Total"]
P1_2["P1: Total"]
P2_2["P2: Total"]
P3_2["P3: Total"]
end
P0_0 <--> P1_0
P2_0 <--> P3_0
P0_0 --> P0_1
P1_0 --> P1_1
P2_0 --> P2_1
P3_0 --> P3_1
P0_1 <--> P2_1
P1_1 <--> P3_1
P0_1 --> P0_2
P1_1 --> P0_2
P2_1 --> P2_2
P3_1 --> P3_2
endIn rounds, every process exchanges partial sums and directly obtains the global aggregate with zero redundant communication.
3.8 Broadcast: MPI_Bcast
Section titled “3.8 Broadcast: MPI_Bcast”A broadcast transfers data belonging to a single source process to every process in the communicator:
int MPI_Bcast( void* data_p /* in/out: buffer to send/receive */, int count /* in: number of elements */, MPI_Datatype datatype /* in: MPI datatype */, int source_proc /* in: rank of root sender */, MPI_Comm comm /* in: communicator */);- On
source_proc,data_pis an input argument containing the data to distribute. - On all other processes,
data_pis an output argument that receives the incoming data.
Refactoring Get_input with MPI_Bcast
Section titled “Refactoring Get_input with MPI_Bcast”Using MPI_Bcast, Program 3.5 is simplified into Program 3.6:
void Get_input(int my_rank, int comm_sz, double* a_p, double* b_p, int* n_p) { if (my_rank == 0) { printf("Enter a, b, and n\n"); scanf("%lf %lf %d", a_p, b_p, n_p); } MPI_Bcast(a_p, 1, MPI_DOUBLE, 0, MPI_COMM_WORLD); MPI_Bcast(b_p, 1, MPI_DOUBLE, 0, MPI_COMM_WORLD); MPI_Bcast(n_p, 1, MPI_INT, 0, MPI_COMM_WORLD);}3.9 Data Partitioning: Block, Cyclic, and Block-Cyclic
Section titled “3.9 Data Partitioning: Block, Cyclic, and Block-Cyclic”Consider parallel vector addition:
To distribute an -element array across processes, three standard partitioning schemes are used:
Table 3.4: Partitions of a 12-component vector among 3 processes ()
Section titled “Table 3.4: Partitions of a 12-component vector among 3 processes (p=3,n=12p = 3, n = 12p=3,n=12)”| Process Rank | Block Partition | Cyclic Partition | Block-Cyclic Partition () |
|---|---|---|---|
| 0 | 0, 1, 2, 3 | 0, 3, 6, 9 | 0, 1, 6, 7 |
| 1 | 4, 5, 6, 7 | 1, 4, 7, 10 | 2, 3, 8, 9 |
| 2 | 8, 9, 10, 11 | 2, 5, 8, 11 | 4, 5, 10, 11 |
- Block Partition: Contiguous chunks of size are assigned to consecutive ranks. Component belongs to process .
- Cyclic Partition: Components are distributed round-robin. Component belongs to process .
- Block-Cyclic Partition: Blocks of fixed size are distributed round-robin.
Parallel Vector Addition Implementation
Section titled “Parallel Vector Addition Implementation”/* Program 3.8: Parallel vector addition function */void Parallel_vector_sum( double local_x[] /* in */, double local_y[] /* in */, double local_z[] /* out */, int local_n /* in */) { for (int local_i = 0; local_i < local_n; local_i++) { local_z[local_i] = local_x[local_i] + local_y[local_i]; }}3.10 Array Distribution Primitives: MPI_Scatter and MPI_Gather
Section titled “3.10 Array Distribution Primitives: MPI_Scatter and MPI_Gather”3.10.1 MPI_Scatter
Section titled “3.10.1 MPI_Scatter”MPI_Scatter splits a contiguous array residing on a root process and sends an equal slice to each process:
int MPI_Scatter( void* send_buf_p /* in: array on root process */, int send_count /* in: number of elements PER PROCESS */, MPI_Datatype send_type /* in: datatype of send elements */, void* recv_buf_p /* out: local receiving buffer */, int recv_count /* in: number of elements received */, MPI_Datatype recv_type /* in: datatype of receive elements */, int src_proc /* in: rank of root process */, MPI_Comm comm /* in: communicator */);send_count Semantics
send_count specifies the number of items delivered to each individual process, not the total size of send_buf_p.
/* Program 3.9: Reading and scattering a vector */void Read_vector(double local_a[], int local_n, int n, char vec_name[], int my_rank, MPI_Comm comm) { double* a = NULL; if (my_rank == 0) { a = malloc(n * sizeof(double)); printf("Enter the vector %s\n", vec_name); for (int i = 0; i < n; i++) scanf("%lf", &a[i]); MPI_Scatter(a, local_n, MPI_DOUBLE, local_a, local_n, MPI_DOUBLE, 0, comm); free(a); } else { MPI_Scatter(a, local_n, MPI_DOUBLE, local_a, local_n, MPI_DOUBLE, 0, comm); }}3.10.2 MPI_Gather
Section titled “3.10.2 MPI_Gather”MPI_Gather is the inverse of MPI_Scatter. It collects local slices from each process and concatenates them in rank order on a root process:
int MPI_Gather( void* send_buf_p /* in: local buffer */, int send_count /* in: elements sent by each process */, MPI_Datatype send_type /* in: datatype of send elements */, void* recv_buf_p /* out: destination buffer on root */, int recv_count /* in: elements received PER PROCESS */, MPI_Datatype recv_type /* in: datatype of receive elements */, int dest_proc /* in: rank of receiving root process */, MPI_Comm comm /* in: communicator */);/* Program 3.10: Gathering and printing a distributed vector */void Print_vector(double local_b[], int local_n, int n, char title[], int my_rank, MPI_Comm comm) { double* b = NULL; if (my_rank == 0) { b = malloc(n * sizeof(double)); MPI_Gather(local_b, local_n, MPI_DOUBLE, b, local_n, MPI_DOUBLE, 0, comm); printf("%s\n", title); for (int i = 0; i < n; i++) printf("%f ", b[i]); printf("\n"); free(b); } else { MPI_Gather(local_b, local_n, MPI_DOUBLE, b, local_n, MPI_DOUBLE, 0, comm); }}3.11 Parallel Matrix-Vector Multiplication with MPI_Allgather
Section titled “3.11 Parallel Matrix-Vector Multiplication with MPI_Allgather”Matrix-vector multiplication computes , where is an matrix and is an -dimensional vector:
In C, 2D matrices are commonly mapped to a contiguous 1D array in row-major order: element is stored at offset .
flowchart TD
subgraph MatVec["Figure 3.11: Matrix-Vector Multiplication Partitioning"]
direction LR
subgraph Matrix["Matrix A (m x n)"]
R0["Process 0: Rows 0 ... m/p - 1"]
R1["Process 1: Rows m/p ... 2m/p - 1"]
Rn["Process p-1: Last block of rows"]
end
subgraph Vector["Vector x"]
X["Allgather ensures full vector x on every rank"]
end
subgraph Output["Output Vector y"]
Y0["Local y calculated by P0"]
Y1["Local y calculated by P1"]
Yn["Local y calculated by P(p-1)"]
end
Matrix --> Output
Vector --> Output
endTo compute , Process requires row of and the entire vector . Since each process initially holds only a local block of (local_x), processes call MPI_Allgather to assemble the complete vector across all ranks before multiplying:
int MPI_Allgather( void* send_buf_p /* in: local slice */, int send_count /* in: elements sent per process */, MPI_Datatype send_type /* in: send datatype */, void* recv_buf_p /* out: buffer receiving all slices */, int recv_count /* in: elements received PER PROCESS */, MPI_Datatype recv_type /* in: receive datatype */, MPI_Comm comm /* in: communicator */);Complete Implementation: Mat_vect_mult
Section titled “Complete Implementation: Mat_vect_mult”/* Program 3.12: Parallel Matrix-Vector Multiplication */void Mat_vect_mult( double local_A[] /* in: local rows of matrix A */, double local_x[] /* in: local portion of x */, double local_y[] /* out: local portion of y */, int local_m /* in: rows per process */, int n /* in: total columns */, int local_n /* in: components of x per p */, MPI_Comm comm /* in: communicator */) { double* x; int local_i, j;
/* Allocate buffer for full vector x */ x = malloc(n * sizeof(double));
/* Concatenate all local_x pieces into complete vector x on all processes */ MPI_Allgather(local_x, local_n, MPI_DOUBLE, x, local_n, MPI_DOUBLE, comm);
/* Local matrix-vector multiplication */ for (local_i = 0; local_i < local_m; local_i++) { local_y[local_i] = 0.0; for (j = 0; j < n; j++) { local_y[local_i] += local_A[local_i * n + j] * x[j]; } }
free(x);} /* Mat_vect_mult */