Skip to content

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.


In our naive point-to-point reduction, Process 0 performed p−1p - 1 sequential receives and additions:

  • Time complexity on Process 0 was O(p)\mathcal{O}(p).
  • Workers 1,2,…,p−11, 2, \dots, p - 1 became idle immediately after transmitting their values.

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
  end

The tree reduction requires only ⌈log⁡2(p)⌉\lceil \log_2(p) \rceil phases:

  • For p=8p = 8: 3 phases vs. 7 sequential steps.
  • For p=1024p = 1024: 1010 phases vs. 10231023 sequential steps (an acceleration factor of over 100x).

Alternative tree topologies (such as pairing Process ii with i+p/2i + p/2 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.


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 */
);
OperatorMeaning
MPI_MAXMaximum
MPI_MINMinimum
MPI_SUMSum
MPI_PRODProduct
MPI_LANDLogical AND (&&)
MPI_BANDBitwise AND (&)
MPI_LORLogical OR (||)
MPI_BORBitwise OR (|)
MPI_LXORLogical XOR
MPI_BXORBitwise XOR (^)
MPI_MAXLOCMaximum value and its rank location
MPI_MINLOCMinimum value and its rank location

Using MPI_Reduce, the p−1p-1 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);

Collective communications differ fundamentally from point-to-point primitives:

  1. Global Participation: Every process in the communicator must call the collective function. Matching an MPI_Reduce call on rank 0 with an MPI_Recv on rank 1 will cause the program to hang or crash.
  2. Argument Compatibility: Arguments must be consistent. For example, all processes must pass the exact same dest_process to MPI_Reduce.
  3. Buffer Arguments on Non-Destination Processes: output_data_p is only used on dest_process, but all processes must still pass a valid pointer argument (or NULL).
  4. 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”
TimeProcess 0Process 1Process 2
0a = 1; c = 2;a = 1; c = 2;a = 1; c = 2;
1MPI_Reduce(&a, &b, ...)MPI_Reduce(&c, &d, ...)MPI_Reduce(&a, &b, ...)
2MPI_Reduce(&c, &d, ...)MPI_Reduce(&a, &b, ...)MPI_Reduce(&c, &d, ...)
  • Result on Process 0: Call 1 stores 1+2+1=41 + 2 + 1 = 4 into variable bb; Call 2 stores 2+1+2=52 + 1 + 2 = 5 into variable dd.

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
  end

In log⁡2(p)\log_2(p) rounds, every process exchanges partial sums and directly obtains the global aggregate with zero redundant communication.


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_p is an input argument containing the data to distribute.
  • On all other processes, data_p is an output argument that receives the incoming data.

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:

z=x+y  ⟹  zi=xi+yifor 0≤i<n\mathbf{z} = \mathbf{x} + \mathbf{y} \implies z_i = x_i + y_i \quad \text{for } 0 \le i < n

To distribute an nn-element array across pp processes, three standard partitioning schemes are used:

Table 3.4: Partitions of a 12-component vector among 3 processes (p=3,n=12p = 3, n = 12)

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 RankBlock PartitionCyclic PartitionBlock-Cyclic Partition (b=2b = 2)
00, 1, 2, 30, 3, 6, 90, 1, 6, 7
14, 5, 6, 71, 4, 7, 102, 3, 8, 9
28, 9, 10, 112, 5, 8, 114, 5, 10, 11
  • Block Partition: Contiguous chunks of size n/pn / p are assigned to consecutive ranks. Component ii belongs to process ⌊i/local_n⌋\lfloor i / \text{local\_n} \rfloor.
  • Cyclic Partition: Components are distributed round-robin. Component ii belongs to process i(modp)i \pmod p.
  • Block-Cyclic Partition: Blocks of fixed size bb are distributed round-robin.
/* 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”

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);
}
}

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 y=Ax\mathbf{y} = A\mathbf{x}, where AA is an m×nm \times n matrix and x\mathbf{x} is an nn-dimensional vector:

yi=∑j=0n−1Ai,j⋅xj=Ai,0x0+Ai,1x1+⋯+Ai,n−1xn−1y_i = \sum_{j=0}^{n-1} A_{i, j} \cdot x_j = A_{i, 0}x_0 + A_{i, 1}x_1 + \dots + A_{i, n-1}x_{n-1}

In C, 2D matrices are commonly mapped to a contiguous 1D array in row-major order: element (i,j)(i, j) is stored at offset i×n+ji \times n + j.

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
  end

To compute yiy_i, Process qq requires row ii of AA and the entire vector x\mathbf{x}. Since each process initially holds only a local block of x\mathbf{x} (local_x), processes call MPI_Allgather to assemble the complete vector x\mathbf{x} 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 */
);
/* 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 */