Computational
Mechanics
wccm-eccm-ecfd2014.org
All fifteen pages →About this site

An independent publication about the finite element method and the computation of continuum mechanics.

Direct against iterative

Parallel, and the partition

Splitting a mesh across processes is its own graph problem, and load balance decides the wall clock.

Section 03Medium readAll pages

Rows of black server racks with blinking indicator lights in a data center
Lead figure

Splitting a mesh across processes is its own graph problem, and load balance decides the wall clock.

Photo: Dutch national supercomputer "Huygens" - 8183833489 · Wikimedia Commons

The partition problem

Every large finite element or finite volume solve eventually hits the same wall: no single machine has enough memory or enough time. The answer is domain decomposition — cut the mesh into subdomains, assign one or more to each MPI process, solve local problems, and communicate boundary data with neighbours. The arithmetic follows almost automatically once the partition is decided. The partition itself is far from automatic.

The fundamental tension is that you want each subdomain to carry the same number of degrees of freedom — so every process does equal work — while minimising the number of element faces that sit on inter-process boundaries, because those faces generate communication. These two goals pull against each other. A partition that equalises load tends to have jagged, irregular interfaces with long perimeters; a partition with clean minimal cuts tends to leave some processes with much less work than others. Graph partitioning ↗ is the formal name for this trade-off: represent the mesh as a graph where nodes are elements and edges connect adjacent elements, then partition the graph to balance node weights while minimising cut edges.

Practical codes rely on multilevel k-way partitioning algorithms, which coarsen the graph through a sequence of aggregations, partition at the coarsest level, then project the partition back up through refinements. The packages METIS and its parallel counterpart ParMETIS, developed by George Karypis and Vipin Kumar at the University of Minnesota, have been the standard implementation since the mid-1990s and are used internally by most industrial and research CFD and FEM codes. The quality of the partition they produce has a direct, measurable effect on solver time: a five-percent imbalance in element count can cost far more than five percent in wall time if one slow process holds up a global synchronisation at every iteration.

An aisle of compute racks in a machine room
The machine the partition is for

Load balance decides the wall clock: the slowest subdomain sets the pace for every other process at each synchronisation point.

Photo: Virginia Tech - data center · Wikimedia Commons

Communication, ghosts, and the iterative loop

Once partitioned, each process stores not only its own elements but also a layer of ghost (or halo) elements from neighbouring partitions. These ghosts carry field values that cross the boundary — the velocity at a face owned by a neighbour, the displacement at a shared node. At every iteration the process updates its interior, then exchanges ghost data with neighbours via non-blocking MPI sends and receives. How much data moves, and how often, is set largely by the length of the partition interface.

For iterative solvers this exchange is cheap if the interface is small relative to the interior. For direct solvers the situation is less benign: methods such as MUMPS or SuperLU_DIST must factor a globally coupled system, and the fill-in generated during factorisation is not confined to subdomain boundaries. Scalability suffers once the factored matrix no longer fits in distributed memory in a useful layout, which is one of the reasons iterative solvers dominate large parallel runs.

The arithmetic follows almost automatically once the partition is decided.

Preconditioner choice compounds this. A block Jacobi preconditioner applies one independent preconditioner per subdomain with no inter-process communication during the apply step — cheap, but its convergence degrades as process count rises. Additive Schwarz methods extend each subdomain by one or more layers of overlap with neighbours, improving convergence at the cost of more communication. Algebraic multigrid ↗ can be applied in parallel and often provides iteration counts that are almost independent of process count, but constructing the coarse levels requires its own distributed graph operations.

Dynamic load imbalance — arising from adaptive mesh refinement, where some partitions accumulate new elements faster than others — demands periodic repartitioning during the run. Repartitioning mid-solve moves data between processes and interrupts the solver, so codes typically tolerate a growing imbalance up to a threshold before triggering it. Choosing that threshold is empirical; there is no universal answer, and the cost of repartitioning versus the cost of running imbalanced is problem- and machine-specific.

A terminal showing a solver residual trace on a dark screen
Cut where the coupling is weakest

Partitioning is a graph problem — minimise the interface, keep the pieces even, and do it faster than the solve it is preparing for.

What the partition decides, ultimately, is whether a solver that converges in theory can be made to run in practice. The mathematics lives in the linear algebra; the time lives in the graph.

Related