Multi-node training (Beta)
View MarkdownBeta: 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.
Run a two-node job
Section titled “Run a two-node job”"""A two-node DDP smoke run. torchrun reads its rendezvous from the PET_* variables Nodus sets on every rank."""import os
import torchimport torch.distributed as distfrom 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()nodus run --name torchrun-2-node --gpu H100 --nodes 2 --launcher torchrun --max-cost 2.00 -- torchrun train.pytorchrun 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.
Topology
Section titled “Topology”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:
$ nodus run --dry-run --gpu H100 --total-gpus 16 --launcher torchrun -- torchrun train.pydistributed.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.
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.pyLaunchers
Section titled “Launchers”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.
Failures and restarts
Section titled “Failures and restarts”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.
The assembly bound
Section titled “The assembly bound”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.
Beta caveats
Section titled “Beta caveats”- 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.