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

One DataLoader iteration: who does what, and where a batch can waitSamplermain process: indicesIndex queuesone per workerWorker processes 0..N-1dataset[i] for each indexdecode, augment (CPU)collate_fn -> batchtensors into shared memorycap: prefetch_factorbatches in flight per workerround-robinResult queuehandles, not bytesReorder bufferin_order=TruePin threadpage-locked copyTrain loopnext(it)GPU: H2D copy + stepnon_blocking=True on a pinned batchData wait = time the train loop blocks in next(it). If it is > 0 every step, the GPU starves.One slow sample stalls the reorder buffer for every batch queued behind it (head-of-line blocking).
The DataLoader's moving parts. Workers are separate processes; the train loop sees only what reaches the end of the reorder buffer.

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.

  1. 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.
  2. 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.
  3. 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

SymptomCause inside the loaderFix
Worker dies with a bus error mentioning shared memoryBatches live in /dev/shm; containers often default to 64 MBRaise the container's shm size or share the host IPC namespace; lower prefetch_factor
Data wait p50 near zero, p99 of secondsHead-of-line blocking on a slow or corrupt sampleLog slow indices, bound reads with a timeout, fix the files
Epoch has N times the expected samplesUnsharded IterableDataset in N workersShard by rank and get_worker_info()
Augmentations identical across workersA random generator object created in __init__ is copied to every workerReseed it in worker_init_fn from info.seed
Cannot re-initialize CUDA in forked subprocessDataset code touches CUDA in a forked workerKeep workers CPU-only, or use the spawn start method
Loader hangs at start-up or after many epochsFork while another thread holds a lock, or a worker stuck in I/OSet timeout so a hang becomes an error; avoid threads before forking
Startup takes minutes and fails on picklingSpawn or forkserver start method pickles the dataset per workerKeep 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

ChoiceGainsCosts
More workersMore parallel CPU work while cores are freeMemory per worker; contention once cores are saturated
Higher prefetch_factorAbsorbs bursty sample costShared and pinned memory; no gain in the average rate
in_order=FalseHides slow samplesNon-deterministic order; harder resume and debugging
persistent_workers=TrueNo per-epoch restartWorkers hold memory between epochs; dataset changes need a new loader
Offline preprocessingLowest CPU cost per stepStorage, a rebuild step, less augmentation freedom
GPU-side decode and augmentFrees CPU coresGPU 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

  1. Add per-step data wait (p50 and p99) and step time to your training logs, for every rank.
  2. Run the loader-only and model-only measurements on the production machine type and write both numbers down.
  3. Compute CPU demand: per-sample CPU time times required samples per second. Compare it with free cores.
  4. If demand exceeds cores, cut work per sample (pre-resize, cheaper augmentation, cached tokenisation) before adding workers.
  5. Set pin_memory, persistent_workers and a thread-limiting worker_init_fn; size num_workers from the budget.
  6. Log any sample slower than 10 times the median, with its index, and fix the files it finds.
  7. For iterable datasets, assert that samples per epoch equal the dataset size across all ranks and workers.
  8. Check the container's shared memory size and pin the multiprocessing start method explicitly.
Key takeaway: A DataLoader is a set of worker processes, queues and shared memory that must deliver one batch per step on average; queues absorb bursts but never raise the average rate. Prove an input bottleneck with loader-only, model-only and in-loop data wait measurements, budget CPU cores per sample before adding workers, cut work per sample when cores run out, and watch the p99 of data wait for head-of-line stalls from slow samples.