A training job is a pipeline whose last stage happens to be very expensive. Storage hands bytes to CPU workers, which decode, augment and collate them; the batch crosses PCIe into GPU memory; only then do kernels run. If any earlier stage is slower than the GPU, the GPU waits, and the bill shows it as GPU-hours that bought nothing.
This article treats the input pipeline as a throughput budget you can measure and fix: detecting a data stall, DataLoader workers, pinned memory and a side CUDA stream for the copy, GPU decode with NVIDIA DALI, and storage format. A worked example sizes an image job, and a checklist closes it.
The pipeline as a throughput budget
Every stage has a rate, measured in samples per second, and the pipeline runs at the rate of its slowest stage. The GPU stage is fixed by the model and batch size: a step that takes 120 ms on a batch of 256 consumes about 2,133 samples per second. Everything upstream must sustain that, with headroom.
Each stage is limited by a different resource, so the fixes differ. Reading is bounded by storage throughput and request rate, since every small file costs a metadata lookup or an object-store request. Decode and augmentation are bounded by CPU cores. The host-to-device copy is bounded by PCIe bandwidth and by whether the source memory is page-locked. Queues absorb jitter but cannot raise the average rate of a stage that is too slow.
Two consequences follow. More DataLoader workers help only when the CPU stage is the bottleneck and spare cores exist. And optimising the model makes stalls worse: teams often discover their input pipeline after enabling mixed precision, when the step gets faster and throughput barely moves.
Measuring a data stall
Before changing anything, measure how long each iteration waits for data versus how long the GPU computes. GPU work is asynchronous, so step time uses CUDA events with a synchronise only at logging intervals.
import time, torch
def train_epoch(loader, model, opt, log_every=50):
it = iter(loader)
wait_s, step_events = 0.0, []
for i in range(len(loader)):
t0 = time.perf_counter()
x, y = next(it) # blocks if no batch is ready
x = x.cuda(non_blocking=True)
y = y.cuda(non_blocking=True)
wait_s += time.perf_counter() - t0
start = torch.cuda.Event(enable_timing=True)
end = torch.cuda.Event(enable_timing=True)
start.record()
loss = torch.nn.functional.cross_entropy(model(x), y)
opt.zero_grad(set_to_none=True)
loss.backward()
opt.step()
end.record()
step_events.append((start, end))
if (i + 1) % log_every == 0:
torch.cuda.synchronize()
gpu_ms = sum(s.elapsed_time(e) for s, e in step_events) / len(step_events)
wait_ms = 1000 * wait_s / len(step_events)
print(f"step {i+1}: gpu {gpu_ms:.1f} ms, data wait {wait_ms:.1f} ms")
wait_s, step_events = 0.0, []If data wait is a few percent of GPU time, the pipeline is fine. If it is large, the job is input-bound. To find the stage, benchmark the loader alone in samples per second, then with augmentation disabled, then reading pre-decoded tensors; the removal that helps most names the bottleneck. Cores pinned at 100% point to decode, high I/O wait to storage.
CPU workers in the PyTorch DataLoader
The PyTorch DataLoader parallelises the CPU stages with worker processes. Each worker calls the dataset's __getitem__ for its indices, collates a batch and puts it on a shared queue; processes sidestep the Python global interpreter lock.
from torch.utils.data import DataLoader
loader = DataLoader(
dataset,
batch_size=256,
shuffle=True, # or a DistributedSampler per rank
num_workers=12, # CPU decode/augment parallelism per GPU process
pin_memory=True, # batches land in page-locked host memory
prefetch_factor=4, # batches queued per worker (only with workers > 0)
persistent_workers=True, # keep workers alive between epochs
drop_last=True, # fixed shapes, no ragged final batch
)Start num_workers from the cores per GPU process, minus a couple for the main process, and measure; more workers than cores adds contention, not throughput. Each worker holds its own copy of the dataset object's state.
prefetch_factor is the batches each worker loads ahead; deeper queues absorb bursty reads at the cost of host memory. persistent_workers avoids re-forking workers at every epoch boundary. pin_memory copies finished batches into page-locked memory, explained next.
Keep per-sample work lean: crop before resizing, and normalise on the GPU. Sending uint8 instead of float32 cuts the bytes crossing PCIe by four.
The host-to-device copy: pinned memory and streams
The GPU's DMA engines need memory that cannot move. From pageable memory, the driver first copies into a pinned staging buffer, adding a CPU memcpy and making the transfer synchronous with the host. From pinned memory, the DMA reads directly and the copy can run asynchronously.
In PyTorch that means pin_memory=True on the loader together with .to(device, non_blocking=True). To overlap the copy of batch n+1 with compute on batch n, copy on a side stream:
MEAN = torch.tensor([0.485, 0.456, 0.406], device="cuda").view(1, 3, 1, 1)
STD = torch.tensor([0.229, 0.224, 0.225], device="cuda").view(1, 3, 1, 1)
class CudaPrefetcher:
"""Copies the next batch on a side stream while the current one computes."""
def __init__(self, loader, device="cuda"):
self.loader, self.device = loader, device
self.stream = torch.cuda.Stream()
def __iter__(self):
it = iter(self.loader)
nxt = self._load(it)
while nxt is not None:
torch.cuda.current_stream().wait_stream(self.stream)
x, y = nxt
x.record_stream(torch.cuda.current_stream()) # memory safety across streams
y.record_stream(torch.cuda.current_stream())
nxt = self._load(it) # start the next copy now
yield x, y
def _load(self, it):
try:
x, y = next(it)
except StopIteration:
return None
with torch.cuda.stream(self.stream):
x = x.to(self.device, non_blocking=True)
y = y.to(self.device, non_blocking=True)
x = x.float().div_(255).sub_(MEAN).div_(STD) # normalise on the GPU
return x, ywait_stream stops the model reading a half-copied tensor; record_stream stops the caching allocator reusing the memory while compute-stream kernels still read it. Pinned memory is finite, so size prefetch queues from measurement.
Moving decode and augmentation onto the GPU
When the node lacks the cores, move decode and augmentation onto the GPU. NVIDIA DALI builds a pipeline graph whose operators run on CPU or GPU; its mixed decoder parses on the CPU and decodes on the GPU. This pipeline returns cropped, flipped, normalised tensors already in GPU memory.
from nvidia.dali import pipeline_def, fn, types
from nvidia.dali.plugin.pytorch import DALIGenericIterator
@pipeline_def
def train_pipe(file_root, shard_id, num_shards):
jpegs, labels = fn.readers.file(
file_root=file_root, random_shuffle=True,
shard_id=shard_id, num_shards=num_shards, name="Reader")
images = fn.decoders.image(jpegs, device="mixed", output_type=types.RGB)
images = fn.random_resized_crop(images, size=224)
images = fn.crop_mirror_normalize(
images, dtype=types.FLOAT, output_layout="CHW",
mean=[0.485 * 255, 0.456 * 255, 0.406 * 255],
std=[0.229 * 255, 0.224 * 255, 0.225 * 255],
mirror=fn.random.coin_flip())
return images, labels.gpu()
pipe = train_pipe(file_root="/data/train", shard_id=rank, num_shards=world,
batch_size=256, num_threads=4, device_id=local_rank)
pipe.build()
loader = DALIGenericIterator(pipe, ["data", "label"], reader_name="Reader")
for batch in loader:
x, y = batch[0]["data"], batch[0]["label"]The trade is GPU time and memory for CPU time, which matters on a model already at its memory limit. The gain is largest with few cores per GPU, large images or heavy augmentation. DALI brings its own sharding and iterator semantics, so check partial-batch and epoch handling, and judge it by lower data wait, not lower CPU usage.
Storage format, shuffling and caching
Storage format often decides throughput more than worker count. Millions of small files on a network file system or object store are limited by request rate, not bandwidth. Packing samples into large shards, such as tar files of a few hundred megabytes read sequentially, turns millions of random reads into thousands of streaming ones; WebDataset and TFRecord follow this idea.
Shards make shuffling two-level: shuffle shard order per epoch, and samples within a bounded buffer. Small buffers over sorted shards leave batches correlated, so shuffle once before sharding. Assign shards to ranks deterministically so each sample is seen once per epoch and restarts can resume.
Cache on local NVMe so only the first epoch pays the remote read. GPUDirect Storage moves data from NVMe straight into GPU memory, which helps with GPU decode but not with a CPU decode bottleneck.
Worked example: feeding an 8-GPU image job
Take an image classifier on an 8-GPU node with 64 cores. Profiling gives these assumptions: one GPU completes a step on a batch of 256 in 85 ms, so each GPU consumes about 3,000 images per second and the node about 24,000. Average JPEG size is 110 KB. Decode plus random crop, resize and flip costs 6 CPU-ms per image on one core, measured by timing the dataset's item method in a loop.
| Budget | Arithmetic | Requirement |
|---|---|---|
| CPU decode + augment | 24,000 images/s x 6 ms | 144 core-seconds per second: 144 cores |
| Cores available | 64 minus about 8 for training processes | about 56 cores: CPU-bound by more than 2x |
| Storage read | 24,000 x 110 KB | about 2.6 GB/s sustained, more on the first epoch |
| H2D copy as uint8 | 24,000 x 3 x 224 x 224 bytes | about 3.6 GB/s for the node, about 450 MB/s per GPU |
| H2D copy as float32 | 4 x the uint8 figure | about 14.5 GB/s for the node |
The CPU stage needs about 144 cores and the node has about 56, so no DataLoader tuning reaches the target: the node tops out near 56/144 of it, about 9,300 images per second, and the GPUs wait more than half the time. The fixes are pre-decoding to the training resolution offline, GPU decode with DALI, or a node shape with more cores per GPU.
The copy budget shows why normalisation belongs on the GPU: float32 would quadruple PCIe traffic for nothing. The storage figure says shard and cache locally. After the fix, rerun the stall timer and aim for data wait under about 5% of step time.
Failure modes
These failures recur:
- Worker memory blow-up. Python lists of millions of objects defeat copy-on-write through reference counting, so memory grows per worker. Use NumPy arrays or memory-mapped files.
- Duplicate random augmentation. Workers forked with the same random state produce identical crops. PyTorch seeds its own generator per worker, but an augmentation library with its own generator may need seeding in a worker init function, derived from the per-worker torch seed.
- Epoch-boundary stall. Without
persistent_workers, each epoch restarts workers and refills queues. - Oversubscribed CPU. Workers plus NumPy or OpenCV threads inside each fight training for cores. Limit library threads to one per worker.
- Uneven ranks in distributed training. One slow rank's loader delays every collective. Check per-rank data wait.
- Silent sample loss. Dropped partial shards or restart-from-zero resumes change what the model sees. Log samples per epoch.
- Pinned-memory exhaustion. Huge prefetch queues of pinned buffers starve the host of locked pages, and allocation becomes slow or fails.
Trade-offs
| Choice | Gains | Costs |
|---|---|---|
| More DataLoader workers | CPU-stage throughput, if cores are free | Memory per worker; core contention with training |
| Pre-decode / pre-resize offline | Lowest per-epoch CPU cost | Storage size; fixed resolution limits augmentation |
| GPU decode (DALI) | Removes the CPU decode bottleneck | GPU time and memory; a second data API to maintain |
| Sharded formats | Sequential reads, high request efficiency | Two-level shuffling; resharding when data changes |
| Local NVMe cache | Later epochs avoid remote storage | Disk capacity; cache warm-up time |
| Deep prefetch queues | Absorbs jitter | Host and pinned memory; hides, does not fix, a slow stage |
The order of operations that usually wins: measure, fix the storage format, move normalisation to the GPU, tune workers, and only then consider GPU decode. For how the copy interacts with host topology, see PCIe host interconnect architecture. For what the GPU does with the batch once it arrives, see anatomy of one GPU training step, and for the memory the batch competes for, the GPU memory hierarchy.
What to do next
- Add the data-wait versus GPU-time timer to your training loop and record both numbers for one hundred steps.
- Benchmark the loader alone in samples per second, then with augmentation off, to identify the bottleneck stage.
- Set
pin_memory=Trueandnon_blocking=Truetogether, and move normalisation to the GPU so uint8 crosses PCIe. - Size
num_workersfrom available cores per GPU, cap library threads at one per worker, and enable persistent workers. - Write your budget table: samples per second needed, CPU-ms per sample, bytes per sample read and copied.
- If the CPU budget exceeds the cores you have, pre-decode offline or prototype a DALI pipeline and compare data wait.
- Convert small-file datasets to shards, shuffle before sharding, and cache on local NVMe.
- Check per-rank data wait in distributed jobs, and log samples seen per epoch.