Parallel / Distributed Jobs
The terms parallel and distributed can mean different things depending on the person and context. In these guides, parallel jobs mean collections of tasks that can run at the same time independently of one another. Distributed jobs, on the other hand, are those that use applications specially written to run across multiple memory spaces (typically nodes) thereby allowing them to scale up on clusters. Distributed jobs often use MPI or another distributed programming language like Chapel.
Fuzzball supports parallel jobs through “task arrays” and distributed jobs that run on either MPI or PGAS-based networks. For custom distributed workloads that do not fit into these frameworks, Fuzzball also provides a generic multinode implementation that gives you direct control over process coordination across nodes.
A multinode job is scheduled as a group: every rank starts together or none does. It also fails as a group. Each rank reports its own exit, so a rank that is killed, runs out of memory, or whose container crashes fails the whole job within a second or so, and the reported error names the rank, its node, and the exit status. The remaining ranks are then stopped and their resources released, rather than the job continuing a rank short.
This applies to every multinode implementation – openmpi, mpich, gasnet, pmix, and
generic alike – because it is a property of the multinode block rather than of the launcher.
Two cases are slower or invisible. A rank whose node stops answering the control plane is
detected by its lease heartbeat going unanswered, which takes minutes rather than a second. A rank
whose command merely hangs is not detected at all, since its container keeps reporting as
healthy. A rank stopped politely with SIGTERM exits cleanly and is not a failure – that is how
the platform tears a job down.