Fuzzball Documentation
Toggle Dark/Light/Auto mode Toggle Dark/Light/Auto mode Toggle Dark/Light/Auto mode Back to homepage

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.

Rank Failure

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.