Wiring is not coherence
A companion piece in this series toured the physical plant beneath an AI datacenter — the substation, the coolant loops, the network fabric — and argued that the building, not the accelerators, is the engineering object. That argument stops at the rack door. Once power and cooling are delivered and the fabric is lit, a separate and much less photographed engineering problem begins: turning several thousand independently failing machines into something that behaves, for the seconds or months a job runs, like one machine.
That problem is solved by software, not steel, and it has four load-bearing pieces. A scheduler decides which jobs get to run, and for distributed training it must decide this as an indivisible, all-or-nothing admission rather than machine by machine. A checkpoint pipeline decides how much computed work a failure is allowed to erase, and its write frequency is a real optimization with a real answer, not a rule of thumb. A failure-domain design decides how far a single fault is allowed to spread before it is contained — the blast radius. And a telemetry pipeline decides how quickly any of this is even noticed. Underneath all four sits one more piece that has to serve two incompatible workloads at once: the storage architecture that feeds training data in continuously and catches synchronized checkpoint writes in bursts.
This article works through each in order, using production systems papers from the operators who actually run fleets at this scale, and closes on the number — fleet utilization — in which all four either show up or hide.
Gang scheduling: the all-or-nothing admission problem
A distributed training job is not a set of independent tasks that happen to run at the same time. Under data, tensor, or pipeline parallelism, every rank in the job depends on every other rank producing its shard of the forward and backward pass on roughly the same schedule; a job with half its ranks running and half queued does not run at half speed, it does not run at all. The scheduling literature calls admitting such a job as one unit gang scheduling, and it is the first structural fact that separates an AI cluster’s scheduler from a general-purpose batch scheduler.
Gang scheduling interacts badly with two things clusters otherwise want: high utilization and low queuing delay. Microsoft’s study of two months of production traffic on its Philly clusters found that gang-scheduling and data-locality constraints directly shape both, reporting on “the effect of gang scheduling and locality constraints on queuing” and “the effect of locality on GPU utilization” as two of the three central issues its trace analysis was built to explain, alongside failures during training [12]. The mechanism is simple to state and expensive to fix: a job needs a specific count of simultaneously free, often topologically close, accelerators, and until that exact set becomes free, the whole job sits in queue no matter how much aggregate idle capacity the cluster is carrying elsewhere.
A minimal model makes the difficulty explicit. If each of a fleet’s ranks is independently healthy and idle with probability
This is a simplification — real fleets are not independent, and real schedulers do not wait for a spontaneous coincidence — but it explains why gang admission gets combinatorially harder, not linearly harder, as jobs grow, and why every production scheduler in this space is built around avoiding literal simultaneous-availability matching rather than performing it. Google’s Borg, one of the earliest cluster managers to operate at this scale, achieves its utilization through “admission control, efficient task-packing, over-commitment, and machine sharing,” and explicitly uses “scheduling policies that reduce the probability of correlated failures” [10] — packing and correlation-aware placement are both ways of making the effective
At the largest scale the same problem reappears one level up, across datacenters rather than across racks. Meta’s MAST schedules ML training globally and reports that before its deployment, “the most overloaded region had a GPU demand-to-supply ratio of 2.63 for high-priority workloads,” a ratio MAST’s placement reduced to 0.98, “effectively eliminating the overload” [2]. MAST’s three design principles — temporal decoupling, scope decoupling, and exhaustive search — are, in effect, three different ways of narrowing an intractable global matching problem down to one small enough to actually solve before the answer is stale [2].
Checkpointing: what actually gets written, and how often
Checkpointing exists because a training job that runs for weeks will, with near certainty, be interrupted before it finishes. A checkpoint captures enough state — weights, optimizer moments, and often RNG and dataloader position — to resume from that point rather than from zero. The interesting engineering question is not whether to checkpoint but how often, and that question has a real answer rather than a folk one.
Treat checkpointing as a classic renewal-reward tradeoff — this is my own worked derivation, using standard reasoning from fault-tolerant computing rather than any claim from the sources below. Let
which is minimized at
Two consequences follow directly, and both are visible in how production systems have actually evolved. First, because
That is exactly the direction three separate production systems moved. Meta’s Check-N-Run, built for large recommendation models where checkpoints “take a snapshot of an ML model and store it in a non-volatile memory,” attacks
Meta’s account of training the 405-billion-parameter Llama 3 model gives the concrete numbers behind all of this. Each GPU’s checkpointed state runs one to four gigabytes, and the stated design goals are exactly
Failure-domain design: blast radius as an architectural choice
Not every fault should be allowed to cost the whole fleet. Failure-domain design is the discipline of deciding, in advance, how large a group of machines is allowed to go down together — the blast radius — and it trades two costs against each other: smaller domains contain damage more tightly but add placement and coordination overhead, while larger domains are simpler to schedule against but let one bad power shelf, one bad top-of-rack switch, or one corrupted driver update take out proportionally more in-flight work.
Meta’s large-scale reliability study, drawing on eleven months of data across two multi-tenant clusters totaling more than 150 million A100 GPU-hours and four million jobs, makes the stakes of domain sizing explicit: it finds that “while large jobs are most vulnerable to failures, smaller jobs make up the majority of jobs” in the fleet, and builds a taxonomy of failure types together with a fitted model projecting mean time to failure across GPU scales specifically to reason about this asymmetry [8]. A domain boundary drawn without that asymmetry in mind protects the common case — small jobs — at exactly the moment the rare, expensive case — a large job spanning many domains — is most exposed.
Borg’s contribution to this problem is placement policy rather than a fixed partition: it reports using “scheduling policies that reduce the probability of correlated failures,” which is domain design applied dynamically at scheduling time rather than fixed once in the datacenter’s floor plan [10]. MAST’s geo-distributed placement operates at the largest domain granularity available — the datacenter or region itself — and its demand-to-supply rebalancing is, read one way, an automated response to the fact that concentrating too much high-priority load in one region turns that region into an oversized failure and capacity domain simultaneously [2].
The response to a fault inside a domain is where the domain design pays off or fails to. The Shanghai AI Laboratory’s characterization of six months of large-model development traffic on its Acme datacenter introduces “fault-tolerant pretraining, which enhances fault tolerance through LLM-involved failure diagnosis and automatic recovery” as a direct response to the “frequent hardware failures” and “imbalanced resource utilization” the trace exposed [7]. Automatic recovery is only cheap, in wall-clock terms, if the domain that failed is small enough that the rest of the job’s ranks were never touched — which ties failure-domain sizing directly back to the checkpoint-interval mathematics above: a smaller blast radius does not just limit hardware damage, it limits how many ranks have to roll back to the last checkpoint at all.
Telemetry: turning silicon behavior into a decidable signal
None of the above works without knowing, quickly and specifically, what just broke. Fleet telemetry is the pipeline that continuously pulls low-level signals off every accelerator and every chassis — temperatures, error-correcting-code counts, link and retrain errors, power draw, clock throttling — and turns them into something a scheduler or an operator can act on before a slow degradation becomes a hard failure that a checkpoint has to absorb.
NVIDIA’s Data Center GPU Manager is the accelerator-level instance of this pipeline. Its documentation describes “active health monitoring, comprehensive diagnostics, system alerts and governance policies” delivered as “job-level statistics and continuous GPU telemetry” at low overhead, built around what NVIDIA calls “a group-centric philosophy to node level GPU management” so that operators can reason about many GPUs per host and many hosts per job as one unit rather than as thousands of individual devices [13]. That group-centric framing matters for the same reason gang scheduling does: a fleet’s operational unit is the job’s whole footprint, not any single card inside it.
Above the accelerator, the chassis and out-of-band management layer has its own telemetry standard. DMTF’s Redfish specification defines a RESTful management interface built to be “suitable for a wide range of devices, from stand-alone servers, to composable infrastructures, and to large-scale cloud environments,” and its schema includes a dedicated metric-report stream delivered over server-sent events alongside an event-subscription service for asynchronous state-change notification [14]. Where DCGM answers “what is this GPU doing,” Redfish answers “what is this chassis doing” — power state, fan speed, component presence — in a vendor-neutral form that lets one fleet-management system talk to hardware from many manufacturers without a bespoke integration for each.
Telemetry only earns its keep when it is wired to automated response rather than a screen a human has to watch. MegaScale’s account of training a 175-billion-parameter model across more than ten thousand GPUs describes building “a set of diagnosis tools to monitor system components and events deep in the stack, identify root causes, and derive effective techniques to achieve fault tolerance and mitigate stragglers” [1], and the Acme trace paper’s fault-tolerant pretraining system is explicitly described as using diagnosis output to drive “automatic recovery” rather than paging an engineer [7]. In both cases the telemetry pipeline is the sensing half of a closed loop whose acting half is the scheduler and the checkpoint system described above — three of the four disciplines in this article, wired together.
Storage: feeding data in and catching checkpoints on the way out
One storage layer has to serve two workloads that want opposite things from it. Training-data ingestion is read-heavy, mostly steady, and tolerant of moderate latency because it is pipelined ahead of compute. Checkpoint writing is the opposite: it is a write-heavy burst in which every rank in a job — potentially tens of thousands of them — tries to push state out at nearly the same instant, and the job is stalled, or degraded, until enough of that write has landed somewhere durable.
Meta’s Tectonic filesystem is the general answer to running mixed, exabyte-scale workloads on shared infrastructure rather than maintaining a separate specialized system per tenant. The FAST 2021 paper describes Tectonic as Facebook’s “exabyte-scale distributed filesystem,” designed to consolidate large tenants that previously ran service-specific storage systems onto shared multitenant instances at performance comparable to those specialized systems, while allowing individual tenants to specialize their own operations for their own workload shape [6]. A checkpoint-writing tenant and a data-ingesting tenant sharing one Tectonic instance is exactly the kind of workload heterogeneity the design exists to absorb without either tenant needing its own dedicated hardware.
The two checkpoint-specific systems already discussed attack the same contention problem from the other direction: by shrinking what has to reach durable storage at all, or by not sending it to durable storage in the first place. Check-N-Run’s incremental, quantized checkpoints reduce the volume a shared storage tier has to absorb per write by up to an order of magnitude [3], and GEMINI removes the synchronized burst from the durable-storage path entirely for the common case, landing checkpoints in host CPU memory and only falling back to remote storage when memory-resident copies cannot cover a failure [4]. Read together, Tectonic is the substrate built to tolerate two contending workloads sharing capacity, and the checkpoint systems are the mechanism that keeps one of those two workloads from growing large enough to make that tolerance necessary in the first place. Neither approach alone would be sufficient at fleet scale; a storage architecture built only to be generously multitenant would still choke on an unshrunk checkpoint burst, and a shrunk checkpoint with nowhere resilient to land would still be a single point of failure.
Utilization: the number that summarizes everything above
Fleet-level utilization is not a fifth mechanism alongside scheduling, checkpointing, failure domains, and telemetry. It is the scoreboard on which all four either show a gain or hide a loss, and the production numbers make that legible.
MegaScale reports 55.2% model FLOPs utilization training a 175-billion-parameter model across 12,288 GPUs, a 1.34× improvement over Megatron-LM attributed to its full-stack co-design, including the diagnosis and fault-tolerance tooling described above [1]. That is an efficiency number for the computation itself, conditional on the job staying up. Meta’s reliability study proposes a companion metric it calls Effective Training Time Ratio specifically because MFU alone cannot distinguish a job that ran efficiently the whole time from one that ran efficiently but spent a third of its wall-clock time recovering from failures [8] — and Llama 3’s reported greater-than-90% effective training time, held despite 419 unexpected interruptions in 54 days, is that second number made concrete for one specific, disclosed production run [9].
Utilization loss also shows up before a job ever starts computing, in the scheduler’s queue. The Philly trace analysis ties gang-scheduling and locality constraints directly to queuing delay and to GPU utilization loss from placement mismatches [12], and MAST’s demand-to-supply rebalancing from 2.63 down to 0.98 in the most overloaded region is a queuing-side utilization gain achieved entirely through better global placement, with no change to any individual job’s efficiency [2]. None of these four numbers — MFU, effective training time ratio, queuing delay, regional demand-supply balance — is optional context for the others. A fleet can score well on any one of them while losing badly on another, which is exactly why operators who run at this scale report several of these figures rather than treating any single one as sufficient.
Predictions, with the observations that would falsify them
These are forecasts, separated from the sourced analysis above. Horizon: 15 August 2029.
One. Near-per-iteration checkpointing, or an equivalent failure-triggered scheme that removes the fixed-interval decision entirely, will become the default at frontier training scale rather than a specialized optimization, because the renewal-reward mathematics above makes coarser intervals increasingly wasteful as job-level failure rates rise with scale. Disconfirmed if published frontier-scale training systems in 2029 still checkpoint on intervals of an hour or longer as standard practice.
Two. Vendor-neutral telemetry interfaces in the pattern of Redfish will extend further into accelerator-level monitoring, driven by fleets that mix accelerators from more than one manufacturer. Disconfirmed if 2029 fleet-management stacks at multi-vendor sites remain built on separate, non-interoperable per-vendor telemetry systems with no common schema.
Three. A reliability-adjusted companion to model FLOPs utilization — something in the shape of effective training time ratio — will appear as a standard disclosure alongside MFU in system reports for large training runs, because MFU alone is demonstrably unable to separate computational efficiency from fleet reliability. Disconfirmed if leading system reports in 2029 still publish MFU or equivalent throughput figures with no accompanying uptime or effective-time metric.
Four. Failure-domain sizing will increasingly be published as an explicit design parameter — a stated blast-radius policy — rather than held as undisclosed internal scheduler configuration, following the pattern already visible in the geo-distributed placement work cited above. Disconfirmed if operators publishing detailed reliability studies in 2029 continue to omit any stated policy for how large a domain is allowed to fail together.
None of these requires a hardware discontinuity. Each follows from constraints already visible in the sourced material above: failure rates that scale with node count, a checkpoint-interval optimum that responds to those rates, and a growing gap between what MFU measures and what an operator actually needs to know.
What to take away
A cluster is coherent because four pieces of software make it so, not because its wiring does. The scheduler decides admission as one indivisible unit, because a distributed training job is not divisible. The checkpoint pipeline decides how much work a failure erases, and that decision has an actual optimum set by a write cost and a failure rate, not a habit. The failure-domain design decides how far damage spreads before it is stopped, trading coordination overhead against blast radius. The telemetry pipeline decides how fast any of the other three even find out something is wrong. And underneath all four, one storage architecture has to serve a steady read-heavy pipeline and a bursty, synchronized write pipeline without either starving the other.
Utilization is not a fifth thing to manage. It is the number that remains after the other four have already had their effect, which is why an operator who reports only aggregate GPU-hours sold has told you almost nothing, and an operator who reports MFU, effective training time, queuing delay, and domain-failure rate together has told you how the whole system actually behaves. The physical plant sets the ceiling on what a site can deliver. This layer decides how much of that ceiling a training run actually gets to keep.