A different kind of reliability problem

Mainframe engineering pursued reliability by making the single machine trustworthy: redundant power supplies, error-correcting memory, hot-swappable parts, enough headroom that any one component could be serviced without stopping the job. Cluster systems engineering — the discipline this article follows — made the opposite bet. In the early 2000s, engineers building web-scale search systems on thousands of commodity machines concluded that no amount of single-machine hardening would save them, because at that population size something was statistically always broken. The rational response was not to prevent failure. It was to write software that assumed failure continuously and kept producing correct results anyway.

That decision is dated specifically to a cluster of publications between 2003 and 2006, and it is the founding act of everything that follows: five decades of practice in how a fleet of unreliable parts is scheduled, coordinated, watched, checkpointed and built so the aggregate keeps working. The hardware changed completely between the first chapter of this story and the last. The engineering problem did not.

Failure as the default assumption, not the exception

The clearest statement of the new assumption is architectural rather than rhetorical. Ghemawat, Gobioff and Leung, describing the Google File System at the 2003 Symposium on Operating Systems Principles, built a storage system for “hundreds of terabytes of storage across thousands of disks on over a thousand machines,” accessed concurrently by hundreds of clients — explicitly because ordinary file systems assumed hardware that was mostly reliable, and this hardware was not [1]. GFS answered with automatic chunk replication and continuous integrity checking rather than an attempt to make any individual disk trustworthy: the file system, not the disk, was the unit of reliability.

ADVERTISEMENT

A companion paper the same year made the scale explicit. Barroso, Dean and Hölzle described Google’s search infrastructure as running on “clusters of more than 15,000 commodity-class PCs,” chosen deliberately over smaller numbers of high-end servers because price-performance favored volume over individual robustness, with fault-tolerant software absorbing the resulting failure rate [3]. This is the paper that names the trade explicitly: buy cheap, unreliable units in bulk, and spend the engineering budget on the software layer that survives their failure instead of on parts that resist it.

The programming model that made this survivable at the application layer arrived in 2004. Dean and Ghemawat’s MapReduce let engineers “without any experience with parallel and distributed systems” run computations across thousands of machines, with the runtime handling partitioning, scheduling and, centrally, machine failures — by 2004 Google was running more than a thousand MapReduce jobs a day over many terabytes of input [2]. The mechanism was almost embarrassingly simple relative to its consequence: because a map or reduce task read from and wrote to durable storage rather than holding state only in a live process, a failed task could simply be re-executed elsewhere. Failure was not caught and handled as an exception. It was absorbed by re-running a unit of work that had never assumed it would get exactly one uninterrupted attempt.

That property — a failed unit of work can simply be redone, because no single machine’s continuity was ever a precondition for correctness — is the structural idea underneath almost everything that follows in this history, including where it eventually stops holding.

Coordination becomes its own layer

A fleet of machines that can each fail independently still needs some of them to agree on facts: who is the primary replica, which lease is current, which lock is held. Burrows’s 2006 paper on the Chubby lock service names the trade-off precisely: Chubby was built to provide “coarse-grained locking as well as reliable (though low-volume) storage for a loosely-coupled distributed system,” deliberately prioritizing “availability and reliability” over raw throughput, and by the time of publication several instances had each run for over a year serving tens of thousands of concurrent clients [4]. Chubby’s own advertised use — helping systems like GFS elect a primary and agree on it durably — is a tell about how the field’s thinking evolved in just three years: by 2006, coordination among failure-prone machines was no longer folded into each application. It had become a shared, separately engineered layer, because getting distributed agreement right once and reusing it was cheaper than each system inventing its own answer to “which replica do I trust right now.”

A wide archive plan-chest drawer holding a nine-track tape reel, a removable disk pack and a rack-mount RAID caddy in a row, one caddy's carrier only half seated on its felt cradle
Figure 1. Replicated, chunked storage did not replace the tape reel and the disk pack; it added a third specimen built on the same conviction that any one of them could fail without the data being lost.

The scheduler becomes the fleet’s operating system

Storage and coordination solved where data lives and who agrees on what. A separate and equally old problem is which job runs on which machine, right now, out of everything competing for the fleet. Google’s answer, described publicly in 2015 after roughly a decade of internal operation, was Borg: a system that “runs hundreds of thousands of jobs, from many thousands of different applications, across a number of clusters each with up to tens of thousands of machines,” achieving high utilization through “admission control, efficient task-packing, over-commitment, and machine sharing with process-level performance isolation” [5]. The paper’s own framing is instructive: a scheduler at this scale is not a convenience for allocating spare cycles. It is the mechanism by which a fleet’s aggregate utilization — how much of the hardware anyone paid for is doing useful work at any given moment — becomes a property that can be engineered and improved, rather than an accident of whichever team happened to grab a machine first.

ADVERTISEMENT

The field’s own sense that this had become a distinct discipline shows up in the parallel codification of warehouse-scale computing as a named unit of design. Barroso, Hölzle and Ranganathan’s The Datacenter as a Computer, now in a third edition, treats the datacenter not as a room full of separate machines but as a single designed computer whose “processor” is the fleet-wide scheduler and whose reliability properties have to be engineered at the level of the whole warehouse rather than any one rack [7]. That the same institutional lineage of scheduling ideas was worth re-explaining in three editions across roughly a decade is itself evidence of how quickly the discipline’s own vocabulary and assumptions kept needing to be restated for a widening audience.

By 2016, the scheduling ideas developed inside Borg were explicitly being generalized beyond any one operator. Burns, Grant, Oppenheimer, Brewer and Wilkes described “lessons learned from three container-management systems over a decade” — Borg, its successor Omega, and the externally released Kubernetes — as a single continuous line of thinking about cluster scheduling, with each system correcting specific weaknesses observed in its predecessor’s production use [8]. The detail worth noting historically is not any single technical choice but the direction of travel: scheduling knowledge that had been proprietary internal infrastructure in 2003 was, within roughly thirteen years, being published as reusable, portable, open systems engineering.

Telemetry has to catch the failures that never crash

Everything so far concerns a binary notion of failure: a machine is up or it is down, a task completed or it must be re-run. Dean and Barroso’s 2013 paper “The Tail at Scale” identified a different, harder failure mode. In a system built from thousands of components, each with a small independent probability of a slow response rather than a crashed one, “temporary high latency episodes which are unimportant in moderate size systems may come to dominate overall service performance at large scale” — because a single slow component among thousands can delay an entire aggregate response even though nothing actually failed in the conventional sense [6]. Their proposed mitigations, hedged requests and tied requests that let a system route around a component that is merely lagging, only make sense once telemetry itself has been engineered to notice degradation rather than only outage. A monitoring system built to watch for crashes will not see this failure mode at all; it has to be built to watch a distribution, not a boolean.

This is the point in the history where “failure domain” stops meaning only “which machines share a power feed or a network switch” and starts also meaning “which requests share exposure to the same slow, still-technically-alive component.” Both meanings persist in cluster systems engineering today, and conflating them is a documented source of avoidable outages.

A period-correct console terminal open on an examination stand with a faint trace mid-scroll on its small screen, a service manual propped open beside it with one page still turning
Figure 2. A component does not have to stop to hurt a service; a fleet's telemetry had to learn to notice the machine that only slowed down.

Standardizing how the fleet gets built

Scheduling and telemetry assume a fleet already exists. Getting a fleet built identically, many times, at multiple sites, is a separate and older-fashioned engineering problem: construction and procurement standardization. Facebook’s April 2011 announcement of the Open Compute Project made this explicit by publishing the hardware specifications behind its first custom-built facility — covering servers, racks and building design — stating the goal of treating datacenter hardware design the way open-source software treated code, so that the same specification could be built again by other operators rather than reinvented at each new site [14]. Read historically, the significant claim is not any particular design choice in the 2011 specification. It is the idea that a fleet is not usefully understood as one building, repeated by trial and error; it is one specification, built many times, with the specification itself treated as an engineering artifact subject to revision, versioning and reuse — the same discipline construction engineering had long applied to any standardized structure, now applied to a computing fleet.

A construction blueprint for a modular server hall half-unrolled across an archive desk, held down at one corner by a standardized rack-frame specimen, the far half still curled
Figure 3. A fleet is not one building repeated; it is one specification built many times, and the specification is as much an artifact of this history as any rack.

GPU clusters break the scheduling assumptions built for web services

Everything described so far was built to schedule short, largely independent units of work — a map task, a search query, a web request — across a fleet where any one machine’s temporary absence barely mattered to the whole. Deep learning training jobs on GPU clusters do not look like that, and the mismatch is documented directly in the first large published trace of such a cluster.

ADVERTISEMENT

Jeon and colleagues’ 2019 analysis of Microsoft’s Philly cluster, covering a 75-day trace from October to December 2017 across 96,260 jobs on fourteen virtual clusters, found that around 30% of jobs were killed or finished unsuccessfully due to failures, that even single-GPU jobs averaged only about 52% GPU utilization, and that distributed jobs spanning multiple servers degraded further — one measured configuration fell to roughly 28% utilization — because gang-scheduling and data-locality constraints meant a job needed a specific, simultaneous set of co-located machines rather than any available capacity [9]. None of Borg’s core assumptions — that jobs are largely independent, that partial progress on one machine has standalone value, that a scheduler can pack many small heterogeneous jobs onto shared hardware — held for this workload. A distributed training job wants a fixed set of GPUs, all present at once, all advancing in lockstep, for as long as the job runs; losing one of them does not shrink the job’s throughput proportionally, it stops the job.

Gandiva, published at OSDI the year before the Philly trace analysis, was a direct response to that mismatch. Xiao and colleagues built a scheduler that exploited the predictable, cyclical structure of deep learning training — GPU memory and compute use following a repeating pattern aligned to each mini-batch — to time-slice multiple jobs onto the same GPUs and migrate jobs toward better-fitting machines using introspection on their live progress rate, reporting a 26% improvement in aggregate utilization on a real 180-GPU workload [10]. Gandiva and the Philly trace analysis together mark the moment cluster scheduling split into two lineages: one still optimized for many small, independent, interruptible jobs, and a second built specifically around synchronous, gang-scheduled, topology-sensitive training runs. Frontier-scale training infrastructure descends entirely from the second lineage.

A 2010s-era GPU accelerator tray being set into a specimen row on the examination stand beside an older network interface card, its edge connector not yet seated
Figure 4. A GPU training job does not behave like the web requests the earlier schedulers were built for; it wants a fixed set of machines, all at once, all in step, for as long as it takes.

Checkpointing crosses from convenience to survival mechanism

If a training job cannot simply be re-run task-by-task the way a MapReduce job could, recovering from a failure depends entirely on how recently and how cheaply its state was saved. Checkpointing research therefore had to become a first-class systems problem in its own right rather than a background convenience, and its mathematics is old. A 2024 re-derivation of the classical result on checkpoint scheduling shows that the loss-minimizing interval between checkpoints is proportional to the square root of the product of checkpoint save time and mean time to failure,

τopt≈2 δ M, \tau_{\mathrm{opt}} \approx \sqrt{2\,\delta\,M}, ↗

with δ\delta↗ the time cost of writing one checkpoint and MM↗ the mean time between failures — and the paper notes that its simplified derivation reproduces the same interval as the original 1974 analysis it revisits [15]. The relationship exposes the exact trade-off cluster engineers have faced at every scale since: a checkpoint that is slow to write has to be taken less often, which increases the average work lost per failure, while a fleet whose components fail more frequently needs checkpoints taken more often regardless of how expensive each one is. As MM↗ shrank with fleet size, δ\delta↗ had to shrink to match, or lost work would grow without bound.

Meta’s Check-N-Run system, described at NSDI in 2022, is a direct instance of that pressure applied to production recommendation-model training. Because checkpoint frequency was bottlenecked by storage write bandwidth and by network bandwidth to remote storage, Eisenman and colleagues built incremental checkpointing that saves only the modified portion of a model together with quantization that shrinks checkpoint size without degrading accuracy, together reducing required write bandwidth by 6 to 17 times and storage capacity by 2.5 to 8 times across real production models [13]. Check-N-Run is the hinge point in this history: it treats checkpointing not as a batch-job convenience inherited from earlier systems but as a scarce resource to be engineered directly, in a system that still predates the frontier-scale synchronous training runs described next.

Frontier-scale training turns fault tolerance into a continuous operating condition

At the scale reached by 2024, checkpointing and failure recovery stopped being an occasional concern and became a continuously operating subsystem of training itself, documented in two independent operator reports.

ByteDance’s MegaScale, presented at NSDI 2024, describes training a 175-billion-parameter model across 12,288 GPUs at 55.2% model FLOPs utilization — a 1.34-times improvement over the Megatron-LM baseline it compares against — using a full-stack co-design of algorithms, communication overlap and diagnostic tooling [11]. Its checkpointing is deliberately two-stage: workers first write GPU state to host memory in a few seconds using pinned memory and optimized serialization, after which training resumes almost immediately, while a background process asynchronously moves that state to distributed storage without blocking further progress. Recovery is similarly engineered rather than incidental — a single designated worker rereads a shared state partition from storage and rebroadcasts it to the rest of its group, reducing load roughly linearly with group size. Paired with a heartbeat monitoring system, millisecond-level network diagnostics and lightweight self-check tests run across the fleet, the reported result on a separate multi-week production run of more than 10,000 GPUs was over 100 restarts due to failures, more than 90% of software and hardware faults automatically identified and fixed, detection-plus-diagnosis in under 10 minutes, recovery to pre-crash training progress within 15 minutes of the latest checkpoint, and an overall effective training time rate above 90% [11].

Meta’s account of pretraining its 405-billion-parameter Llama 3 model gives an independently sourced, comparably precise picture from a different operator. Training ran on up to 16,384 H100 GPUs, and across one 54-day snapshot of that run there were 466 total job interruptions: 47 planned, tied to automated maintenance such as firmware and kernel upgrades, and 419 unexpected. Of the unexpected interruptions, about 78% were attributed to confirmed or suspected hardware issues — GPU failures accounted for 148 interruptions (30.1% of the unexpected total) and GPU HBM3 memory failures for 72 more (17.2%), with software bugs and network issues making up most of the remainder — and despite that rate of roughly one interruption every three hours, the team reports achieving higher than 90% effective training time, with significant manual intervention required only three times across the entire snapshot [12]. The paper attributes checkpoint storage to Tectonic, Meta’s distributed file system, running at this scale across roughly 7,500 servers with about 240 petabytes of usable storage and sustained throughput on the order of 2 terabytes per second, with peak throughput considerably higher — the same distributed-storage philosophy first published as GFS two decades earlier, rebuilt at a throughput GFS’s original authors did not have reason to anticipate [12].

A modern checkpoint-storage blade being lowered into the archive drawer beside the nine-track tape reel and disk pack, its connector edge not yet mated
Figure 5. The newest specimen in the drawer answers the same question the oldest one did: where does the work go when the machine underneath it disappears.

What changed, and what stayed the same

The following synthesis is analysis built on the sourced facts above, not a claim attributed to any single source. Two things are true at once about this history, and neither is the whole story on its own.

The philosophy did not change. Every system described here, from GFS in 2003 to MegaScale and the Llama 3 infrastructure report in 2024, makes the identical founding move: treat failure as a statistical certainty rather than an emergency, and put the engineering effort into surviving it cheaply rather than preventing it entirely. That is a continuous, twenty-year line of thought, not a sequence of unrelated inventions.

What changed is the unit of recovery, and it changed because the workload changed. MapReduce could recover from a failure by re-running one independent task, because a map or reduce task’s correctness never depended on any other task’s continued existence — the workload was, in the literal sense, embarrassingly parallel. Gradient-synchronous training on thousands of tightly coupled GPUs has no equivalent independently re-executable unit: every rank in a collective operation is waiting on every other rank, so a single failed GPU does not shrink the job, it stops it, and recovery means restoring the entire cohort’s state together from a shared checkpoint rather than re-running one small piece of work. That structural difference is why checkpoint-write latency, once a background housekeeping concern, became a quantity worth engineering down to single-digit seconds at the scale documented in MegaScale, and why an operator now reports interruption counts and causal attribution with the same rigor a earlier generation reported job-completion statistics.

The other change is compression of timescale rather than any change in kind. A one-thousand-machine GFS cluster in 2003 could treat failure as a background statistical fact, absorbed invisibly by chunk replication. A 16,384-GPU training run in 2024 experienced an interruption roughly every three hours — a rate high enough that fault tolerance stopped being a subsystem invoked occasionally and became, in practical terms, a continuously active operating condition of the job for its entire duration. The engineering response to that compression — sub-15-second checkpoint writes, sub-15-minute full-cohort recovery, automated fault classification — is this history’s most recent chapter, and it is a direct, traceable descendant of the 1974 checkpoint-interval mathematics and the 2003 decision to stop being surprised by failure.

Predictions, with the observations that would falsify them

These are forecasts, clearly separated from the sourced history above. Horizon: 15 August 2030.

One. Reported effective-training-time percentages and interruption-cause breakdowns, in the style of the MegaScale and Llama 3 disclosures, will become a standard element of frontier training reports rather than an occasional disclosure, because comparing infrastructure claims without them will have become indefensible. Disconfirmed if leading operators in 2030 still publish training-run announcements with no failure or effective-time statistics at all.

Two. Checkpoint-write latency for frontier-scale synchronous training will continue falling toward the low single-digit seconds per step regardless of further growth in model or cluster size, because the Young-Daly relationship between checkpoint cost and failure rate makes any slower write economically unsustainable as clusters grow. Disconfirmed if published checkpoint-write times for state-of-the-art training runs in 2030 are materially slower, in wall-clock terms, than those reported for 2024 systems at comparable GPU count.

Three. The scheduling lineage that split at Gandiva and the Philly trace analysis — general multi-tenant cluster scheduling versus synchronous gang-scheduled training — will remain two distinct disciplines with separate tooling, rather than converging on one scheduler serving both well. Disconfirmed if a single scheduling system becomes the dominant choice for both general cloud workloads and large synchronous training jobs at major operators by 2030.

Four. Published fleet-construction specifications, in the Open Compute Project tradition, will extend further into software-defined operational practice — standardized runbooks and failure taxonomies, not only hardware — because the reliability-engineering value of shared vocabulary applies as much to failure classification as to physical design. Disconfirmed if failure-cause taxonomies remain operator-specific and mutually incompatible across major cluster operators through 2030.

None of these requires a discontinuity in the underlying technology. Each is a continuation of a pattern already visible across the sourced material above.

What to take away

Cluster systems engineering is not an invention of the GPU era; it is a twenty-plus-year-old discipline that the GPU era inherited and then placed under far more pressure. The founding move, dated to 2003 through 2006, was to stop engineering individual machines for trust and start engineering fleets for statistical survival — in storage, in coordination, in scheduling, and eventually in how the physical fleet itself gets specified and built. GPU-based deep learning training did not discard that discipline; it broke one of its load-bearing assumptions, that a failed unit of work could simply be re-run in isolation, and every system described in this history’s second half — Gandiva’s introspective scheduling, Check-N-Run’s incremental checkpoints, MegaScale’s two-stage checkpoint pipeline, Meta’s Tectonic-backed recovery at Llama 3 scale — is a documented, dated response to that break. The chips will keep changing. The question this discipline exists to answer — what does the system do in the next several seconds after a component it was depending on disappears — has not changed since 2003, and there is no sourced reason in this history to expect it to stop being the central question next.