# Multi-node training (Beta)

> Run one training job across several machines with torchrun, Ray or your own launcher, and understand what it costs to assemble them.

Source: https://nodus-platform-site.pages.dev/docs/guides/multi-node/
Build revision: 211ad9f836655b1c3a2668c4693e442471f28614

**Beta:** this feature may change.

Beta

Multi-node training is in Beta and on for every organization, with no access request and no purchase needed. Beta gangs have at most 8 nodes. When members on different providers connect over their public addresses (the `public` path), traffic between them, including gradients and the rendezvous, is **not encrypted**; each member accepts it only from the other members’ addresses. Keep `network: Colocated` if your data must not cross the internet unencrypted.

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

train.py

```python
"""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

```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.

## 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:

Terminal window

```console
$ 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)

```sh
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
```

## 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

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

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

* 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.
