A training job that needs 64 GPUs is not 64 independent requests. It needs all 64 at once, on nodes that can talk to each other quickly, and it needs a plan for the hour when one of them fails. GPU orchestration is the set of systems that make those decisions: which jobs run, in which order, on which devices, and what happens when hardware breaks. Getting it wrong costs real money, because an idle accelerator costs as much per hour as a busy one.
This article is an overview of the cluster level. It explains the layers from job submission down to the node, the two dominant stacks (Kubernetes with a batch queue, and Slurm), gang and topology-aware placement with tested code, failure handling, and a worked example on a 64-GPU cluster. Scheduling inside one GPU (streams, MPS and MIG) is covered in GPU scheduling architecture; this article treats a GPU, or a slice of one, as the unit being handed out.
The layers of orchestration
Orchestration decomposes into questions that different components answer. Keeping them separate in your head is the fastest way to debug a stuck job.
- What exists? Inventory: which nodes have which GPUs, their model, memory, health and interconnect. On Kubernetes this comes from a device plugin or a Dynamic Resource Allocation driver; on Slurm from
gres.confand the node definitions. - Who may use it? Quota and fairness: team limits, priorities, borrowing idle capacity and preemption.
- Can the whole job start now? Admission and gang scheduling: a distributed job starts only when every worker can be placed.
- Where exactly? Placement: pick nodes that minimise communication cost and fragmentation.
- What if it breaks? Health, remediation and restart from checkpoint.
Inventory: how a GPU becomes schedulable
On Kubernetes the classic path is the device plugin API. A plugin running on each node, such as NVIDIA's k8s-device-plugin, registers with the kubelet and advertises GPUs as an extended resource, nvidia.com/gpu. Extended resources are integers and cannot be overcommitted, so a pod asks for whole devices and the scheduler counts them like CPU cores. The NVIDIA GPU Operator packages the pieces a node needs: the driver, the container toolkit, the device plugin, GPU Feature Discovery (which adds node labels such as the GPU product) and the DCGM exporter for metrics.
apiVersion: v1
kind: Pod
metadata:
name: trainer-0
spec:
nodeSelector:
nvidia.com/gpu.product: NVIDIA-H100-80GB-HBM3 # label value comes from GPU Feature Discovery
containers:
- name: trainer
image: registry.example.com/train:2026.10
resources:
limits:
nvidia.com/gpu: 8 # whole devices; requests default to limitsCounting devices is too coarse for modern hardware, which is why Kubernetes added Dynamic Resource Allocation. DRA graduated to stable in Kubernetes v1.34 with the resource.k8s.io/v1 API. Drivers publish devices and their attributes in ResourceSlice objects, administrators define DeviceClasses, and workloads ask for devices through ResourceClaims (or ResourceClaimTemplates) that can select on attributes rather than just a count. The scheduler then allocates specific devices. If your cluster and driver support it, DRA is the better foundation for sharing, partitioned devices and attribute-based selection; check your vendor driver's maturity before depending on it. Always verify the label value in the example above on your own nodes rather than copying it.
Gang scheduling: all or nothing
The default Kubernetes scheduler places pods one at a time. That is fine for web services and dangerous for distributed training. Take two free 8-GPU nodes and two 16-GPU jobs, each with two 8-GPU workers. If the scheduler places one worker of each, all 16 GPUs are held by jobs that cannot start, nothing is left for either second worker, and nothing progresses until a timeout fires. This is partial-gang deadlock, and at scale it quietly burns a large share of the fleet.
Gang scheduling makes allocation all-or-nothing. Three common ways to get it:
- Kueue. Jobs are created suspended and queued in a LocalQueue that points at a ClusterQueue holding quota per ResourceFlavor (for example, one flavor per GPU type). Kueue admits a workload only when its entire request fits the quota, then unsuspends it. ClusterQueues in a cohort can borrow each other's idle quota. Kueue also offers topology-aware scheduling.
- Volcano and coscheduling plugins. A PodGroup declares
minMember; the scheduler binds none of the pods until that many can be bound. - Slurm. A job's allocation is atomic by design:
sbatch --nodes=4 --gpus-per-node=8either gets all four nodes or waits. The backfill scheduler lets small jobs run in holes that will not delay the reserved start of larger ones.
Quota admission and placement are different checks. A job can be within quota and still not fit because free GPUs are scattered across nodes. That is fragmentation, and placement policy controls it.
Topology-aware placement
Placement has two goals that pull against each other: put a job's workers close together, because collective communication across racks is slower than inside a node or an NVLink domain, and pack small jobs tightly, so whole nodes stay free for the next big job. The sketch below captures both: multi-node jobs take whole nodes, preferring the tightest single rack that fits and spilling across as few racks as possible; single-node jobs use best fit, choosing the node they leave least free.
from dataclasses import dataclass
@dataclass
class Node:
name: str
rack: str
free: int = 8
healthy: bool = True
def place_gang(gpus, nodes, per_node=8):
"""All-or-nothing: return {node: gpus} for the whole job, or None and allocate nothing."""
ok = [n for n in nodes if n.healthy]
if gpus <= per_node: # single node: best fit
fits = [n for n in ok if n.free >= gpus]
if not fits:
return None
best = min(fits, key=lambda n: (n.free - gpus, n.name))
return {best.name: gpus}
if gpus % per_node:
raise ValueError("multi-node jobs take whole nodes")
need = gpus // per_node
racks = {}
for n in sorted(ok, key=lambda n: n.name):
if n.free == per_node:
racks.setdefault(n.rack, []).append(n)
one = [r for r in racks.values() if len(r) >= need]
if one: # tightest single rack that fits
return {n.name: per_node for n in min(one, key=len)[:need]}
spill = [n for r in sorted(racks.values(), key=len, reverse=True) for n in r]
if len(spill) < need:
return None
return {n.name: per_node for n in spill[:need]} # fewest racks, fullest firstReal schedulers express the same idea with richer topology (NVLink domain, leaf switch, spine), network rails and scores rather than hard rules, but the shape is the same: filter by feasibility, then rank by locality and by how much contiguous capacity the choice leaves behind.
When hardware fails
At scale, hardware failure is routine. A large synchronous job fails when any one of its GPUs, links or hosts fails, so its failure rate grows roughly linearly with its size. Orchestration has to turn a failure into a short pause instead of a lost day. The loop has four steps.
- Detect. DCGM health watches and diagnostics, driver Xid events (for example Xid 79, GPU fallen off the bus), ECC error counters, and application signals such as NCCL timeouts. See GPU Hardware Faults for what each signal means.
- Isolate. Cordon the node (Kubernetes) or drain it (Slurm) so nothing new lands on it, and remove the bad device from inventory.
- Restart. Tear down the whole gang, requeue it, and resume from the latest checkpoint. Partial restarts of a synchronous job rarely work.
- Repair and return. Run diagnostics, reset or replace, then uncordon only after the node passes a burn-in test.
def on_health_event(node, event, cluster, queue):
if event.kind in {"xid_fatal", "ecc_uncorrectable", "dcgm_health_fail"}:
cluster.cordon(node) # stop new placements first
for job in cluster.jobs_on(node):
cluster.stop_gang(job) # every worker, not only the local ones
queue.requeue(job, resume_from=job.latest_checkpoint())
cluster.open_repair_ticket(node, evidence=event)
elif event.kind == "nccl_timeout":
queue.requeue(event.job, resume_from=event.job.latest_checkpoint())
cluster.suspect(event.job.nodes) # quarantine only after repeat offencesThe checkpoint interval sets how much work a failure costs. Young's approximation gives an interval near sqrt(2 * C * M), where C is the time to write a checkpoint and M is the job's mean time between failures. Asynchronous checkpointing shrinks C and allows a shorter interval; see GPU Checkpointing Deep Dive.
Worked example: a 64-GPU cluster
A cluster has 8 nodes of 8 GPUs, n1 to n4 in rack r1 and n5 to n8 in rack r2. Four jobs arrive in priority order: train-A (32 GPUs), train-B (16 GPUs) and four 2-GPU eval jobs. Running place_gang on this state gives the following.
| Job | Placement | Why |
|---|---|---|
| train-A, 32 | n1, n2, n3, n4 | both racks fit 4 whole nodes; the tie goes to r1 |
| train-B, 16 | n5, n6 | r2 is the only rack with whole nodes left |
| eval-1 to eval-4, 2 each | all four on n7 | best fit keeps packing the partly used node |
| (free) | n8 whole | a whole node remains for the next 8-GPU job |
A spreading policy, such as the Kubernetes scheduler's default least-allocated scoring, could put two evals on each of n7 and n8, leaving 8 GPUs free but no whole node, so a waiting 8-GPU job would starve. Now n3 raises Xid 79. The health loop cordons n3, stops all four train-A workers and requeues it. The healthy whole nodes are now n1, n2 and n4 in r1 plus n8 in r2, so the restart lands on n1, n2, n4 and n8, spanning racks. That works but runs slower if inter-rack bandwidth is lower. The operator's choice is to accept the slower placement, wait for n3's repair, or let train-A run at 24 GPUs if the framework supports elastic restarts.
Assume a checkpoint takes 2 minutes to write and train-A's mean time between failures is 10 hours. Young's approximation gives sqrt(2 x 2 x 600), about 49 minutes, so checkpoint every 45 to 50 minutes and expect to lose roughly half an interval of work, plus restart time, per failure.
Kubernetes versus Slurm
| Concern | Kubernetes + Kueue | Slurm |
|---|---|---|
| Unit of work | pods grouped into a Workload | jobs and job steps |
| Gang semantics | admission of the whole workload | native: allocation is atomic |
| Quota and fairness | ClusterQueue quota, cohorts, borrowing | partitions, QOS, fair-share |
| Device inventory | device plugin or DRA | GRES |
| Strength | shared platform with serving and pipelines | mature batch scheduling, backfill |
| Weak spot | more moving parts to integrate | less natural for long-running services |
Many organisations run both: Slurm for large pretraining reservations, Kubernetes for serving, evaluation and pipelines. The concepts in this article map across, and GPU Scheduling With Slurm covers the Slurm side in detail.
Failure modes
- Partial gangs. Pods bound one at a time hold GPUs while waiting for siblings. Use gang admission everywhere distributed jobs run.
- Fragmentation. Plenty of free GPUs, no whole nodes. Track the count of fully free nodes as a first-class metric, not just total free GPUs.
- Stale inventory. A GPU that failed after allocation still counts as healthy and keeps attracting jobs that crash at start. Feed health events back into inventory automatically.
- Zombie allocations. A job whose rank 0 died while the others wait in a collective holds the whole allocation. Set collective timeouts and a progress watchdog.
- Preemption storms. Aggressive preemption without checkpoints throws away hours of work. Preempt only jobs that checkpoint, and give them a grace period to save.
Trade-offs
Strict topology rules improve step time but lengthen queue waits; soft scores do the opposite. Quota borrowing raises utilisation but makes preemption, and therefore checkpoint discipline, mandatory. GPU sharing through time-slicing or MIG raises utilisation for small inference and notebook workloads but adds interference or rigid partitions. Measure queue wait, GPU allocation and actual utilisation from DCGM together; optimising one alone moves the cost into another.
What to do next
- Draw your own stack as the five layers above and name the component that answers each question.
- Check that every multi-node job path uses gang admission, and test it by submitting two jobs that together exceed capacity.
- Add a dashboard for fully free nodes, queue wait by job size, and allocated versus utilised GPUs.
- Wire DCGM and Xid events into automatic cordon and requeue, and rehearse one injected node failure.
- Set each large job's checkpoint interval from its measured failure rate and checkpoint time.
- Write down your preemption policy, including the grace period, and tell every team that borrows quota.