Skip to content

Multi-node training (Beta)

View Markdown

Beta: this feature may change before general availability

A distributed Job runs your command on several machines at once, called a gang. Nodus acquires every member, connects them, checks that they can reach each other, and only then starts your command on all of them together. You write an ordinary training script; the launcher finds its peers from environment variables Nodus sets.

train.py
"""A two-node DDP smoke run. torchrun reads its rendezvous from the PET_* variables Nodus sets on every rank."""
import os
import torch
import torch.distributed as dist
from torch.nn.parallel import DistributedDataParallel as DDP
backend = "nccl" if torch.cuda.is_available() else "gloo"
dist.init_process_group(backend)
local_rank = int(os.environ["LOCAL_RANK"])
device = torch.device("cuda", local_rank) if backend == "nccl" else torch.device("cpu")
if backend == "nccl":
torch.cuda.set_device(device)
# Every process contributes 1, so the sum is the world size.
one = torch.ones(1, device=device)
dist.all_reduce(one)
model = DDP(torch.nn.Linear(16, 1).to(device), device_ids=[local_rank] if backend == "nccl" else None)
opt = torch.optim.SGD(model.parameters(), lr=0.1)
for step in range(20):
x = torch.randn(32, 16, device=device)
loss = (model(x) - x.sum(dim=1, keepdim=True)).pow(2).mean()
opt.zero_grad()
loss.backward()
opt.step()
if dist.get_rank() == 0:
print(f"world_size={int(one.item())} node={os.environ['NODUS_NODE_RANK']} "
f"transport={os.environ['NODUS_GANG_TRANSPORT']} loss={loss.item():.4f}")
dist.destroy_process_group()
run.sh
nodus run --name torchrun-2-node --gpu H100 --nodes 2 --launcher torchrun --max-cost 2.00 -- torchrun train.py

torchrun needs no flags: its node count, rendezvous address and local address come from the PET_* variables on every rank. The run prints world_size=2 from rank 0 when both nodes joined.

Say how big the gang is with exactly one of:

Field (nodus run flag) Meaning
distributed.nodes (--nodes) The number of machines, 2 to 8
distributed.totalGPUs (--total-gpus) The total GPU count; Nodus picks the machine shape. A total one machine can hold becomes a single-node Job
distributed.gpusPerNode (--gpus-per-node) GPUs per machine: 1, 2, 4 or 8. With nodes it defaults to the GPU count of resources.gpu; with totalGPUs Nodus picks it unless you set it

Every member gets the same GPU count. When resources.gpu.type lists several accelerator types, members may get different ones, each billed at its own rate. The resolved shape is written once to status.topology and never changes on a restart, so the world size stays the same for the whole run. A dry run shows it, with the hold for every member (gangHoldUSD) and the assembly bound:

Terminal window
$ nodus run --dry-run --gpu H100 --total-gpus 16 --launcher torchrun -- torchrun train.py

distributed.network bounds where members may be:

network Members are placed
Colocated (default) With one provider in one region, on its private network
Regional With any providers inside one region class
Global Anywhere. The run shows the WANBound condition, because training between distant machines is limited by the network

distributed.transport says which network paths you accept. Direct (default) admits every path that is not relayed: a provider’s private network, the members’ public addresses between providers (unencrypted in the Beta), and direct encrypted paths. Auto also admits relayed paths, which carry traffic through Nodus relays: they work between providers that cannot reach each other directly, but they are slow and fit small models only. The estimate warns RelayedLowBandwidth when a relayed path is possible.

run.sh (across providers)
nodus run --name torchrun-2-node-xp --gpu H100 --nodes 2 --launcher torchrun \
--network global --transport auto --startup-timeout 15m --max-cost 2.00 -- torchrun train.py
launcher What each rank runs
Plain (default) Your command once per node. RANK, LOCAL_RANK=0, LOCAL_WORLD_SIZE=1 and WORLD_SIZE (the node count) make env:// initialization work for one process per node
Torchrun Your torchrun command on every node; rank 0 hosts the rendezvous. Nodus owns restarts, so PET_MAX_RESTARTS=0
Ray A Ray cluster: rank 0 is the head, the others join it, and your command runs on rank 0 once all nodes are up, with RAY_ADDRESS set
Verl The Ray launcher plus NODUS_VERL_OVERRIDES (trainer.nnodes, trainer.n_gpus_per_node) to append to your verl command
Accelerate, Deepspeed The environment contract plus ACCELERATE_* or /etc/nodus/hostfile. Accepted in the Beta, qualified later

Every rank gets the same variables except the rank-specific ones:

Variable Value
NODUS_NODE_RANK, NODE_RANK The Nodus rank of this node; rank 0 hosts the rendezvous
NODUS_NUM_NODES, NNODES, NODUS_GPUS_PER_NODE The node count and GPUs per node
NODUS_NODE_IPS Every member’s address in rank order
MASTER_ADDR, MASTER_PORT Rank 0’s address and 29500
NODUS_GANG_EPOCH, NODUS_GANG_TRANSPORT The current epoch and its path: private, public, direct or relayed
NODUS_RESTORE_URI The checkpoint to resume from after a restart, when one exists

NODUS_NODE_RANK is Nodus’s rank. torch assigns its own global RANK inside torchrun and may order nodes differently. Setting any NODUS_*, PET_* or MASTER_* variable, NODE_RANK or NNODES in your spec is rejected. Your NCCL_* settings are kept, except the few the network path decides.

If a member’s machine is lost or reclaimed, Nodus stops the whole gang’s current epoch at once and restarts it: surviving machines are kept, only the lost ranks get new machines, and every rank starts again at the next epoch from the latest gang checkpoint (NODUS_RESTORE_URI). A process left over from the old epoch cannot join the new rendezvous. recovery.maxAttempts counts these restarts (default 3). The first failure’s cause is the gang’s reason, so a crash on one rank that brings down the others reports that crash.

If a machine is refused or never appears while the gang is still assembling, only that rank gets another machine (up to two per assembly) and the registered ranks keep waiting at the start barrier. A suspend or cancel in progress is never turned into a restart: a machine lost during it simply completes the stop.

If your command exits non-zero on any rank, the run fails; a zero exit on every rank succeeds it.

Each rank’s attempts carry the nodus.dev/rank label, so nodus get attempts -l nodus.dev/job=<name>,nodus.dev/rank=1 lists rank 1’s attempts across restarts, with each attempt’s placement, boot or restore time and cost.

Assembly waits at most distributed.startupTimeout (default 15 minutes, 5 to 60) for every member to become ready and pass the network check. If it does not, every acquired member is released, and Nodus tries again up to maxAssemblyRetries times (default 2) before the run fails with GangAssemblyTimeout.

You pay for each member from the moment its provider starts billing until it is confirmed deleted, including time spent waiting for the other members. The estimate shows the most a failed assembly can cost, the assembly bound: the sum of the members’ rates plus up to two replacement machines, times the startup timeout plus teardown, times the number of tries. Nodus pays, not you, when assembly fails for a Nodus reason: a network check that fails on a path Nodus lists as qualified, an address collision, or an outage of the Nodus mesh.

  • Public paths carry traffic between members’ public addresses in clear text. Only the gang’s own members may connect to a member’s rendezvous and collective ports, but the bytes are not encrypted on the way.
  • Relayed paths carry every byte through a Nodus relay, are billed per relayed GiB, and suit small models only.
  • Relayed gangs have at most 2 members until larger relayed gangs are qualified.
  • On relayed paths your image must use dynamically linked glibc programs; otherwise the run fails with ShimNotLoaded.
  • Provider pairs are added as each is qualified; a pair that is not qualified is never offered. The measured TCP throughput and round-trip time per path class are published here as pairs are qualified.