A training cluster is bought by the GPU hour and billed whether or not the GPUs are computing. The GPU data loading bottleneck is the case where they are not: the accelerator finishes a step and waits for the next batch to arrive from the input pipeline. This post is about where that wait comes from, why object storage makes it worse, and how to measure it without fooling yourself.
What a GPU data loading bottleneck looks like
A training step has three parts: produce a batch (read, decode, augment, collate), move it to the device, and run the forward and backward passes. Frameworks overlap the first with the third, so the loader works on the next batch while the GPU computes the current one. That overlap hides loading time only while the loader produces batches at least as fast as the GPU consumes them.
When it does not, the GPU idles for the difference on every step. A loader-bound job does not show up as an error anywhere; it shows up as a run that takes longer than the arithmetic says it should, and as a utilisation graph that looks like a comb.
What dataloader workers actually do
In PyTorch's DataLoader, and in its equivalents elsewhere, a worker is a process that takes a list of sample indices, calls the dataset for each one, and hands the collated batch back over a queue. Each call typically opens a file, reads and decodes it, applies augmentations and converts the result to a tensor.
Three things limit what workers can deliver. Decode and augmentation are CPU work, and a node has a fixed ratio of cores to GPUs, so adding workers past the core count only adds contention. Each worker's reads are mostly serial, so a worker waiting on storage cannot be rescued by more CPU. And the prefetch queue is finite: a stall long enough to drain it reaches the GPU no matter how many workers exist.
The usual first fix is more workers, and it works exactly until one of those three limits is the real constraint. Finding out which one requires measurement, not a larger number in a config file.
Random access, throughput and IOPS
Training reads are random by design: shuffling every epoch means consecutive samples come from unrelated places in the dataset, so the pipeline issues many small reads at scattered offsets rather than one sequential stream. The storage figure that predicts loader performance is therefore operations per second at the sample size, not megabytes per second on a sequential benchmark.
Media differ most on exactly this axis. Memory and SSD handle random reads well; a disk array with high sequential throughput can collapse under random reads of small files, because each one costs a seek. The same dataset can keep a GPU busy from SSD and starve it from HDD, with identical throughput on paper.
Whether the working set fits in local memory decides whether the page cache saves you. If it fits, the second epoch reads from memory and the problem disappears after the first; if not, every epoch pays the full cost, and the first epoch's numbers are the ones to plan from.
Where object storage latency bites
When samples live in a bucket, every open becomes a network request, and the per-request latency of the object store is paid once per sample. The loader now has to hide a delay that is orders of magnitude larger than a local read, and the only tool it has for that is parallelism across workers.
Three patterns make it worse. Small samples mean the request overhead dominates the transfer time, so the store's request rate rather than its bandwidth becomes the ceiling. Listing a large prefix to build the index of samples is itself a long paged operation that runs before the first step. And a cold first epoch pulls the whole dataset across the network while the GPUs wait, then does it again on the next node, because nothing shared remembers what was fetched.
The remedies are structural: pack small samples into shards, put a cache with memory and SSD tiers between the loader and the bucket, and give every node a shared view so that a sample fetched once is served locally after that. The introduction to Runix FS describes that layer; the data pipeline guide covers the data before it reaches storage.
How to measure utilisation honestly
The number most people look at, the utilisation column in nvidia-smi, is a coarse signal: it reports whether a kernel was running during the sample period, not how much of the device that kernel used. A GPU can show high utilisation while running small kernels that leave most of the chip idle, and gaps shorter than the sampling interval disappear.
Better signals, in the order worth collecting:
- Step time with synthetic input. Feed the model the same batch from device memory in a loop and record the step time. That is the floor; the difference between it and the real step time is the loader's share.
- A profiler timeline. The PyTorch profiler and similar tools show data-loading spans next to the kernels. Gaps between kernels that line up with loader activity are the bottleneck, seen rather than inferred.
- Device-level counters. Tools such as NVIDIA DCGM expose metrics like SM activity and tensor-core activity, which separate a kernel being present from the device being busy.
- Loader-side counters. Worker CPU time, time spent waiting on reads, queue depth over time, and the request latency distribution from the storage layer. If the queue is empty when the GPU asks, the loader lost.
Record these for the first epoch and a later one separately, and record time to first batch, where listing and a cold cache show up before any step runs. A pipeline that is fine on the third epoch and starved on the first has a cold-cache problem, which is a different fix from a pipeline starved on every epoch.
What to change, in order of cost
- Match worker count to the CPU cores the job actually has, and measure decode cost per sample first.
- Move decode onto the GPU where a library supports the format.
- Pack small samples into shards, so reads are large and sequential within a shard.
- Put a cache in front of the bucket, so the working set is served from memory or SSD after the first read and shared across nodes.
- Only then buy more storage bandwidth, the most expensive option and the one least likely to be the actual constraint.
Runix FS is designed for the fourth step: a file system and a cache with memory, SSD and HDD tiers in front of the S3-compatible bucket a team already has, so that loaders read through a local path and a block fetched once is served from the cache afterwards. It is in early access, sized and deployed per customer in their own cloud account; we publish no throughput figures of our own, and the quick start shows what a mount involves.
Questions this raises
How do I know whether my training job is bottlenecked on data loading?
Run the model on a synthetic batch held in device memory and compare that step time with the real one. If the real step is much slower and the loader's prefetch queue is empty when the GPU asks for a batch, the input pipeline is the limit.
Does adding more DataLoader workers always help?
No. More workers help until the CPU cores, serial waits on storage or the depth of the prefetch queue become the limit, and past the core count they add contention rather than throughput.
Why is object storage slower for training than local disk?
Each sample read becomes a network request with its own latency and a per-request charge, and training reads are small and random, so the store's request rate rather than its bandwidth sets the ceiling.