Skip to content

Derived Datatypes and Performance Evaluation

In distributed-memory computing, communication across the network interconnect is several orders of magnitude slower than local CPU operations. While a floating-point calculation on modern hardware takes a fraction of a nanosecond, transferring a packet over a network introduces microsecond-level latency.

To maximize parallel efficiency, programs must minimize both the total volume of transmitted data and—crucially—the total number of individual messages sent.


Sending multiple small messages incurs network header and latency overhead on every transfer. Consider transmitting an array of 1,000 doubles:

/* UNOPTIMIZED: 1,000 separate messages */
if (my_rank == 0) {
for (int i = 0; i < 1000; i++)
MPI_Send(&x[i], 1, MPI_DOUBLE, 1, 0, comm);
} else {
for (int i = 0; i < 1000; i++)
MPI_Recv(&x[i], 1, MPI_DOUBLE, 0, 0, comm, &status);
}
/* OPTIMIZED: 1 consolidated message */
if (my_rank == 0) {
MPI_Send(x, 1000, MPI_DOUBLE, 1, 0, comm);
} else {
MPI_Recv(x, 1000, MPI_DOUBLE, 0, 0, comm, &status);
}

In empirical cluster benchmarks, the unoptimized loop of sends takes between 50 to 100 times longer than the single consolidated send.

MPI provides three primary mechanisms to consolidate non-trivial data layouts:

  1. The count argument: Groups contiguous elements of identical primitive types into a single buffer.
  2. Derived Datatypes: Represents arbitrary collections of heterogeneous data types and non-contiguous memory locations.
  3. MPI_Pack and MPI_Unpack: Manually marshals heterogeneous data into contiguous byte buffers before transmission.

3.12.1 Building a Struct Type: MPI_Type_create_struct

Section titled “3.12.1 Building a Struct Type: MPI_Type_create_struct”

In our trapezoidal rule program, Process 0 broadcast three separate input variables: two double values (a, b) and one int (n). Rather than issuing three distinct MPI_Bcast calls, we can bundle them into a single custom MPI datatype.

Formally, an MPI derived datatype is defined as a type map: a sequence of basic datatypes paired with their byte displacements from the start of the memory block:

Type Map={(type0,disp0),(type1,disp1),…,(typek−1,dispk−1)}\text{Type Map} = \{ (\text{type}_0, \text{disp}_0), (\text{type}_1, \text{disp}_1), \dots, (\text{type}_{k-1}, \text{disp}_{k-1}) \}

Suppose variables a, b, and n are arranged in memory on Process 0 as follows:

VariableC TypeMemory AddressByte Displacement from &a
adouble240
bdouble4016
nint4824

The corresponding derived datatype is:

{(MPI_DOUBLE,0),(MPI_DOUBLE,16),(MPI_INT,24)}\{ (\text{MPI\_DOUBLE}, 0), (\text{MPI\_DOUBLE}, 16), (\text{MPI\_INT}, 24) \}

To construct this type in C:

int MPI_Type_create_struct(
int count /* in: number of fields in structure */,
int array_of_blocklengths[] /* in: number of elements per block */,
MPI_Aint array_of_displacements[] /* in: byte offset of each block */,
MPI_Datatype array_of_types[] /* in: basic datatype of each block */,
MPI_Datatype* new_type_p /* out: newly created MPI datatype */
);

Calculating Byte Displacements with MPI_Get_address

Section titled “Calculating Byte Displacements with MPI_Get_address”

Because compilers insert structure padding and memory alignment according to machine architecture, memory offsets should never be hardcoded. Instead, use MPI_Get_address and the address integer type MPI_Aint:

int MPI_Get_address(
void* location_p /* in: pointer to variable or field */,
MPI_Aint* address_p /* out: memory address as integer */
);

Before a derived datatype can be passed to communication routines, it must be committed:

int MPI_Type_commit(MPI_Datatype* new_mpi_t_p);

When no longer needed, internal resources are freed using:

int MPI_Type_free(MPI_Datatype* old_mpi_t_p);

Below is the complete implementation of Build_mpi_type and the refactored Get_input function:

/* Program 3.13: Get_input with a derived datatype */
void Build_mpi_type(
double* a_p /* in */,
double* b_p /* in */,
int* n_p /* in */,
MPI_Datatype* input_mpi_t_p /* out */
) {
int array_of_blocklengths[3] = {1, 1, 1};
MPI_Datatype array_of_types[3] = {MPI_DOUBLE, MPI_DOUBLE, MPI_INT};
MPI_Aint a_addr, b_addr, n_addr;
MPI_Aint array_of_displacements[3] = {0};
/* Query actual physical addresses in memory */
MPI_Get_address(a_p, &a_addr);
MPI_Get_address(b_p, &b_addr);
MPI_Get_address(n_p, &n_addr);
/* Compute byte offsets relative to the start address */
array_of_displacements[1] = b_addr - a_addr;
array_of_displacements[2] = n_addr - a_addr;
/* Build and commit the new composite type */
MPI_Type_create_struct(3, array_of_blocklengths,
array_of_displacements, array_of_types,
input_mpi_t_p);
MPI_Type_commit(input_mpi_t_p);
} /* Build_mpi_type */
void Get_input(int my_rank, int comm_sz, double* a_p, double* b_p, int* n_p) {
MPI_Datatype input_mpi_t;
Build_mpi_type(a_p, b_p, n_p, &input_mpi_t);
if (my_rank == 0) {
printf("Enter a, b, and n\n");
scanf("%lf %lf %d", a_p, b_p, n_p);
}
/* Single broadcast distributes all three variables at once */
MPI_Bcast(a_p, 1, input_mpi_t, 0, MPI_COMM_WORLD);
MPI_Type_free(&input_mpi_t);
} /* Get_input */

3.13 Performance Evaluation of MPI Programs

Section titled “3.13 Performance Evaluation of MPI Programs”

We write parallel programs to achieve faster run times than sequential programs. To quantify parallel performance, we must measure execution time accurately.

  • CPU time (such as measured by the C standard function clock()): Records only the time spent executing user instructions on the CPU core. It omits idle time spent blocking on I/O or waiting for incoming messages across the network.
  • Wall-clock time: Measures the total elapsed time between start and finish, including communication latency and synchronization delays.

MPI provides a portable high-resolution wall-clock timer:

double MPI_Wtime(void);

Accurate Benchmarking with MPI_Barrier and MPI_Reduce

Section titled “Accurate Benchmarking with MPI_Barrier and MPI_Reduce”

Because processes may not start simultaneously and finish at different times, the parallel run time is governed by the slowest process.

To obtain an accurate, reproducible measurement:

  1. Synchronize all processes before timing with MPI_Barrier().
  2. Record local_start = MPI_Wtime().
  3. Execute the code to be benchmarked (excluding disk/console I/O).
  4. Record local_finish = MPI_Wtime().
  5. Compute the global parallel execution time by taking the maximum elapsed time across all ranks using MPI_Reduce(..., MPI_MAX).
double local_start, local_finish, local_elapsed, elapsed;
/* Synchronize processes so they begin the measured block together */
MPI_Barrier(comm);
local_start = MPI_Wtime();
/* Parallel computational kernel (e.g. Mat_vect_mult) */
local_finish = MPI_Wtime();
local_elapsed = local_finish - local_start;
/* Determine the elapsed time of the slowest process */
MPI_Reduce(&local_elapsed, &elapsed, 1, MPI_DOUBLE, MPI_MAX, 0, comm);
if (my_rank == 0) {
printf("Elapsed parallel time = %e seconds\n", elapsed);
}

3.14 Empirical Results: Parallel Matrix-Vector Multiplication

Section titled “3.14 Empirical Results: Parallel Matrix-Vector Multiplication”

The table below presents timing results (in milliseconds) for parallel matrix-vector multiplication across varying square matrix sizes n×nn \times n and process counts pp:

Table 3.5: Run-times of Serial and Parallel Matrix-Vector Multiplication (ms)

Section titled “Table 3.5: Run-times of Serial and Parallel Matrix-Vector Multiplication (ms)”
comm_sz (pp)n=1024n = 1024n=2048n = 2048n=4096n = 4096n=8192n = 8192n=16384n = 16384
14.116.064.02701100
22.38.533.0140560
42.05.118.070280
81.73.39.836140
161.72.65.91971

For an n×nn \times n matrix, the sequential algorithm executes nn multiplications and nn additions per row across nn rows:

Tserial(n)≈2n2 FLOPs  ⟹  Tserial(n)≈an2T_{\text{serial}}(n) \approx 2n^2 \text{ FLOPs} \implies T_{\text{serial}}(n) \approx a n^2

In the parallel algorithm, each process computes only n/pn / p rows, executing n2/pn^2 / p floating-point operations. However, each process must also exchange its vector slice via MPI_Allgather:

Tparallel(n,p)=Tserial(n)p+Tallgather(n,p)T_{\text{parallel}}(n, p) = \frac{T_{\text{serial}}(n)}{p} + T_{\text{allgather}}(n, p)

  • For large nn and small pp: The computation term Tserial(n)/pT_{\text{serial}}(n) / p dominates. Doubling pp roughly halves the total execution time.
  • For small nn and large pp (e.g., n=1024,p=16n = 1024, p = 16): The communication overhead TallgatherT_{\text{allgather}} dominates. Increasing pp yields diminishing returns or plateaus.

Speedup (SS) measures how much faster the parallel program runs compared to the serial baseline:

S(n,p)=Tserial(n)Tparallel(n,p)S(n, p) = \frac{T_{\text{serial}}(n)}{T_{\text{parallel}}(n, p)}

  • Linear Speedup: S(n,p)=pS(n, p) = p. Running on pp cores results in an exact pp-fold acceleration.

Table 3.6: Speedup of Parallel Matrix-Vector Multiplication

Section titled “Table 3.6: Speedup of Parallel Matrix-Vector Multiplication”
comm_sz (pp)n=1024n = 1024n=2048n = 2048n=4096n = 4096n=8192n = 8192n=16384n = 16384
11.01.01.01.01.0
21.81.91.91.92.0
42.13.13.63.93.9
82.44.86.57.57.9
162.46.210.814.215.5

Efficiency (EE) measures the fraction of time cores spend performing useful work:

E(n,p)=S(n,p)p=Tserial(n)p×Tparallel(n,p)E(n, p) = \frac{S(n, p)}{p} = \frac{T_{\text{serial}}(n)}{p \times T_{\text{parallel}}(n, p)}

Linear speedup corresponds to E=1.0E = 1.0. In practice, due to communication and synchronization overhead, efficiency is generally <1.0< 1.0.

Table 3.7: Efficiency of Parallel Matrix-Vector Multiplication

Section titled “Table 3.7: Efficiency of Parallel Matrix-Vector Multiplication”
comm_sz (pp)n=1024n = 1024n=2048n = 2048n=4096n = 4096n=8192n = 8192n=16384n = 16384
11.001.001.001.001.00
20.890.940.970.960.98
40.510.780.890.960.98
80.300.610.820.940.98
160.150.390.680.890.97

A parallel program is defined as scalable if efficiency can be maintained as the number of processes increases:

  • Strong Scalability: The program maintains constant efficiency when the number of processes pp is increased while keeping the problem size nn fixed.
  • Weak Scalability: The program maintains constant efficiency when the problem size nn is increased at the same rate as the number of processes pp (i.e., fixed problem size per core).

As demonstrated in Table 3.7, parallel matrix-vector multiplication is weakly scalable: increasing both nn and pp by a factor of 2 maintains near-ideal efficiencies (E≈0.98E \approx 0.98).