Est.

NVMe Caching Architecture for Object Storage Reads

GPU training wastes half its cycles waiting for data from cheap object storage.

Correspondent · · 11 min read
Cover illustration for “NVMe Caching Architecture for Object Storage Reads”
POSIX Filesystems · October 5, 2026 · 11 min read · 2,541 words

Object storage is cheap, durable, and nowhere near fast enough for a GPU that costs more per hour than most people's rent. That's the core problem this piece is about: a structural latency mismatch between S3-style storage and AI compute, and the NVMe caching architecture built to bridge it.

Object storage latency versus AI compute demands

A single S3 read typically takes somewhere between 50 and 200 milliseconds, depending on object size and conditions. A GPU doing training work burns through a batch in a fraction of that time, then sits idle waiting for the next one to arrive. Multiplying that wait by every batch in every epoch across a cluster full of GPUs makes the idle time stack up fast.

Meta's internal data found that almost 56% of GPU cycles sat stalled waiting for training data to arrive. That points to a structural problem with the request pattern itself, not tuning or bandwidth. Object stores can push enormous aggregate throughput when the request pattern is large, sequential reads. Small, random reads cause the trouble, the kind a computer vision job generates when it opens millions of individual image files one at a time. Each file costs another round trip, and round trips are where the 50-200 millisecond tax gets paid, over and over.

Checkpointing adds a second layer of pain on top. A large model's checkpoint file is big, and when it gets written through the same object-storage path the training reads are using, the two traffic types compete for the same pipe. One of them loses. It's usually the training run.

None of this gets fixed by buying faster GPUs or provisioning more bandwidth. The bottleneck sits in the architecture, in the physical distance and protocol overhead between where data lives and where compute happens. Closing that gap means putting something faster in between.

The three-tier cache stack between S3 and the GPU

The answer that's actually shipped in production AI clusters is a three-tier stack: object storage, NVMe, and GPU memory, each one closer to the compute than the last and each one holding less data than the one before it.

Object storage (S3, GCS, Azure Blob, R2, take your pick) holds the full dataset. It's durable, it's cheap per gigabyte, and it's where everything lives at rest. NVMe, whether it's bolted directly to the compute node or attached over a fabric, holds the active working set, the data a job is actually touching right now, at latency low enough that the GPU doesn't notice the fetch. GPU memory sits at the top of the pyramid holding only the current batch, because GPU memory is small and every gigabyte of it is doing expensive work.

Prefetch is what makes the stack function instead of just sitting there as a diagram. While the GPU chews through the current batch, a background thread is already pulling the next several batches out of object storage and staging them on NVMe. By the time the GPU finishes and asks for more, the data's already local. The GPU never has to wait on a cold read from object storage mid-run, because the system made sure it never would.

This isn't some exotic trick reserved for teams with a storage engineer on call. Plenty of training pipelines and data-loading libraries bake prefetch logic directly into the data loader, so the three-tier pattern happens by default rather than by heroic effort.

The reason the stack has three tiers instead of one comes down to cost. Object storage is cheap precisely because it's slow. NVMe is fast and expensive. GPU memory is both fast and scarce, so it only ever holds the smallest possible slice of data. No single tier can do every job at once, so the architecture spreads the work across all three, matching each tier's cost to how much speed the job actually needs from it.

Diagram: The Three-Tier Cache Stack: From Object Storage to GPU. Visualizes: Show a vertical three-tier pyramid (or layered stack) illustrating the caching architecture between S3-style object storage and the GPU.

Where the cache lives: placement decisions and their consequences

Knowing the stack has three tiers doesn't tell you where the middle one, the NVMe layer, should physically sit. That decision shapes the latency floor, how the system fails, and whether the cache can grow with the workload.

Node-local NVMe is the simplest option: the drive lives in the same box as the GPU, so reads are about as fast as reads get. The cache is private to that node. If ten nodes are all training on overlapping shards of the same dataset, each one keeps its own separate copy, which is a duplicated cache spread across a cluster. Lose the node to preemption or failure, and that node's cache disappears with it.

NVMe over Fabrics, running over high-speed networking like RoCE on 100 gigabit ethernet, offers a different shape. Remote NVMe drives show up to the compute node looking like local PCIe devices. Multiple nodes can read from one shared cache instead of each keeping a private copy. NVMe/RoCE leans on RDMA to move data without routing it through the OS kernel or the CPU, so latency stays well below what a TCP-based setup delivers. NVMe/TCP trades some of that speed away in exchange for running on standard network gear instead of specialized RDMA-capable hardware.

Which one fits depends on the job. Distributed training, where dozens of workers are all pulling from the same dataset shards, benefits from the shared model: one cache, many readers, no duplication. Inference, or single-node training where the working set is stable and predictable, usually doesn't need that coordination and does fine with a local cache.

The obvious objection to shared NVMe-oF is that it adds networking complexity and requires RoCE-capable NICs and switches that not every shop has lying around. For multi-GPU clusters, the coordination it buys is worth the setup cost. For smaller deployments, local NVMe paired with asynchronous checkpoint writes back to object storage is simpler and gets the job done without the extra hardware.

Placement also carries a compliance dimension. Cache nodes should sit in the same region as the compute using them, so data doesn't quietly cross a border it wasn't supposed to cross. In regulated industries, that can decide whether an audit goes fine or doesn't.

Eviction policy and working-set sizing as the lever for cache efficiency

Putting NVMe in the right place solves half the problem. Deciding what the cache keeps and what it throws away matters just as much, because a cache full of the wrong data delivers little benefit.

Most caching systems default to LRU, least recently used, which works well when requests come in unpredictably, the way web traffic does. Training workloads are epoch-driven instead: the full dataset gets read repeatedly in shuffled order, so a shard that looks "least recently used" right now might be scheduled for another read in twenty minutes. LRU evicts it anyway, and the system pays for a shard it just threw away to pull it right back from object storage.

Frequency-aware eviction, LFU, or hybrid policies built around epoch boundaries handle this better. They recognize that a shard getting read once per epoch across many epochs is worth keeping, even if it hasn't been touched in the last few minutes.

Policy only matters if the cache is sized to match the job. The goal is to keep the active epoch's shards, the current minibatches and whatever augmentation cache the job needs, on NVMe, while the next epoch's shards get staged on a warmer tier ahead of time. Get that staging wrong and the cache either sits half-empty or churns constantly, pulling in data it's about to evict.

Small-file datasets break this picture in a different way. When a job is opening millions of tiny files, the real bottleneck often isn't the data itself but the metadata lookup, the inode or index information needed to find each file before reading it. A cache that only holds data pages and ignores that lookup path hasn't actually solved the slowdown. It's cached the wrong thing.

Write-path behavior: why checkpoint traffic needs its own lane

Everything so far has been about reads. Writes, specifically checkpoint writes, cause a separate and commonly overlooked problem, and the usual mistake is routing them through the exact same NVMe path the training reads depend on.

Checkpoint files for large models are big, and writing one mid-training competes directly with the prefetch pipeline trying to stage the next batch. The GPU ends up stalling because the write it's sharing a queue with is hogging the path, not because the data it needs isn't ready.

The fix is to give checkpoint traffic its own lane: a distinct NVMe tier or a dedicated write buffer, so reads and writes aren't fighting over the same queue depth. A well-built system also replicates the write before confirming it's done, then flushes it asynchronously to object storage in the background. Training throughput stays high because the GPU doesn't wait on the slow part, and durability holds because the object store, the actual source of truth, eventually gets a copy. A node failure in this model loses only whatever hadn't been flushed yet, bounded by the flush interval rather than by the whole cache.

Lifecycle policies extend the same logic further out: checkpoints move from NVMe to warm SSD to object storage to cold archive as they age, so only the most recent checkpoints, the ones actually needed for a quick rollback, sit on the expensive tier.

The Dell Technologies and NVIDIA NIXL integration shows what aggressive write-path optimization looks like in practice. It let vLLM and LMCache offload KV cache to Dell ObjectScale over RDMA, writing directly into GPU memory through cuObject. That replaced a substantial chunk of GPU time spent rebuilding context with a much shorter storage retrieval instead. The accelerator spends its time running inference rather than reconstructing context it technically already had access to.

GPU-direct object access and the cache boundary

RDMA-based GPU-direct access to object storage is redrawing where the cache boundary sits, by letting object payloads stream straight into GPU memory and skip the CPU, and in some cases skip the local NVMe staging layer too.

Object storage access used to mean a slow HTTP or REST round trip: data arrived as a byte stream into CPU memory first, then got copied over to GPU memory, with local NVMe commonly sitting in between as scratch space to smooth the handoff. NVIDIA cuObject changes the mechanics of that path. It uses RDMA operations instead of TCP transfer, which lets data move from S3-compatible object storage directly into GPU memory without CPU kernel involvement for the payload. The CPU simply isn't on the critical path for the big transfers anymore.

The xio-sig initiative pushes this further by defining a shared RDMA wire protocol, so GPU-direct object access works across different storage server implementations instead of locking a cluster into one vendor's hardware. Google Cloud, already a maintainer for cuFile within xio-sig, is weighing whether to expand that role to cover cuObject as well, and Microsoft has signaled it intends to join the xio-sig Board. Those moves matter because a shared protocol is only as useful as how many storage vendors actually build to it.

None of this retires NVMe caching. A dataset that gets read across many epochs still benefits enormously from local caching, because without it every epoch pays the object-store round trip all over again. GPU-direct access earns its keep in different territory: single-pass or low-repeat-read work, inference, and KV-cache rehydration, where there's no second epoch coming along to justify staging data on NVMe first.

Google Cloud Storage Rapid is a parallel move in the same direction, a purpose-built object storage tier aimed at cutting blocked GPU time and speeding up data loading for multi-modal training, with faster checkpoint restores built in. Rapid Cache delivers those gains without requiring any code changes. Rapid Bucket asks more of the user, requiring adoption of a new zonal bucket and API in exchange for its benefits.

Requirements for a cloud-agnostic caching layer across S3, GCS, R2, and Azure Blob

A caching layer tied to one object store just swaps old lock-in for new lock-in. Genuine cloud-agnosticism requires the cache to present a stable, consistent interface to compute, while the object store underneath stays swappable.

S3 compatibility has become the default API shape across providers, which helps, but matching the API isn't the whole job. A caching layer still has to account for differences in consistency models, request-rate limits, and egress pricing between S3, GCS, R2, and Azure Blob, because those differences show up as bugs or surprise bills if the cache doesn't handle them.

The interface that matters most for compute is POSIX, a plain filesystem, not the object API. Models and the tooling built around them are trained heavily on bash and ordinary file manipulation, and agents interacting with data overwhelmingly reach for filesystem semantics over custom SDKs. A caching layer aiming to avoid forcing a migration needs to support POSIX fully, not partially: atomic rename, file locking, mmap, fsync, hard links, symlinks, sparse files. Leaving any of that out means existing code needs rewriting to work at all, which quietly cancels the "drop this in, no migration needed" pitch.

The object store has to stay the source of truth throughout. A caching layer built on these principles holds no persistent copy of data outside the customer's own bucket, so access can be revoked at any time and the data never actually left where the customer put it. That's a different architecture from systems that pull data into a proprietary store of their own.

Storage elasticity rounds out the requirements: metering based on what's actively cached rather than the size of the whole dataset. AI workloads often touch a working set that's a small fraction of the full corpus, and which fraction that is shifts from run to run. Paying for the active slice, not the entire archive, matches how the workload actually behaves.

How to decide which architecture fits a given workload

No single NVMe caching setup works for every job. The right one follows from how the workload reads data, how much cache loss it can tolerate, and how it handles writes, and getting that match right decides whether the cache holds up once real production traffic hits it.

Large-scale, distributed, multi-epoch training wants current epoch shards and optimizer state on NVMe, local or over a fabric, with the broader dataset staged on a warm tier ahead of when it's needed. Checkpoint writes need their own separate path, paired with an asynchronous flush back to object storage for fault tolerance.

Online, latency-sensitive inference wants model weights, embeddings, and hot KV caches sitting on NVMe or in memory directly. Object storage retrieval latency is too slow for that hot path to tolerate, and GPU-direct access over RDMA earns its place here specifically for rehydrating KV caches quickly.

Batch or offline inference can get away with warmer object-store retrieval paired with scheduled prefetch into NVMe ahead of the run. Running a full NVMe fleet for that kind of job is overkill when the latency targets are loose to begin with.

Small-file datasets, like a vision job touching millions of individual images, need the cache to cover metadata and lookup paths, not just data pages, or the architecture ends up fast at reading files it's slow to find.

More in POSIX Filesystems