Four choices, not one architecture
“AI datacenter systems engineering” sounds like it names a single discipline with a settled toolkit. It does not. Underneath the phrase are a handful of recurring design decisions that every operator of a large training cluster has to make, and on which serious engineers genuinely disagree — not because someone has not yet found the right answer, but because the right answer depends on the shape of the workload, the failure economics of the hardware, and how much spare capacity the operator is willing to pay for in advance.
This article works through four of those decisions using the primary systems-research literature and the engineering reports operators have themselves published: whether a scheduler should be a single authority or a negotiation among peers; whether fault tolerance should mean freezing a job’s state and restarting it, or running redundant copies of the work so a failure costs nothing to recover from; whether training data and checkpoints belong on a parallel filesystem built for throughput or an object store built for durability and elasticity; and whether a fleet should hunt for failing hardware before it is scheduled work, or wait for a job to fail and diagnose afterward. In each case the goal is not to declare a winner. It is to show the actual tradeoff — the numbers, the cost structure, the workload assumption — that makes each side of the argument a reasoned position rather than a preference.
A caution before starting. None of the systems compared below were built to solve the same problem at the same scale, so a benchmark number from one is not a fair stand-in for a benchmark number from another; the comparisons that follow are architectural and economic, not a leaderboard.
Centralized arbitration versus decentralized negotiation
The oldest and most consequential of the four decisions is who gets to decide where a job runs.
Google’s Borg, described in the systems-conference paper that remains the reference account of the design, is unambiguously centralized: a single logical scheduler admits, places, and manages “hundreds of thousands of jobs, from many thousands of different applications, across clusters with up to tens of thousands of machines” [1]. Centralization here buys a global, consistent view of the cluster — the scheduler can pack bin-fitting decisions, priority preemption, and resource isolation against one authoritative picture of what is free, rather than against a picture that might already be stale by the time two independent deciders act on it.
That consistency has a cost, and Google’s own follow-on research is explicit about it. Omega, published two years before the Borg paper, argues that “increasing scale and the need for rapid response to changing requirements are hard to meet with current monolithic cluster scheduler architectures,” which “restricts the rate at which new features can be deployed” and “decreases efficiency and utilization” as a single scheduling loop becomes the bottleneck through which every decision must pass [2]. Omega’s answer was not to abandon centralization but to soften it: multiple schedulers operate in parallel against one shared, replicated view of cluster state, using optimistic concurrency control to resolve the conflicts that arise when two schedulers reach for the same resources at once. It is a shared-state design, not a decentralized one — every scheduler still sees the whole cluster — but it trades some of Borg’s strict serialization for parallelism.
Sparrow goes further and abandons the shared view entirely. Its decentralized architecture uses batch sampling, in which each of many independent schedulers probes a small random subset of worker machines rather than consulting a global picture, combined with late binding, in which a task is not committed to a worker until that worker actually becomes free [3]. On a 110-machine cluster the authors report performance “within 12% of an ideal scheduler” with none of the central-bottleneck risk that motivated Omega — but Sparrow was built and evaluated for a specific workload shape: very large numbers of very short tasks, the pattern typical of interactive data-analytics queries, where scheduling latency itself is the thing being optimized.
That last detail is the hinge the whole comparison turns on. Large AI training jobs are the near-opposite workload: not many short tasks but one enormous, long-running, tightly synchronized job whose thousands of ranks must be placed with attention to network topology and locality, and where the cost of a scheduling decision made in milliseconds is irrelevant against a job that will run for weeks. Empirical work on production GPU clusters bears this out directly: an analysis of Microsoft’s Philly clusters found that “deep learning frameworks require gang scheduling,” which “reduc[es] the flexibility of scheduling and mak[es] the jobs themselves inelastic to failures at runtime,” and that locality — how close together in the network a job’s GPUs are placed — has a first-order effect on both queueing delay and achieved utilization [4]. A scheduler optimized to minimize per-task placement latency, Sparrow’s explicit design target, is optimizing the wrong variable for this workload.
Consistent with that, the schedulers hyperscale operators have built specifically for AI training have moved toward more centralization and more global scope, not less. Meta’s MAST performs scheduling “at a global scale” across geo-distributed datacenters using what its authors call temporal decoupling, scope decoupling, and exhaustive search, and reports reducing the demand-to-supply ratio in the most overloaded region from 2.63 to 0.98 for high-priority workloads by giving one system visibility across regions that users previously had to choose between manually [5]. Microsoft’s Singularity likewise centralizes scheduling decisions across “a global fleet of AI accelerators,” making every job “preemptable, migratable, and dynamically resizable by default,” so that “a live job can be dynamically and transparently preempted and migrated to a different set of nodes, cluster, data center or a region and resumed exactly from the point where execution was preempted” [6]. That is a vendor’s own description of its system’s capability, not an independently measured result, and it should be read as such — but the architectural choice it reflects, favoring one system with global visibility over many independent local deciders, tracks the same logic Jeon and colleagues’ trace analysis suggests: at training scale, coordination quality dominates decision latency.
Checkpoint-and-restart versus redundant execution
The second decision is what to do when a GPU, a node, or a network link actually fails mid-job — and at the scale modern training runs operate, failure is not an edge case but a certainty within the run.
The dominant approach, synchronous checkpoint-and-restart, periodically writes the full training state to durable storage and, on failure, rolls the whole job back to the most recent checkpoint. Meta’s infrastructure team, analyzing eleven months of data across two large research clusters totaling more than 150 million A100 GPU-hours, quantifies exactly why this approach becomes harder to sustain as jobs scale. They report that “the mean-time-to-failure (MTTF) of 1024-GPU jobs is 7.9 hours — roughly two orders of magnitude lower than 8-GPU jobs (47.7 days),” and their projections put MTTF at 1.8 hours for 16,384-GPU jobs and 0.23 hours for 131,072-GPU jobs [7]. Failure rate does not merely rise with scale; it rises fast enough that the interval between checkpoints has to shrink correspondingly, and every checkpoint write is itself dead time in which the job produces no forward progress.
The authors frame that tradeoff with the classical Young–Daly optimal-checkpoint-interval model, which balances the fixed cost of writing a checkpoint against the expected cost of losing work to a failure:
where
One answer to that requirement is to make the checkpoint itself faster rather than abandon the approach. Gemini, built by researchers at Rice University and Amazon, checkpoints to the CPU memory of nearby hosts instead of remote persistent storage, using “a provably near-optimal checkpoint placement strategy” together with traffic scheduling that interleaves checkpoint transfers with ongoing training communication, and reports failure recovery “more than 13 times faster than existing solutions” in the authors’ evaluation [8]. Amazon’s SageMaker HyperPod applies the same principle in a managed product, storing frequent checkpoints in CPU RAM for fast recovery while “periodically” replicating to Amazon S3 for durability [15] — treating checkpoint-and-restart as a problem to be optimized rather than replaced.
The competing answer abandons synchronous restart altogether in favor of redundancy: keep more than one copy of the computation running, so a failure never stops forward progress at all. Oobleck generates multiple pipeline-parallel configurations with replicated pipeline segments in advance and, on a node failure, switches to a smaller configuration built from the surviving replicas rather than restarting from a checkpoint, a design its authors say “mathematically guarantees that available resources can always be fully utilized after failures occur,” reporting recovery up to 29.6 times faster than the checkpoint-based systems they compare against in their own benchmarks [9]. Bamboo applies a related idea to cheap, preemptible cloud instances specifically: nodes perform “redundant computations … over not only its own layers but also over some layers in its neighbor,” so that when a preemptible instance is reclaimed its neighbor already holds a working copy of its state, which the authors report yields 3.7 times higher training throughput and 2.4 times lower cost than relying on on-demand instances for the same job [10].
The honest comparison is a cost-structure tradeoff, not a technical superiority claim in either direction. Checkpoint-and-restart, even optimized to Gemini’s in-memory speed, pays its cost only when a failure actually happens — but per Kokolis and colleagues’ own numbers, at extreme scale that cost is paid often and the mean time between payments keeps shrinking. Redundant execution pays a standing tax on every single step of training, in the extra GPU-hours the shadow computation consumes, whether or not a failure ever occurs — but that tax is fixed and predictable, and it converts an open-ended tail risk into a known, budgeted overhead. Unicron, from Alibaba, sits between the two: it does not run redundant compute continuously, but adds “in-band error detection for real-time error identification without extra overhead” plus “a dynamic cost-aware plan generation mechanism” that reconfigures a still-running job around a failure faster than a full checkpoint restart would, reporting up to 1.9 times higher training efficiency than the checkpoint-restart baselines it compares against on a 128-GPU cluster [11]. Which of the three postures — pure checkpoint-restart, pure redundancy, or fast reconfiguration around a partial failure — is worth its overhead depends on how expensive spare capacity is relative to how expensive downtime is for that specific job, and reasonable infrastructure teams reach different answers to that question. Consistent with this, Shanghai AI Laboratory’s six-month characterization of LLM development workloads across two clusters totaling 4,704 A100 GPUs reports building their own fault-tolerant pretraining approach specifically because generic checkpoint-restart did not fit the failure pattern they actually observed in production [12] — a reminder that the right answer is also a function of the operator’s own measured failure distribution, not a property of the technique in the abstract.
Parallel filesystems versus object storage
The third decision is where training data and checkpoints physically live, and it is less a binary choice than the first two — in practice, most large operators end up with both, at different tiers.
The case for a purpose-built distributed filesystem is throughput and latency under exactly the access pattern GPU training produces: many workers reading large numbers of files, or writing large checkpoints, concurrently and repeatedly. Meta’s Tectonic, described in a systems-conference paper on the design, “consolidates large tenants that previously used service-specific systems into general multitenant filesystem instances” at a scale where “each cluster can be multiple exabytes … large enough to serve an entire Facebook-sized data center,” achieved by partitioning metadata across many independent metadata servers rather than routing every request through one bottleneck [13]. Amazon’s FSx for Lustre makes the equivalent commercial case: the product is described as delivering “sub-millisecond latencies” and “up to hundreds of gigabytes per second of throughput,” explicitly positioned for “ML data loading, model checkpointing, inferencing operations, and key-value caching” [16].
Object storage optimizes for a different set of properties: durability, elastic capacity that does not require pre-provisioning a filesystem’s worth of dedicated hardware, and a simpler operational model at rest. The tension between the two is not hypothetical or vendor-marketing color — it is the reason an industry benchmark exists specifically to measure it. MLCommons introduced the MLPerf Storage suite as “the first artificial intelligence/machine learning benchmark suite that measures the performance of storage for machine learning workloads,” built around an emulation mechanism that reproduces realistic training I/O patterns without requiring the accelerators themselves, precisely because storage system choice measurably changes achieved training throughput and no single number from either architecture generalizes across workloads [14].
In production practice, the two are increasingly not a choice but a stack. Amazon’s own guidance describes FSx for Lustre file systems that “link” to S3 buckets, transparently presenting S3 objects as files and lazily loading data from S3 only as it is actually needed, so that a training job can begin “without needing to download the full … training dataset from S3” while the durable copy of record stays in the cheaper, more elastic object store [16]. SageMaker HyperPod’s managed tiered checkpointing follows the identical pattern in the other direction, for writes rather than reads: checkpoints land first in fast CPU memory, replicate asynchronously across nearby nodes, and only then get written back to S3 for long-term durability [15]. The genuine disagreement here is not filesystem versus object store as competing final answers — it is how much tiering complexity is worth taking on. A smaller operation may reasonably decide a single, sufficiently large parallel filesystem is simpler to run correctly than a two-tier pipeline with its own consistency and staleness questions; a hyperscale operator may reasonably decide the throughput and cost gap is too large to leave on the table. Both are defensible engineering judgments, not a resolved question.
Proactive health-checking versus reactive failure handling
The fourth decision concerns the fleet as a whole, not any single job: should an operator actively hunt for failing hardware before it is ever handed work, or wait for something to break and then diagnose it?
Meta’s reliability study is the clearest and most quantified account of the proactive position, and it is worth reporting its numbers precisely because the effect size is large. The team runs health checks “periodically scheduled to run every five minutes,” screening for GPU errors, filesystem mount failures, and service problems, under an explicit stated philosophy: “a core philosophy underlying our training cluster design is to strive for no second job failure from a bad node — a failed node is a poor scheduling candidate” [7]. Applying this system to identify chronically faulty “lemon” nodes — nodes that fail repeatedly and disproportionately, comprising in one of their clusters just 1.2% of the machine footprint but roughly 13% of daily job submissions — and removing them from scheduling ahead of time, the authors report the failure rate for large jobs (512 or more GPUs) fell from 14% to 4% [7]. NVIDIA’s DCGM tooling operationalizes the same proactive philosophy as standardized, vendor-supplied infrastructure: it documents three distinct check tiers — non-invasive background checks that run continuously alongside live jobs, fast “prologue” checks that verify a GPU’s readiness in the seconds before a job is scheduled onto it, and deeper “epilogue” diagnostics that run for a few minutes after a job fails or a GPU’s health looks suspect — explicitly framed as catching “unresponsive GPUs, corrupted firmware, thermal escapes” before they cause a job failure rather than after [17].
Set against this is a genuinely different, and much older, tradition in general-purpose site reliability engineering, and it is not a naive or under-resourced alternative — it is Google’s own long-standing published doctrine for its production services. The Site Reliability Engineering book’s chapter on monitoring is explicit that black-box, symptom-triggered monitoring “has the key benefit of forcing discipline to only nag a human when a problem is both already ongoing and contributing to real symptoms,” and the book’s broader alerting philosophy insists that every page a human receives must be “actionable” and reflect something “actively or imminently user-visible” — a deliberate rejection of complex predictive or proactive alerting on the grounds that it degrades signal quality and burns out on-call engineers with false positives [18].
These are not simply two teams that have not talked to each other; they are two philosophies calibrated to different failure economics, and characterizing that difference is more useful than declaring one obsolete. General-purpose service reliability, the world the SRE book was written for, typically has redundancy that absorbs a single failing component without a user ever noticing, so a false-positive proactive alert costs an engineer’s attention for nothing — a real and recurring expense across a fleet of thousands of independently-monitored components. Large synchronous training jobs invert that arithmetic. A single degraded GPU among thousands of tightly coupled, gang-scheduled ranks does not fail gracefully in isolation; per Jeon and colleagues’ characterization of gang scheduling’s rigidity, it can silently slow or corrupt an entire run of tens of thousands of GPU-hours before any symptom clear enough to trigger a reactive page ever appears [4] [7]. Under those economics, a five-minute proactive health check that occasionally flags a node unnecessarily is cheap insurance against a failure mode whose reactive alternative is discovering the corruption only when the job’s loss curve or a collective-operation timeout finally makes it visible, hours or days later. The disagreement, in other words, is not about which philosophy is correct in general — it is about how expensive a false positive is relative to how expensive an undetected failure is, and that ratio is a property of the workload, not a universal constant.
Where the disagreement is real
Pulled together, a pattern runs through all four comparisons. In each pairing, one side optimizes for a world of many small, independent, loosely coupled units of work, and the other optimizes for a world of few, enormous, tightly synchronized ones — and the reason serious engineers land on different sides is that both worlds genuinely exist inside the same industry, sometimes inside the same company.
Sparrow’s decentralized scheduling and the SRE book’s reactive monitoring both come from a lineage built around large numbers of short, recoverable, loosely coupled tasks, where a bad decision or a missed check affects one small thing and self-corrects quickly. Borg’s centralization, MAST’s global view, Meta’s proactive lemon-node hunting, and the redundant-execution systems all come from a lineage built around jobs so large, so synchronous, and so expensive per hour that a single missed signal or a single suboptimal placement can waste an amount of compute that would fund an entire team’s worth of careful engineering. Neither lineage is wrong. The disagreement that is real, and worth naming rather than resolving, is about which lineage a given piece of infrastructure actually belongs to — and that is a judgment call operators are still making differently from each other in production today.
What follows for someone actually building this
None of the four axes reduces to a checklist, but the sourced material above supports a few honest heuristics rather than rules.
Scheduler centralization earns its coordination overhead in proportion to how synchronous and how network-topology-sensitive the workload is; a cluster running many small, independent jobs has less to gain from a single global arbiter than one running a handful of enormous, gang-scheduled training runs. Checkpoint-and-restart remains the more capital-efficient default when spare, currently-idle capacity is scarce or expensive, because it pays its cost only on failure; redundant execution becomes more attractive specifically on preemptible or otherwise cheap capacity, where the “wasted” shadow computation was inexpensive to begin with, as Bamboo’s own economics illustrate directly [10]. Storage tiering is worth its complexity in rough proportion to cluster size, because the throughput gap between a parallel filesystem and direct object-store access widens as concurrent readers multiply, while the operational cost of running a second storage tier is closer to fixed. And proactive health-checking earns its false-positive cost specifically where job restarts are expensive and gang-scheduled; it is a much harder sell for infrastructure where individual task failures are cheap and self-healing already.
Predictions, with what would falsify them
These are forecasts, kept separate from the sourced findings above. Horizon: August 2029.
One. Tiered storage architectures — a fast parallel-filesystem or in-memory cache layer in front of durable object storage, for both training data and checkpoints — will be the default reference design for new large training clusters, rather than either a pure-filesystem or pure-object-store design. Disconfirmed if the majority of newly built large clusters standardize on a single storage tier for both data ingestion and checkpointing.
Two. Proactive, pre-job and in-job GPU health-checking will become near-universal specifically for large, gang-scheduled training clusters, even as general-purpose cloud services continue to rely on the SRE book’s reactive, symptom-triggered philosophy as their default. Disconfirmed if large training operators are still predominantly relying on purely reactive, post-failure diagnosis for fleets of thousands of GPUs.
Three. Redundant or replicated execution will remain a minority pattern concentrated on preemptible and opportunistic capacity, while optimized checkpoint-and-restart remains the default fault-tolerance strategy for dedicated, reserved training capacity. Disconfirmed if redundant execution becomes the standard fault-tolerance approach on dedicated, non-preemptible large training clusters rather than checkpoint-restart.
What to take away
Every one of these four comparisons resolves the same way on inspection: not into a winner, but into a clearer picture of which variable each approach is actually optimizing, and which workload makes that variable the one that matters. A scheduler built for millisecond task latency and a scheduler built for coordinated placement across tens of thousands of synchronized ranks are not competing answers to the same question; they are answers to different questions that both happen to be called “scheduling.” The same is true of checkpoint-restart against redundancy, of parallel filesystems against object stores, and of proactive health-checking against reactive paging. Treating any one of these pairs as a solved question, with a single right answer transferable across every cluster, is the mistake — the actual engineering skill is diagnosing which side of each pairing a given workload’s failure economics put it on, and building accordingly.