All essays
GuideGUIDEFEB 2026

GPU Data Pipeline Optimization: Data Loading, Preprocessing, and Caching for Training

Optimize data pipelines for GPU training. PyTorch DataLoader tuning, NVIDIA DALI, WebDataset, LMDB, data augmentation on GPU, prefetching, and cache strategies for H100/B200 clusters at 50+ GB/s throughput.

01

WHY THE DATA PIPELINE IS THE #1 TRAINING BOTTLENECK

Data loading is the most common performance bottleneck in GPU training, accounting for 20-50 percent of training wall-clock time in production. An H100 GPU can consume 3.3 TB/s of data through its internal memory, but even a 128-GPU cluster reading from a parallel filesystem at 120 GB/s aggregate provides only 0.9 GB/s per GPU. The gap between memory bandwidth and storage throughput is 3,600x at the single-GPU level. The practical constraint is different: each GPU processes approximately 1-5 GB/s of preprocessed data during training (batch size 32, 256 KB sample size, 120 samples per GPU-second for ResNet-152; 500 KB per token sequence, 4K token context, 0.5 GB/s for a 70B LLM).

The bottleneck chain has four tiers: filesystem latency (5-20ms for first byte on parallel FS), decompression (CPU-limited, 2-4 GB/s per core for JPEG or zstd), augmentation (CPU-bound random crop/flip/color jitter operations), and host-to-GPU transfer (PCIe Gen5 at 64 GB/s per H100, i.e., not usually the bottleneck). Each tier must be measured independently: `strace -c` shows filesystem read syscall latency, `perf stat` shows CPU cycles spent in decompression libraries, and PyTorch's `DataLoader` profiling shows `data_load_time` vs `compute_time` in the training loop. The target is `data_load_time / compute_time < 0.1`, meaning the GPU spends less than 10 percent of time waiting for data.

Bottleneck TierTypical Throughput LimitMeasurement MethodBottleneck Indicator
Filesystem I/O (first byte)0.5-5 GB/s per node (NFS)strace -c; iostat -x 1iowait > 5% on CPU, read_avg > 10ms
Decompression (zstd/JPEG)2-4 GB/s per CPU coreperf top (libjpeg-turbo, zstd)CPU core near 100% under compressor function
Image Augmentation (random ops)0.5-2 GB/s per CPU corePyTorch profiler (queue latency)queue_latency > 10ms, GPU idle
Host-to-GPU transfer64 GB/s (PCIe Gen5 x16)nvidia-smi pcie bandwidthGPU memcpy time > 2% of step
GPU compute (Preprocess)DALI: 10-40 GB/s per GPUDALI profiler outputDALI throughput < GPU consume rate
02

PYTORCH DATALOADER TUNING: WORKERS, SHUFFLING, AND PIN MEMORY

The PyTorch DataLoader provides the most accessible optimization lever: `num_workers`. Increasing workers from 0 (single-process) to 4-8 reduces data loading time by 3-5x on systems with sufficient CPU cores and memory bandwidth. The optimal worker count is `2 * num_physical_cores - 2` for training nodes where the GPU training loop uses most CPU cores (PyTorch eager mode with batch normalization). For H100 DGX nodes (128 CPU cores, 2 TB RAM), 8-16 workers per GPU is typical. The `prefetch_factor=3` parameter enqueues 3 batches per worker ahead of the training step, providing a buffer that absorbs I/O variability.

Three additional PyTorch DataLoader parameters that matter for GPU training. `pin_memory=True` copies tensors to page-locked (pinned) memory before transfer to GPU, enabling asynchronous GPU transfers and reducing host-to-GPU copy time from 200 microseconds to 50 microseconds per batch. `persistent_workers=True` keeps worker processes alive between epochs, saving 2-5 seconds of fork overhead per epoch for small datasets. For shuffle, `shuffle=True` with `worker_init_fn` seeded by `torch.initial_seed()` ensures each worker shuffles independently without overlapping indices. The `generator=torch.Generator().manual_seed(seed)` pattern reproduces shuffle order across runs for debugging determinism. The downside: `num_workers > 0` increases host memory usage by approximately 2 GB per worker on datasets with large samples (4K images or sequences).

DataLoader ParameterOptimal SettingImpactTradeoff
num_workers8-16 per GPU3-5x throughput increase vs single processRAM: +2 GB/worker, CPU: +8-16 cores
prefetch_factor2-330-50% reduction in data wait timeRAM: batches stay in CPU memory longer
pin_memoryTrue50-75% reduction in host->GPU copy timeRAM: pinned memory is non-swappable, allocates ~batch size
persistent_workersTrue (epochs > 1)Eliminates 2-5s worker startup per epochRAM: workers stay live, ~500 MB/worker resident
batch_samplerDistributedSampler with drop_last=TrueExactly equal batch distribution across GPUsLoses last partial batch per GPU
timeout60-120 secondsPrevents training hang on slow I/ODelays error detection by 60s per failed batch
03

NVIDIA DALI: GPU-ACCELERATED DATA PIPELINES

NVIDIA DALI (Data Loading Library) moves data preprocessing from CPU to GPU, eliminating the CPU bottleneck that the PyTorch DataLoader cannot address. DALI defines preprocessing pipelines as directed acyclic graphs executed on the GPU: decode JPEG directly from GPU memory (nvJPEG, 30-50 GB/s throughput on H100), apply random crop/flip/color jitter as GPU CUDA kernels (instead of CPU OpenCV), and concatenate samples into GPU-resident batches that are fed directly to the model without CPU round-trip. DALI supports four levels of GPU offload: Level 0 (CPU pipeline, similar to PyTorch DataLoader), Level 1 (CPU decode + GPU augmentation), Level 2 (GPU decode + GPU augmentation), and Level 3 (GPU filesystem I/O via GDS).

Benchmarks from NVIDIA show DALI Level 2 on a single H100 achieving 32,700 images/second throughput for ResNet-50 training (224x224, random crop/flip, JPEG decode) versus 8,900 images/second with PyTorch DataLoader (8 workers). The 3.7x throughput improvement comes from moving JPEG decode and augmentation from CPU (12 CPU cores saturated) to GPU (25 percent H100 utilization). For LLM training on sequence data (tokenized text), the DALI advantage is smaller (1.2-1.5x) because text preprocessing is CPU-light (no JPEG decode, no complex augmentation). The DALI pipeline is defined in Python: `pipe = nvidia.dali.pipeline.Pipeline(batch_size=4096, num_threads=4, device_id=0)` with operators like `fn.readers.file()`, `fn.decoders.image()`, `fn.resize()`, `fn.crop_mirror_normalize()`, and `fn.cast()`.

04

WEBDATASET AND SHARDED DATA FORMATS

WebDataset is the standard sharded data format for GPU training, solving the file-per-sample problem that causes metadata server overload on parallel filesystems. Instead of storing 1.4 million individual JPEG files for ImageNet training (creating 1.4M inode operations per epoch), WebDataset packs samples into `.tar` archives (shards) of 256 MB-1 GB each. Each shard contains 10,000-50,000 samples. The GPU cluster's 256-512 training workers each read one shard file sequentially, reducing filesystem metadata operations from 1.4M per epoch to 64 per epoch (one open per shard). WebDataset's `wds.WebDataset` integrates with PyTorch: `dataset = wds.WebDataset("datasets/imagenet-{000000..000064}.tar").shuffle(10000).decode("torch").to_tuple("jpg", "cls")`.

Shard size optimization balances per-worker throughput against load balancing. `shard_size = total_dataset_bytes / (num_gpus * prefetch_batches * batch_size * sample_bytes)`. For a 1.4 TB training dataset on 512 GPUs, 64 shards of 22 GB each give each GPU 2-4 shards. The dataset is replicated across parallel filesystem OSTs or WEKA nodes for broadcast read during epoch start. MosaicML's StreamingDataset (now part of Databricks) extends WebDataset with deterministic shard ordering, shard metadata caching, and local SSD caching for cloud-based training. `local=/mnt/cache` flag in StreamingDataset automatically downloads and caches shards to local NVMe on first epoch, reducing epoch 2+ read latency from 20ms filesystem first-byte to 100 microsecond local NVMe.

Data FormatOpen/Read PatternFiles per 1M SamplesEpoch 1 Time (1K GPU)Epoch 2+ TimeUse Case
Individual files (JPEG)1 open per sample (1.4M opens)1,400,00055 min45 minLegacy, not recommended for GPU training
WebDataset tar (256 MB)1 open per shard (128 opens)128 (shards)32 min28 minImage training, multimodal datasets
Mosaic StreamingDataset1 open per shard, local cache128 (shards)35 min12 min (local NVMe)Cloud training, multi-epoch runs
LMDB (mmap)1 open per database, mmap1 (database)28 min26 minSmall- to medium-sized datasets (< 100 GB)
TFRecord1 open per shard via tf.data128 (shards)38 min32 minTensorFlow native, TPU training
05

TIERED CACHING: LOCAL NVME, HOST RAM, AND GPU MEMORY

Tiered caching reduces data pipeline latency by storing frequently accessed samples in progressively faster but smaller storage tiers. The three-tier GPU cache hierarchy: Tier 1 (GPU VRAM, 80-144 GB per H100, ~3.3 TB/s bandwidth, 10-100 nanosecond latency) stores the current and prefetched batches; Tier 2 (Host RAM, 512 GB-2 TB per DGX node, ~50 GB/s bandwidth, 100-200 nanosecond latency) stores cached samples for the current epoch; Tier 3 (Local NVMe SSD, 2-15 TB per node, ~7 GB/s read, ~30 microsecond latency) stores shard files for the entire dataset. NVIDIA's GPUDirect Storage (Tier 0) bypasses host RAM entirely for direct GPU loading from NVMe or fabric storage.

The most effective caching strategy for multi-epoch GPU training is local NVMe cache on first epoch read, then host RAM page cache for subsequent epochs. Implementation: PyTorch's `torchdata` with `OnlineCache` reads samples from the parallel filesystem on first pass, writes to local NVMe (`/mnt/cache/dataset-{uuid}/`), and reads from local NVMe on subsequent epochs. The cache hit rate for epoch 2+ is 100 percent with local NVMe. The `torchdata` cache state machine: `cache_policy="warm"` (write-through on first epoch), `cache_policy="read"` (read-only from cache on epoch 2+). For RAM caching, Linux's page cache automatically caches files read from NVMe: `echo 3 > /proc/sys/vm/drop_caches && vmtouch -t /mnt/cache/imagenet-shard-*.tar` preloads all shards into page cache (if sufficient RAM) for zero-wait epoch transitions.

Filed under
PyTorch DataLoaderNVIDIA DALI OptimizationGPU Data PreprocessingWebDataset TrainingData Cache StrategiesGPU Data AugmentationTraining I/O Optimization