What is Data Parallelism?
Data parallelism is the simplest way to train on several GPUs: every GPU holds a full copy of the model, each one processes a different slice of the training batch, and the GPUs average their gradients before every weight update so the copies stay identical. It speeds up training by processing more examples per step. It does nothing to help a model that does not fit on one GPU, because each GPU still stores everything.
How it works
One training step, with N GPUs:
- Each GPU gets 1/N of the batch and runs the forward and backward pass on its slice, producing its own gradients.
- The GPUs run an all-reduce on the gradients, so every GPU ends up with the average over all N slices. This is the one communication step, and it is what NCCL is used for.
- Every GPU applies the same optimizer update to its own copy, so all copies remain identical.
The effective batch size is the per-GPU batch times the number of GPUs, so 8 GPUs each processing 4 sequences train on 32 sequences per step. Making the batch larger this way is why data parallelism raises throughput: it does not make a single step faster, it processes more examples per step.
In PyTorch this is DistributedDataParallel. Its design paper (Li et al., 2020) describes two tricks that hide most of the cost: gradients are grouped into buckets, and the all-reduce for a bucket starts as soon as its gradients are ready, overlapping with the rest of the backward pass.
Worked example: Llama 3.1 8B
Llama 3.1 8B has 8.03 billion parameters (computed from its published config.json: 32 layers, hidden size 4096, feed-forward size 14,336, 8 key/value heads, 128,256-token vocabulary). Train it on 8 GPUs in BF16 gradients.
- Gradient size: 8.03 billion x 2 bytes = 16.06 GB.
- Ring all-reduce traffic per GPU: 2 x (8 - 1) / 8 x 16.06 = 28.1 GB sent and the same received. This is the cost per optimizer step, independent of how many examples each GPU processed.
- Time at theoretical peak, using half of a bidirectional link rate per direction as the NCCL page does: NVLink at 900 GB/s gives 450 GB/s each way and about 0.06 s. PCIe Gen 4 at 64 GB/s gives 32 GB/s each way and about 0.88 s. With overlap, part of that hides behind compute, but a step that computes for less than the all-reduce takes is limited by the network.
- Memory per GPU: the ZeRO paper (Rajbhandari et al., 2019) counts 16 bytes per parameter for mixed-precision training with Adam: 2 for the weights, 2 for the gradients and 12 for the fp32 weight copy plus the two Adam moments. For 8.03 billion parameters that is 128.5 GB of model state on every GPU, before activations. That exceeds an 80GB H100, and even fits a 141GB H200 only with almost nothing left over.
The last bullet is data parallelism's real limit. Adding GPUs does not reduce per-GPU memory, because every GPU stores the same thing. FSDP fixes exactly that by sharding the model state across the data-parallel GPUs while keeping the same batch-splitting scheme.
Versus the other parallelism types
- FSDP is data parallelism with the weights, gradients and optimizer state sharded across the GPUs. Same data split, far less memory, more communication.
- Tensor parallelism splits each layer's matrix multiplications across GPUs. It cuts memory and latency but communicates inside every layer.
- Pipeline parallelism puts different layers on different GPUs and passes activations between them.
Large runs combine them. Megatron-LM's paper reports an 8.3-billion-parameter model trained with 8-way model parallelism inside each server and data parallelism across servers, up to 512 V100 GPUs, at 74% of linear scaling on that setup.
What it means when you pick a GPU
Data parallelism scales almost for free when the model plus optimizer state fit on one card and each GPU gets enough work between all-reduces. Fine-tuning a 7B-class model with a parameter-efficient method fits that profile, since the all-reduce only has to move the small set of trainable weights.
For full training, check two things. First, memory: 16 bytes per parameter means a 1-billion-parameter model needs about 16 GB of state per GPU, so a 24GB card tops out at roughly 1.5 billion parameters (24 / 16) and an 80GB card at roughly 5 billion, both before counting activations. Past that, shard with FSDP or split the model. Second, the link: all-reduce traffic is proportional to model size, so gradients of a large model across PCIe-only cards (such as the L40S or RTX 4090, which have no NVLink) will leave the GPUs waiting. Between servers, see NVLink vs InfiniBand. For multi-GPU nodes, see GPU clusters and our guide to AI compute clusters. Aquanode rents GPUs by the hour, so you can test whether your step time is compute-bound before you commit to a larger cluster.
Building on GPUs? Aquanode runs the workload.
Deploy on H100, H200, B200, A100 and MI300X across a multi-provider marketplace, without racking your own hardware or committing to one cloud's spec sheet.
See also
FSDP (Fully Sharded Data Parallel)
FSDP shards a model's weights, gradients and optimizer state across GPUs and gathers each layer only when needed, so models too big for one GPU can train.
Tensor Parallelism
Tensor parallelism splits each layer's matrix multiplications across several GPUs that work on every token together. It needs NVLink-class bandwidth.
Pipeline Parallelism
Pipeline parallelism puts different layers of a model on different GPUs and streams micro-batches through them. It needs little bandwidth but leaves idle gaps.
NCCL (NVIDIA Collective Communications Library)
NCCL is NVIDIA's library for all-reduce and other multi-GPU communication. PyTorch uses it to sync GPUs, and NVLink vs PCIe decides how fast it runs.
VRAM
VRAM is the memory attached to a GPU that holds the data it works on, and it caps which AI models fit. VRAM vs RAM, how to check yours, and how much AI needs.
BF16 (bfloat16)
BF16 is a 16-bit float with FP32's 8 exponent bits but only 7 mantissa bits. It is the default for training and costs 2 bytes per model parameter.