How unit 5 is examined
This unit covers MPI message passing on distributed-memory machines; the marks sit in non-blocking communication and in reducing communication overhead (both 7 marks, May 2022).
Message passing
<span style="display:inline-block;padding:.16em .6em;border:1.5px solid currentColor;border-radius:999px;font-size:.68em;font-weight:700;letter-spacing:.06em;text-transform:uppercase;opacity:.75">Not asked since 2022</span>
Definition. <mark>Message passing is a programming model in which processes with separate private memories cooperate only by explicitly sending and receiving messages.</mark>
Key points.
- Each process has its own address space, so no variable is shared and data moves only through send and receive calls.
- MPI (Message Passing Interface) is the standard library specification for this model, with bindings for C and Fortran.
- All processes run the same program (SPMD) and are told apart by their rank, from
MPI_Comm_rank, within a communicator such asMPI_COMM_WORLD. MPI_Initstarts the environment,MPI_Comm_sizegives the number of processes, andMPI_Finalizeends it.- It scales to clusters because memory is not shared, but the programmer must place the data and the communication by hand.
Message and point to point communication
<span style="display:inline-block;padding:.16em .6em;border:1.5px solid currentColor;border-radius:999px;font-size:.68em;font-weight:700;letter-spacing:.06em;text-transform:uppercase;opacity:.75">Not asked since 2022</span>
Definition. <mark>Point-to-point communication is the transfer of a message between exactly one sender and one receiver, using MPI_Send and MPI_Recv.</mark>
Key points.
- A call is
MPI_Send(buf, count, datatype, dest, tag, comm)andMPI_Recv(buf, count, datatype, source, tag, comm, &status). - A message is identified by its source, destination, tag and communicator, and the receive matches on these.
MPI_ANY_SOURCEandMPI_ANY_TAGallow a receive to match any sender or tag, andstatusreports what actually arrived.MPI_SendandMPI_Recvare blocking: the call returns only when its buffer is safe to reuse.- Two processes that both send first and receive later can deadlock, so one must receive first or use
MPI_Sendrecv.
collective communication
<span style="display:inline-block;padding:.16em .6em;border:1.5px solid currentColor;border-radius:999px;font-size:.68em;font-weight:700;letter-spacing:.06em;text-transform:uppercase;opacity:.75">Not asked since 2022</span>
Definition. <mark>Collective communication is an operation in which all processes of a communicator take part in one call, such as broadcast, scatter, gather or reduce.</mark>
Key points.
MPI_Bcastcopies data from a root process to all others, andMPI_Scattersplits the root's array into pieces, one per process.MPI_Gathercollects one piece from each process at the root, andMPI_Allgathergives the result to everyone.MPI_Reducecombines values with an operation (MPI_SUM,MPI_MAX) and delivers the result to the root;MPI_Allreducedelivers it to all.MPI_Barriermakes every process wait until all have reached it.- Every process in the communicator must call the collective, and implementations use tree algorithms, so it is faster than hand-written send loops.
non blocking point-to-point communication
<span style="display:inline-block;padding:.16em .6em;border:1.5px solid currentColor;border-radius:999px;font-size:.68em;font-weight:700;letter-spacing:.06em;text-transform:uppercase;opacity:.75">Low weight</span>
Definition. <mark>Non-blocking communication (MPI_Isend, MPI_Irecv) returns immediately and lets computation overlap the transfer; completion is checked later with MPI_Wait or MPI_Test.</mark>
| Point | Blocking (MPI_Send/Recv) |
Non-blocking (MPI_Isend/Irecv) |
|---|---|---|
| Return | After buffer is safe to reuse | Immediately, with an MPI_Request |
| Overlap | None, process idles | Computation overlaps communication |
| Completion | Implicit | MPI_Wait or MPI_Test |
| Buffer | Reusable on return | Must not be touched until completed |
| Deadlock | Likely if both send first | Avoided |
| Use | Simple exchange | Halo exchange while computing interior |
Asked: [7 marks] (May 2022) What is non-blocking point-to-point communication? How it differs from blocking communication?
virtual topologies
<span style="display:inline-block;padding:.16em .6em;border:1.5px solid currentColor;border-radius:999px;font-size:.68em;font-weight:700;letter-spacing:.06em;text-transform:uppercase;opacity:.75">Not asked since 2022</span>
Definition. <mark>A virtual topology is a logical arrangement of processes, such as a Cartesian grid or graph, that maps the program's communication pattern onto ranks.</mark>
Key points.
MPI_Cart_createbuilds a grid communicator with given dimensions and periodicity.MPI_Cart_coordsandMPI_Cart_rankconvert between rank and grid coordinates.MPI_Cart_shiftfinds the neighbour ranks in a direction, which suits stencil and halo exchange.- It makes code clearer and lets the system reorder ranks to match the physical network.
MPI performance tools
<span style="display:inline-block;padding:.16em .6em;border:1.5px solid currentColor;border-radius:999px;font-size:.68em;font-weight:700;letter-spacing:.06em;text-transform:uppercase;opacity:.75">Not asked since 2022</span>
Definition. <mark>MPI performance tools measure and visualise where an MPI program spends time in computation, communication and waiting.</mark>
Key points.
MPI_Wtimegives wall-clock time for manual timing of code sections.- Tracing tools such as Vampir, Scalasca, TAU and Intel Trace Analyzer record events and show timelines.
- Profilers such as mpiP give a summary of time per MPI call and per rank.
- They reveal load imbalance, late senders and long waits, which guide optimisation.
communication parameters
<span style="display:inline-block;padding:.16em .6em;border:1.5px solid currentColor;border-radius:999px;font-size:.68em;font-weight:700;letter-spacing:.06em;text-transform:uppercase;opacity:.75">Not asked since 2022</span>
Definition. ==Communication time of a message is modelled by latency and bandwidth: $T = T_{lat} + n/B$.==
Key points.
- Latency $T_{lat}$ is the fixed start-up time per message, independent of its size.
- Bandwidth $B$ is the data rate of the link, so $n/B$ is the transfer time of $n$ bytes.
- Small messages are latency-dominated and large messages are bandwidth-dominated.
- Effective bandwidth $n/T$ approaches $B$ only for large $n$.
impact of synchronizations sterilizations and contentions
<span style="display:inline-block;padding:.16em .6em;border:1.5px solid currentColor;border-radius:999px;font-size:.68em;font-weight:700;letter-spacing:.06em;text-transform:uppercase;opacity:.75">Not asked since 2022</span>
Definition. <mark>Synchronization forces processes to wait for each other, and contention occurs when several messages compete for the same link, both adding to communication time.</mark>
Key points.
- Barriers and blocking calls make fast processes idle until the slowest arrives, so load imbalance becomes lost time.
- Serialization happens when a step must run one process at a time, which limits speedup.
- Contention on shared links or a single root (as in gather) reduces the bandwidth each message gets.
- Fewer barriers, balanced work and spread-out communication reduce these effects.
reductions in communication overhead
<span style="display:inline-block;padding:.16em .6em;border:1.5px solid currentColor;border-radius:999px;font-size:.68em;font-weight:700;letter-spacing:.06em;text-transform:uppercase;opacity:.75">Low weight</span>
Definition. <mark>Communication overhead is the time spent transferring data and waiting instead of computing; it is cut by sending less, less often, and hiding it.</mark>
Techniques.
- Aggregation: combine many small messages into one to pay latency once.
- Overlap: use non-blocking calls to compute while data moves.
- Locality: partition data so that most accesses are local and neighbours are close.
- Prefer collectives, and avoid needless barriers.
| Point | Non-blocking | Asynchronous |
|---|---|---|
| Meaning | Call returns at once, buffer not yet reusable | Transfer proceeds independently of the program |
| Completion | Checked with Wait/Test |
May need no check by the sender |
| Level | Interface semantics | Progress mechanism |
Asked: [7 marks] (May 2022) How to reduce the communication overhead? Differentiate non-blocking and asynchronous communication.
Last-minute revision
- MPI: processes with private memory that communicate by explicit messages, all running one program (SPMD).
- Rank identifies a process in a communicator;
MPI_COMM_WORLDholds all. MPI_SendandMPI_Recvare blocking; both sending first can deadlock.MPI_IsendandMPI_Irecvreturn a request; finish withMPI_WaitorMPI_Test.- Non-blocking allows overlap of computation and communication; do not touch the buffer before completion.
- Collectives: Bcast, Scatter, Gather, Reduce, Allreduce, Barrier; all processes must call them.
- Cartesian topology:
MPI_Cart_create,MPI_Cart_shift. - Communication time $T = T_{lat} + n/B$.
- Reduce overhead by aggregation, overlap and locality.
MPI_Wtimemeasures time; Vampir, TAU and Scalasca trace.
Memory hooks
- Isend and Irecv: the I is for Immediate return.
- Wait or Test before reuse of the buffer.
- Overhead cure ALO: Aggregate, Locality, Overlap.
- Latency is the fixed cost per message, bandwidth the cost per byte.
- Collective means everyone calls it.
Coverage checklist
- Message passing: private memory, send/receive, ranks.
- Message and point to point communication: Send/Recv, tags, deadlock.
- collective communication: Bcast, Scatter, Gather, Reduce.
- non blocking point-to-point communication: Isend/Irecv, Wait/Test, comparison table (May 2022 7 marks).
- virtual topologies: Cartesian grid functions.
- MPI performance tools: Wtime, tracing and profiling tools.
- communication parameters: latency, bandwidth, $T = T_{lat} + n/B$.
- impact of synchronizations sterilizations and contentions: waiting, serialization, link contention.
- reductions in communication overhead: aggregation, overlap, locality, non-blocking vs asynchronous (May 2022 7 marks).