Est.

Checkpoint Storage for Distributed Training on S3

Features Editor · · 11 min read
Cover illustration for “Checkpoint Storage for Distributed Training on S3”
Object Storage · September 30, 2026 · 11 min read · 2,409 words

Checkpoint storage on S3 rarely appears on the risk list until a training run at real scale hits it head-on, and by then it consumes the budget instead of trimming a footnote from it. A checkpoint isn't just "the model." It's the model weights, the optimizer state, the gradient accumulators, and the scheduler state, all bundled together so training can resume from a known-good point if something breaks. That bundling is where the surprise starts.

Adam-class optimizers store per-parameter first and second moment estimates, making optimizer state roughly 4x the weight size alone. So a model that sounds modest in parameter count turns into something much heavier on disk than anyone budgeted for.

Put numbers on it. A typical modern LLM checkpoint runs 350 to 500 GB of model state. A 70-billion-parameter model checkpoints at around 980 GB, essentially a terabyte, every single time it saves Efficient Training of Large Language Models on Distributed Infrastruc…. Push into trillion-parameter territory and a single checkpoint event can generate 15 TB or more. That's a terabyte-scale write operation, not a file. That's a small data center's worth of writes, all at once.

Frequency is not optional, either. At the scale Meta ran for Llama 3, hardware failures hit roughly every three hours across a fleet of thousands of GPUs. And checkpoints aren't only insurance against crashes. Teams use them to migrate platforms, to roll back after a loss spike sends training off the rails, or to seed a new objective off a partially trained model. Once it's clear what gets written and how often, the next problem writes itself: all of that data has to land somewhere, all at once, on a clock.

The burst write pattern at distributed training scale

That's the target. Missing it doesn't fail quietly.

The write pattern itself depends on how the training job is parallelized, and both flavors are stressful in their own way. In FSDP training, every GPU rank writes its shard simultaneously, producing massive write concurrency across hundreds or thousands of ranks.

That idle time is the real cost, and it compounds fast. Every second a checkpoint pause runs is multiplied by every GPU sitting there doing nothing, and at large cluster scale, that idle stretch becomes one of the most expensive recurring line items in the entire training run. Not a one-time tax. A recurring one, every single checkpoint interval, for the life of the job.

A concrete example makes the math land. Take a 512-GPU cluster producing 15 TB checkpoints, backed by 100 GBps of storage throughput. That works out to a 150-second write duration, which blows past a 10% overlap target teams typically aim for. And that's before accounting for file count. Modern 3D parallelism setups (splitting work across data, tensor, and pipeline dimensions) can generate hundreds or thousands of distinct files per checkpoint event, and most storage systems were never built to handle that much metadata contention at once. Storage isn't a one-time setup cost here. It's a tax collected on every single checkpoint, for the entire life of the training run. The time constraint requires each checkpoint write to complete in under 5 minutes to avoid meaningful training interruption, necessitating burst write bandwidth of 1–2 GB/s per checkpoint.

S3 Architecture and Failure Modes for Checkpoint Writes

S3 was built to serve high-throughput batch reads and big sequential objects extremely well. Checkpointing asks for the opposite: low-latency, bursty, fine-grained writes, all landing at once. The mismatch isn't a bug somebody forgot to patch. It's structural, baked into what S3 was designed to do in the first place.

Start with per-operation latency.

Then there's throttling. Concentrated writes from a large cluster can trip S3's rate limits, and once that happens, retries kick in and stretch the checkpoint window even further. It's a bit like everyone in a stadium trying to exit through the same single door at once: the door doesn't get faster just because more people are pushing on it.

Teams who assume S3 behaves like a normal filesystem get tripped up by atomic rename, which deserves its own mention. On S3, there's no such thing as an atomic rename. A rename is actually a copy followed by a delete, which doubles the I/O cost of every checkpoint finalization step. Append operations and write-ahead log patterns, common tools for safely finishing a write, are similarly unsupported or handled inconsistently.

None of this erases what S3 is genuinely good at. Those strengths just don't touch the latency. The net effect: an object storage operation that should take seconds can stretch into minutes under checkpoint burst conditions, turning what should be background infrastructure into the thing everyone's staring at on the dashboard. So the obvious next question: does AWS's own high-performance tier fix this? Per-operation latency means a standard S3 PUT incurs 30–40 ms of latency per request, so when checkpointing is a coordinated burst across all training ranks, every GPU waits through that latency. Geo-redundancy and cost are genuine strengths (S3 Standard at ~$23/TB/month is among the cheapest durable storage available), but those strengths do not help with write latency.

What S3 Express One Zone changes

It's a genuine architectural response to the exact problem checkpointing creates, not a rebrand of the same product.

The headline number: each request completes in roughly 4 milliseconds, and a single directory bucket handles high transaction volume without needing prefix partitioning tricks, which substantially cuts down the throttling behavior that plagues standard S3 under checkpoint bursts repost.aws. AWS also folded this into SageMaker directly, announcing in February 2024 that Express One Zone accelerates SageMaker Model Training with consistent single-digit millisecond latency for checkpoints and model outputs.

Cost sits in an interesting middle ground. Nobody sees a free lunch, though. The trade-off is durability: Express One Zone lives in a single Availability Zone. That's the whole point, keeping storage physically close to compute, but it makes the tier a poor fit for compliance-sensitive data or anything that needs multi-AZ durability.

And it doesn't fix everything. The atomic rename issue is still there repost.aws. Per-operation latency drops, but doesn't vanish relative to local NVMe, and for the largest bursts, the gap between 4 milliseconds and the sub-millisecond world of local NVMe still matters repost.aws. Express One Zone raises the ceiling. It doesn't close the gap entirely, and how much that residual gap matters depends heavily on what's happening at the framework level. AWS cites up to 10x faster data access and 80% lower request costs versus S3 Standard, with random-access checkpoint writes and metadata lookups being the workloads where those gains are most legible docs.aws.amazon.com. At $0.11/GB/month (~$110/TB/month), the cost is dramatically cheaper than FSx for Lustre ($143–600/TB/month) while offering substantially lower latency than S3 Standard ($23/TB/month) docs.aws.amazon.com sedai.io MLCommons.

How PyTorch DCP and async checkpointing change write timing

Async checkpointing tries to solve this from a completely different angle: don't make the GPU wait on storage at all. The model state copies to CPU RAM, training resumes immediately, and the actual write to storage happens in a background thread, decoupling GPU availability from when the write actually finishes.

Sharded checkpointing complements this. Instead of gathering everything to a single rank before writing, each node writes its own shard in parallel. AWS added Distributed Checkpoint support to its S3 Connector for PyTorch in November 2024, letting multiple training processes write to S3 in parallel and cutting write time for large foundation model checkpoints.

It gets less tidy, though. A 2026 microbenchmark trained a 1.216-billion-parameter causal decoder across four NVIDIA RTX 5880 GPUs on PyTorch 2.11.0, and the results complicate the "async is always better" story. Standard async using background threads was actually worse, at plus 1,043.5%, because scheduling contention ate the latency gains arxiv.org. Background processes came in even worse at plus 1,090.2%.

The takeaway isn't subtle: async checkpointing isn't automatically an upgrade. A naive implementation introduces scheduling contention that can outweigh whatever latency it was supposed to save. PyTorch's own team addressed this directly, and cached plans combined with reduced GIL contention in the DCP async path deliver a genuine 6x speedup, but that gain only appears on the optimized code path, not the one most teams reach for by default pytorch.org. There's also a memory cost nobody should skip. Async staging requires enough CPU RAM to hold at least one full checkpoint while the write completes, and for a 70B model at 980 GB, that's a non-trivial memory budget Efficient Training of Large Language Models on Distributed Infrastruc…. In sharded checkpointing, each node writes its own shard in parallel rather than gathering to a single rank; PyTorch FSDP supports sharded checkpoints but produces full checkpoints (gathered to rank 0) by default, and wall-clock write time falls proportionally to shard count. Baseline performance is 1.024 s/step arxiv.org. Synchronous checkpointing yields 9.960 s/step (+870.8%) arxiv.org. Cached-process staging yields 7.747 s/step (+656.7%), generating 3.39 TB/hour of write traffic.

Write-back NVMe caching as a storage-layer mitigation for S3 latency

A different fix works at the storage layer instead of the framework layer. Training ranks write to a local NVMe cache at sub-2-millisecond latency, and the cache asynchronously flushes that data to S3 in the background. GPUs unblock the moment the local snapshot completes, not when S3 confirms the write.

Prefix distribution pairs well with this approach. Spreading checkpoint writes across multiple S3 prefixes cuts the request rate hitting any single prefix, which lowers the odds of tripping bucket-level throttling. It's a small architectural tweak with an outsized payoff during burst writes.

GPU Direct Storage adds another layer of speed on top. Without it, data travels NVMe to CPU to PCIe to GPU. With GDS enabled, that path collapses to NVMe straight to PCIe to GPU, cutting out the CPU hop entirely. DeepSpeed's Async I/O module supports this natively through its aio configuration block.

The portability argument affects which architecture choices hold up across providers, and it may outweigh the raw speed numbers. Write-back caching works the same way across clouds, across GPU providers, and on-premises, with the S3 bucket staying the durable source of truth no matter which compute layer happens to be generating checkpoints that week. RDMA transport support for this pattern currently covers read I/O only in some implementations, with write paths still running over TCP and RDMA write support slated for later. A POSIX-mounted filesystem that caches writes locally to NVMe and flushes them asynchronously to the backing bucket delivers this pattern in practice: sub-millisecond local writes for the training job, S3 as the system of record, and no data migration step in between. The 30–40 ms S3 latency problem disappears from the GPU's perspective, as the storage system absorbs it behind the cache boundary MLCommons. The cost position of the NVMe cache tier is estimated by Alluxio at $11–42/TB/month, derived from GPU cloud instance NVMe costs (AWS ien.24xlarge, c5n.18xlarge) plus S3 Standard backend, substantially cheaper than FSx for Lustre at $143–600/TB/month while closing most of the latency gap MLCommons.

Placing S3 correctly in a multi-tier checkpoint pipeline

Zoom out and a cleaner mental model appears. Omdia analyst Dennis Hahn described training infrastructure as a three-stage pipeline, and it's a useful way to think about where each storage type actually belongs. The first stage is ingest and curation: cheap, cost-optimized capacity storage that absorbs raw unstructured data as it comes in. Stage two is data staging and modeling: high-throughput, low-latency performance storage that feeds GPUs during training and carries the active checkpointing load. Stage three is deployment: storage built for production-grade resiliency once the model is done training.

Stage two is where the hot tier lives, and it's usually either NVMe attached directly to GPU nodes or a parallel filesystem like Lustre, BeeGFS, or GPFS, typically delivering 5 to 20 GB/s https://sysart.consulting/insights/checkpoint-model-storage-on-premises/. The rule of thumb is 1 to 2 TB of NVMe per GPU to avoid bottlenecks.

S3 as a durable archive and S3 as a synchronous checkpoint write target are different jobs, and the teams that treat them as the same job are the ones who discover the bottleneck at scale. Teams that treat them as interchangeable are the ones who find the bottleneck the hard way, usually mid-training run.

The industry is starting to test this formally. MLPerf Storage v3.0, released in September 2026, added checkpointing as a first-class workload, testing models at 8B, 70B, 405B, and 1,250B parameters in Llama 3-style configurations, and for the first time added an S3 data access layer alongside the traditional POSIX path. Roughly one-sixth of v3.0 submissions used that S3 layer, which says something: S3 is earning legitimacy as a checkpoint storage path, provided it's architected correctly rather than bolted on as an afterthought. S3's correct position is archiving completed checkpoints, dataset snapshots, and older model versions, where its cost ($23/TB/month standard) and durability are decisive advantages and its write latency is irrelevant.

Diagram: Checkpoint Storage Cost vs. Latency: Three Tiers Compared. Visualizes: Show a ranked comparison of three checkpoint storage tiers on two dimensions — cost per TB/month and approximate write latency — using the exact figures from the…

Evaluating a checkpoint storage approach

The tension doesn't really resolve, it just gets managed. S3's cost and durability sit at the top of the class. Its synchronous write latency does not.

A handful of variables decide where a team should land on that spectrum. Checkpoint size and frequency matter first: a 70B-model checkpoint at 980 GB saving every few hundred steps demands a wildly different storage budget than a smaller model checkpointing once an hour Efficient Training of Large Language Models on Distributed Infrastruc…. Cluster size matters just as much, since the GPU idle cost that justifies NVMe caching scales with cluster size, and at smaller scale the math often favors just tolerating S3's latency.

What the MLPerf Storage v3.0 S3 submissions really signal is that the industry is actively measuring whether object storage, used properly, can close enough of the performance gap to justify its cost advantage, and results will inform architecture choices as they accumulate. A cloud filesystem that mounts S3-compatible buckets as POSIX, caches writes locally to NVMe, and flushes asynchronously back to the bucket. That combination pairs S3's cost floor with NVMe write speeds, runs training code unmodified, keeps the bucket as the single source of truth, and scales capacity without forcing a rebuild of the pipeline every time the cluster grows. Durability requirements mean single-AZ options like S3 Express One Zone (~$110/TB/month, 4 ms latency) are appropriate where multi-AZ durability is not a requirement for active checkpoints docs.aws.amazon.com repost.aws. On the framework side, PyTorch DCP with the optimized async path (6x faster than naive async) changes the calculus independently of the storage tier aws.amazon.com pytorch.org.

Sources

  1. The Storage that feeds AI training and modeling for High-Impact AI
  2. MLCommons Releases New MLPerf Storage v3.0 Benchmark Results - MLCommons
  3. Checkpoint and Model Storage Architecture for On-Premises AI - Sysart Systemic Agile Consulting | SysArt Consulting
  4. Amazon S3 Connector for PyTorch now supports Distributed Checkpoint
  5. Efficient Training of Large Language Models on Distributed Infrastructures: A Survey
  6. Amazon S3 Express One Zone: Key Insights for 2026 | Sedai
  7. Choosing an S3 connector for ML training with S3 Express One Zone | AWS re:Post
  8. Posted On: Feb 27, 2024
Filed underObject Storage

More in Object Storage