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:

  1. Each GPU gets 1/N of the batch and runs the forward and backward pass on its slice, producing its own gradients.
  2. 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.
  3. 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

Submit the job. Everything after that is ours.

Sign up in 60 seconds. Pay for the GPU minutes you actually use.

© 2026 Aquanode. All rights reserved.

All trademarks, logos and brand names are the property of their respective owners.