Unit content
Distributed-memory parallelism and message passing
In a distributed-memory parallel program, each process owns a separate address space. One process cannot read another process's ordinary variables directly; data crosses process boundaries through explicit communication.
A common design partitions a dataset so each process performs most computation on local data. When one partition needs values owned elsewhere, processes exchange messages containing the required data.
Communication cost has at least two components: a startup latency and a size-dependent transfer cost. A simple model is
$$T_{msg}\approx \alpha+\beta n,$$
where $\alpha$ is per-message latency, $n$ is the amount of data and $\beta$ is time per unit of data.
This explains why sending one large block can be much cheaper than sending the same bytes as thousands of tiny messages.
Distributed-memory parallelism can scale beyond one shared-memory machine and avoids cache-coherence traffic across nodes, but the algorithm must explicitly manage partitioning, communication and synchronization.
Good decompositions maximize useful local work, minimize communicated data and, when dependencies allow it, overlap communication with computation.