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.
3.12 MPI-Derived Datatypes
Section titled “3.12 MPI-Derived Datatypes”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.
Data Consolidation Techniques
Section titled “Data Consolidation Techniques”MPI provides three primary mechanisms to consolidate non-trivial data layouts:
- The
countargument: Groups contiguous elements of identical primitive types into a single buffer. - Derived Datatypes: Represents arbitrary collections of heterogeneous data types and non-contiguous memory locations.
MPI_PackandMPI_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:
Suppose variables a, b, and n are arranged in memory on Process 0 as follows:
| Variable | C Type | Memory Address | Byte Displacement from &a |
|---|---|---|---|
a | double | 24 | 0 |
b | double | 40 | 16 |
n | int | 48 | 24 |
The corresponding derived datatype is:
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 */);Committing and Freeing Types
Section titled “Committing and Freeing Types”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);3.12.2 Complete Example: Program 3.13
Section titled “3.12.2 Complete Example: Program 3.13”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.
3.13.1 Wall-Clock Time vs. CPU Time
Section titled “3.13.1 Wall-Clock Time vs. CPU Time”- 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:
- Synchronize all processes before timing with
MPI_Barrier(). - Record
local_start = MPI_Wtime(). - Execute the code to be benchmarked (excluding disk/console I/O).
- Record
local_finish = MPI_Wtime(). - 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 and process counts :
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 () | |||||
|---|---|---|---|---|---|
| 1 | 4.1 | 16.0 | 64.0 | 270 | 1100 |
| 2 | 2.3 | 8.5 | 33.0 | 140 | 560 |
| 4 | 2.0 | 5.1 | 18.0 | 70 | 280 |
| 8 | 1.7 | 3.3 | 9.8 | 36 | 140 |
| 16 | 1.7 | 2.6 | 5.9 | 19 | 71 |
Mathematical Performance Model
Section titled “Mathematical Performance Model”For an matrix, the sequential algorithm executes multiplications and additions per row across rows:
In the parallel algorithm, each process computes only rows, executing floating-point operations. However, each process must also exchange its vector slice via MPI_Allgather:
- For large and small : The computation term dominates. Doubling roughly halves the total execution time.
- For small and large (e.g., ): The communication overhead dominates. Increasing yields diminishing returns or plateaus.
3.15 Speedup, Efficiency, and Scalability
Section titled “3.15 Speedup, Efficiency, and Scalability”Speedup
Section titled “Speedup”Speedup () measures how much faster the parallel program runs compared to the serial baseline:
- Linear Speedup: . Running on cores results in an exact -fold acceleration.
Table 3.6: Speedup of Parallel Matrix-Vector Multiplication
Section titled “Table 3.6: Speedup of Parallel Matrix-Vector Multiplication”comm_sz () | |||||
|---|---|---|---|---|---|
| 1 | 1.0 | 1.0 | 1.0 | 1.0 | 1.0 |
| 2 | 1.8 | 1.9 | 1.9 | 1.9 | 2.0 |
| 4 | 2.1 | 3.1 | 3.6 | 3.9 | 3.9 |
| 8 | 2.4 | 4.8 | 6.5 | 7.5 | 7.9 |
| 16 | 2.4 | 6.2 | 10.8 | 14.2 | 15.5 |
Parallel Efficiency
Section titled “Parallel Efficiency”Efficiency () measures the fraction of time cores spend performing useful work:
Linear speedup corresponds to . In practice, due to communication and synchronization overhead, efficiency is generally .
Table 3.7: Efficiency of Parallel Matrix-Vector Multiplication
Section titled “Table 3.7: Efficiency of Parallel Matrix-Vector Multiplication”comm_sz () | |||||
|---|---|---|---|---|---|
| 1 | 1.00 | 1.00 | 1.00 | 1.00 | 1.00 |
| 2 | 0.89 | 0.94 | 0.97 | 0.96 | 0.98 |
| 4 | 0.51 | 0.78 | 0.89 | 0.96 | 0.98 |
| 8 | 0.30 | 0.61 | 0.82 | 0.94 | 0.98 |
| 16 | 0.15 | 0.39 | 0.68 | 0.89 | 0.97 |
Scalability: Strong vs. Weak
Section titled “Scalability: Strong vs. Weak”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 is increased while keeping the problem size fixed.
- Weak Scalability: The program maintains constant efficiency when the problem size is increased at the same rate as the number of processes (i.e., fixed problem size per core).
As demonstrated in Table 3.7, parallel matrix-vector multiplication is weakly scalable: increasing both and by a factor of 2 maintains near-ideal efficiencies ().