A training job is input-bound when the GPU finishes a step and then waits for the next batch. Utilisation dashboards show it as a sawtooth: bursts of 100% separated by gaps, or a flat 40-60% that no kernel tuning moves. The cause is rarely the GPU. It is a pipeline of Python processes, queues, shared memory and a copy engine that has to produce one batch per step, every step, faster than the model consumes it.
This article opens up the PyTorch DataLoader itself: which process runs which stage, how many batches are really in flight, why one slow sample can stall a whole job, and how to prove with three short measurements whether the loader is your limit before you change anything. The broader pipeline (storage formats, GPU decode, sharding) is covered in GPU data pipelines in depth; here the focus is the loader as a machine and the failure modes that live inside it.
The throughput budget
Start with arithmetic. Call the GPU's compute time per step T_gpu and the time the loader needs to produce one batch T_batch. With N workers running in parallel, the loader delivers a batch every T_batch / N seconds on average. The job is input-bound when T_batch / N > T_gpu. Queues only absorb variance: they let a fast stretch bank batches for a slow one, but they cannot make the average rate higher than the slowest stage.
That is the most common misunderstanding about prefetch_factor. Raising it from 2 to 8 helps when production is fast on average but bursty, such as a few huge images in an otherwise small dataset. It does nothing when every batch is slow, except consume more host memory. The fix for a slow average is less work per sample, more parallel workers, or moving work off the CPU.
Inside the DataLoader
With num_workers=0 everything happens inline in the training process: the loop asks for a batch, the dataset is indexed batch_size times, the samples are collated, and only then does the GPU get work. The GPU is idle for the whole of that time. This is the baseline that every multi-worker setting must beat.
With num_workers=N the main process keeps the sampler and sends lists of indices to the workers' index queues in round-robin order. Each worker is a separate OS process with its own copy of the dataset object. It reads its next index list, calls dataset[i] for each index, runs collate_fn, and places the resulting tensors in shared memory. Only a small handle travels through the result queue; the main process maps the same memory instead of copying bytes through a pipe.
At start-up the loader sends prefetch_factor * num_workers index lists, and after that one new list each time the training loop consumes a batch. So the number of batches in flight is capped at that product, and each one occupies shared memory until it is consumed. With pin_memory=True a thread in the main process copies each finished batch into page-locked host memory, which lets the later .to(device, non_blocking=True) run as an asynchronous DMA transfer instead of a synchronous staged copy.
Finally, the batches have to come out in order. Batch 7 is assigned to a particular worker, and if it is not ready the loader holds batches 8, 9 and 10 in a reorder buffer even if they finished first. Recent PyTorch releases expose in_order=False to relax this, at the cost of a deterministic batch order.
Head-of-line blocking
Head-of-line blocking is the failure mode people miss, because average throughput looks fine. Suppose each batch takes 200 ms to build, except that one sample in 500 is a corrupted file that hits a 10 second read timeout. With 8 workers the average rate is still good, but when the slow batch is next in line, the training loop waits for it even though other workers have finished several batches behind it. Data wait shows a long tail: a p50 near zero and a p99 of seconds.
The signature is distinctive: per-step data wait is usually zero, with occasional spikes far larger than T_batch. The fixes are, in order: find and fix the slow samples (log any __getitem__ slower than a threshold, with its index), bound the slow operation (a read timeout with a fallback sample), and only then consider in_order=False. Out-of-order delivery hides the stall, but it makes runs harder to reproduce and complicates checkpoint resume, because the set of consumed samples is no longer a prefix of the sampler's order.
Proving the loader is the limit
Three measurements separate the cases. Run each for a few hundred steps after warm-up, on the real machine type, because core counts and storage differ between a workstation and a cluster node.
- Loader only. Iterate the DataLoader with no model and no copy. This is the ceiling of the CPU and storage stages, in batches per second.
- Model only. Build one real batch, move it to the GPU once, and run the training step on it repeatedly. This is the ceiling of the GPU, including the optimizer step.
- In-loop data wait. In the real loop, time how long
next(it)blocks, per step, and report p50 and p99 alongside the step time.
If loader-only throughput is below model-only throughput, you are input-bound and no GPU work will help. If both ceilings are comfortably above the real rate, the cost is in the hand-off: the copy, the pin thread, or contention between workers and the training process for the same cores. Here is a compact harness.
import time, statistics, torch
def loader_only(loader, steps=300, warmup=30):
it = iter(loader)
for _ in range(warmup):
next(it)
t0 = time.perf_counter()
for _ in range(steps):
next(it)
return steps / (time.perf_counter() - t0) # batches/s, CPU + storage ceiling
def model_only(model, batch, opt, loss_fn, steps=100, warmup=10):
x, y = (t.cuda(non_blocking=True) for t in batch)
for i in range(warmup + steps):
if i == warmup:
torch.cuda.synchronize(); t0 = time.perf_counter()
opt.zero_grad(set_to_none=True)
loss_fn(model(x), y).backward()
opt.step()
torch.cuda.synchronize()
return steps / (time.perf_counter() - t0) # batches/s, GPU ceiling
def train_with_wait(model, loader, opt, loss_fn, steps=300):
waits, it = [], iter(loader)
for _ in range(steps):
t0 = time.perf_counter()
x, y = next(it) # blocks if no batch is ready
waits.append(time.perf_counter() - t0)
x, y = x.cuda(non_blocking=True), y.cuda(non_blocking=True)
opt.zero_grad(set_to_none=True)
loss_fn(model(x), y).backward()
opt.step()
waits.sort()
return statistics.median(waits), waits[int(0.99 * len(waits))]The step is asynchronous, so a brief wait in next(it) while the GPU still has queued kernels costs nothing; a persistent wait drains the queue and does. For a timeline view of the same thing, see GPU profiling with Nsight, where an input-bound job shows gaps in the CUDA stream that line up with the main thread blocked on a queue.
Worked example: a CPU budget for 8 GPUs
An image classifier trains on one node with 8 GPUs and 64 CPU cores. Each GPU takes a batch of 256 images, and the model-only measurement gives 0.30 s per step, or about 853 images per second per GPU. The job needs 8 x 853 = 6,827 images per second in total.
The loader-only measurement for one worker gives 111 images per second, which means about 9 ms of CPU per image: about 5 ms of JPEG decode, 3 ms of random resized crop and colour jitter, and 1 ms of tensor conversion and collation. The total CPU demand is 6,827 x 9 ms = 61.4 core-seconds per second, so 61 fully busy cores, before counting the 8 training processes, their pin threads and the operating system. The node has 64. The job is input-bound by construction, and measured data wait confirms it: p50 of 40 ms on a 300 ms step, so GPU utilisation is stuck near 85%, with periodic stalls on top.
More workers will not help, because the cores are already saturated. The fixes target work per image. Pre-resizing the dataset offline so that the shorter side is 320 pixels cuts decode time from 5 ms to about 1.5 ms, because decode cost scales with pixel count. Demand falls to about 5.5 ms per image, or 38 cores, leaving room for 5 workers per GPU plus the training processes. Data wait drops to zero at p50 and step time becomes the model's 0.30 s. If the dataset must stay at full resolution, decoding on the GPU is the next option, trading a little GPU time for most of the CPU cost.
Configuring the loader
Once the budget works, configure the loader so it delivers it. The settings below are a reasonable start for a map-style dataset on Linux; the comments say why each one is there.
import os, torch
from torch.utils.data import DataLoader, get_worker_info
def worker_init(worker_id):
# Each worker is a process; stop libraries inside it from spawning thread pools
# that fight the other workers and the training process for the same cores.
torch.set_num_threads(1)
os.environ["OMP_NUM_THREADS"] = "1"
try:
import cv2
cv2.setNumThreads(0)
except ImportError:
pass
# PyTorch seeds torch, random and NumPy's global RNG per worker. A Generator
# created in Dataset.__init__ is copied into every worker with the same state,
# so reseed it here from the per-worker torch seed.
info = get_worker_info()
ds = info.dataset
if hasattr(ds, "reseed"):
ds.reseed(info.seed % 2**32)
loader = DataLoader(
train_set,
batch_size=256,
shuffle=True, # or a DistributedSampler, with set_epoch() each epoch
num_workers=5, # from the core budget, not "as many as possible"
pin_memory=True, # enables async host-to-device copies
persistent_workers=True, # no worker restart and dataset re-init each epoch
prefetch_factor=4, # absorbs jitter; costs shared and pinned memory
drop_last=True, # equal-sized steps across DDP ranks
worker_init_fn=worker_init,
)Iterable datasets need one more step. Every worker gets its own copy of the iterator, so without sharding each sample is produced num_workers times per rank. Split the stream by both the distributed rank and the worker id:
class ShardedStream(torch.utils.data.IterableDataset):
def __init__(self, shard_paths, rank, world_size):
self.paths, self.rank, self.world = shard_paths, rank, world_size
def __iter__(self):
info = get_worker_info()
wid, nw = (info.id, info.num_workers) if info else (0, 1)
# Global reader index across all ranks and workers.
reader, readers = self.rank * nw + wid, self.world * nw
for path in self.paths[reader::readers]:
yield from read_records(path) # your decoder
Failure modes
| Symptom | Cause inside the loader | Fix |
|---|---|---|
| Worker dies with a bus error mentioning shared memory | Batches live in /dev/shm; containers often default to 64 MB | Raise the container's shm size or share the host IPC namespace; lower prefetch_factor |
| Data wait p50 near zero, p99 of seconds | Head-of-line blocking on a slow or corrupt sample | Log slow indices, bound reads with a timeout, fix the files |
| Epoch has N times the expected samples | Unsharded IterableDataset in N workers | Shard by rank and get_worker_info() |
| Augmentations identical across workers | A random generator object created in __init__ is copied to every worker | Reseed it in worker_init_fn from info.seed |
Cannot re-initialize CUDA in forked subprocess | Dataset code touches CUDA in a forked worker | Keep workers CPU-only, or use the spawn start method |
| Loader hangs at start-up or after many epochs | Fork while another thread holds a lock, or a worker stuck in I/O | Set timeout so a hang becomes an error; avoid threads before forking |
| Startup takes minutes and fails on pickling | Spawn or forkserver start method pickles the dataset per worker | Keep dataset state small: paths and offsets, not loaded arrays |
The start method deserves a note. On Linux, Python historically forked workers, which is fast and shares the parent's memory pages. Python 3.14 changed multiprocessing's default start method on Linux away from fork, so code that relied on fork (unpicklable lambdas in a transform, large in-memory arrays captured cheaply) can become slow or break on upgrade. Pass multiprocessing_context explicitly if you depend on one behaviour.
Trade-offs
| Choice | Gains | Costs |
|---|---|---|
| More workers | More parallel CPU work while cores are free | Memory per worker; contention once cores are saturated |
Higher prefetch_factor | Absorbs bursty sample cost | Shared and pinned memory; no gain in the average rate |
in_order=False | Hides slow samples | Non-deterministic order; harder resume and debugging |
persistent_workers=True | No per-epoch restart | Workers hold memory between epochs; dataset changes need a new loader |
| Offline preprocessing | Lowest CPU cost per step | Storage, a rebuild step, less augmentation freedom |
| GPU-side decode and augment | Frees CPU cores | GPU time and memory; a second data path to maintain |
Two related levers sit on the model side. Mixed precision makes steps faster, which makes the loader relatively slower, so re-check data wait after enabling it (mixed precision training). And an input-bound job is paying for idle accelerators by the hour, which is why fixing it often beats buying more GPUs (GPU cost optimisation).
What to do next
- Add per-step data wait (p50 and p99) and step time to your training logs, for every rank.
- Run the loader-only and model-only measurements on the production machine type and write both numbers down.
- Compute CPU demand: per-sample CPU time times required samples per second. Compare it with free cores.
- If demand exceeds cores, cut work per sample (pre-resize, cheaper augmentation, cached tokenisation) before adding workers.
- Set
pin_memory,persistent_workersand a thread-limitingworker_init_fn; sizenum_workersfrom the budget. - Log any sample slower than 10 times the median, with its index, and fix the files it finds.
- For iterable datasets, assert that samples per epoch equal the dataset size across all ranks and workers.
- Check the container's shared memory size and pin the multiprocessing start method explicitly.