Skip to main content

Data vs Tensor vs Pipeline Parallelism: When to Use Each

Summary
Data, tensor and pipeline parallelism split the batch, weight matrices or layers across GPUs. Compare memory, traffic and batch limits, then combine them.

Data, tensor and pipeline parallelism are three ways to spread one training job across GPUs, and each one splits something different. Data parallelism splits the batch, tensor parallelism splits every weight matrix, and pipeline parallelism splits the stack of layers. The split decides what each GPU stores, which exchange it waits on and how large the batch must be. The distributed parallelism page covers each one in depth, with code.

TL;DR

  • Data parallelism (DP) copies the model to every GPU and all-reduces the gradients once per step. Use it whenever the model fits on one GPU. Sharded DP (FSDP, ZeRO-3) stops copying the model state and fits far larger models for 1.5 times the traffic.
  • Tensor parallelism (TP) splits each matrix multiply across GPUs and all-reduces activations four times per layer for every micro-batch. It needs NVLink-class bandwidth, so it stays inside a node.
  • Pipeline parallelism (PP) gives each GPU a run of consecutive layers and sends activations point to point. Its traffic is light, but GPUs sit idle for (p − 1)/(m + p − 1) of each step unless there are many more micro-batches m than stages p.
  • 3D parallelism nests all three: TP inside a node, PP across nodes, and DP over whole copies of that TP × PP grid.

What each one splits

Data parallelism: split the batch

Every GPU holds a full replica and runs forward and backward on its own share of the batch. Before the optimizer step, the GPUs all-reduce their gradients to the average. A ring all-reduce sends and receives about 2(N − 1)/N times the gradient bytes per GPU, nearly independent of the GPU count N (ring all-reduce, step by step). PyTorch DDP hides most of it by reducing gradients in buckets while backward is still running.

The cost is memory. At 16 bytes per parameter, a 7B model needs 112 GB of model state on every GPU before a single activation, more than an 80 GB GPU holds.

Sharded data parallelism keeps the DP pattern but stops the copying. ZeRO (Rajbhandari et al., 2020) shards the optimizer state across the N ranks in stage 1, the gradients too in stage 2, and the weights as well in stage 3, which leaves 16Ψ/N bytes per GPU for Ψ parameters. PyTorch FSDP implements the stage 3 scheme. Before a layer runs, the ranks all-gather its weights and free them afterwards; after backward, a reduce-scatter leaves each rank the summed gradient for its shard alone. Stages 1 and 2 move the same bytes as plain DP. Stage 3 gathers the weights twice, for forward and again for backward, so it moves 1.5 times as much.

Tensor parallelism: split each matrix

Megatron-LM (Shoeybi et al., 2019) splits the matrix multiplies of each transformer layer across t GPUs in pairs, one pair for attention and one for the MLP. The first matrix of a pair is split by columns, so each GPU computes its own attention heads, or its own slice of the MLP's hidden units, with no exchange. The second is split by rows, so each GPU produces a partial sum of the full output, and one all-reduce adds the partials. That gives two all-reduces per layer in forward and two in backward, each over batch × sequence × hidden values.

All t GPUs work on the same micro-batch, so TP leaves the batch size alone. The price is that each all-reduce sits between two dependent operations, where it is hard to hide, and every matmul shrinks by a factor of t, which lowers GPU efficiency as t grows. Sequence parallelism (Korthikanti et al., 2022) swaps each all-reduce for a reduce-scatter and an all-gather that move the same total and also shard the layer-norm and dropout activations.

Pipeline parallelism: split the layers

Pipeline parallelism cuts the network into p stages of consecutive layers, one stage per GPU or per TP group. A stage sends its output activations to the next stage point to point, and in backward it sends the activation gradients back. Only tensors at stage boundaries move, so PP tolerates slower links.

With one batch in flight, all but one stage would wait, so GPipe (Huang et al., 2019) cuts the batch into m micro-batches that follow each other through the stages. The pipeline still has to fill and drain. The idle share of each step, the bubble, is

\text{bubble} = p - 1m + p - 1

Narayanan et al. (2021) write the same idle time as (p − 1)/m of the ideal compute time, and GPipe found it negligible once m ≥ 4p. The bubble is the same under the one-forward-one-backward (1F1B) schedule, which comes from PipeDream (Narayanan et al., 2019) and which Megatron-LM runs in a flushed, synchronous form. 1F1B starts each backward pass as early as possible, so a stage keeps activations for at most p micro-batches instead of m. Megatron-LM's interleaved schedule shrinks the bubble itself by giving each GPU several smaller stages, at the price of more sends.

Eight GPUs, five layouts

Pick a layout, then step through its exchanges; the strip between the nodes says whether a step crosses the slower link. Try the 7B model under DP, then under FSDP.

The 8-way TP layout spans both nodes on purpose, to show the arrangement production setups avoid.

Side by side

Data, tensor and pipeline parallelism compared
What is splitData parallelThe batch. Each GPU holds the whole model (DDP) or a 1/N shard of its state (FSDP, ZeRO).Tensor parallelEvery weight matrix, by columns and then by rows, across t GPUs.Pipeline parallelThe layer stack, into p stages of consecutive layers.
CommunicationData parallelAll-reduce of gradients once per step, overlapped with backward. FSDP: all-gather weights, reduce-scatter gradients.Tensor parallelAll-reduce of activations, 2 per layer in forward and 2 in backward, between dependent matmuls.Pipeline parallelPoint-to-point sends of boundary activations forward and their gradients backward, per micro-batch.
Links it needsData parallelTolerates the inter-node network while backward has enough work to hide the all-reduce.Tensor parallelNVLink or another fast link inside one node or NVLink domain.Pipeline parallelWorks across nodes; only boundary tensors move.
Model state per GPUData parallel16Ψ (DDP); 4Ψ + 12Ψ/N (ZeRO-1); 2Ψ + 14Ψ/N (ZeRO-2); 16Ψ/N (ZeRO-3, FSDP).Tensor parallelAbout 16Ψ/t. Activations inside the split regions also divide by t.Pipeline parallelAbout 16Ψ/p with balanced stages. The first stage holds activations for up to p micro-batches under 1F1B.
Compute efficiencyData parallelFull-size matmuls. Scales nearly linearly while the all-reduce stays hidden.Tensor parallelMatmuls shrink by t and wait on all-reduces, so efficiency falls as t grows, sharply across nodes.Pipeline parallelFull-size matmuls, but GPUs idle for the bubble and the slowest stage sets the pace.
Batch-size constraintData parallelGlobal batch of at least N micro-batches, so adding GPUs grows the batch.Tensor parallelNone: all t GPUs share one micro-batch.Pipeline parallelNeeds m much larger than p micro-batches per step (GPipe: m ≥ 4p).
Library supportData parallelPyTorch DDP and FSDP; DeepSpeed ZeRO stages 1 to 3; the Megatron-LM distributed optimizer.Tensor parallelMegatron-LM (Megatron Core); PyTorch tensor parallel APIs on DTensor; DeepSpeed through Megatron-DeepSpeed.Pipeline parallelPyTorch torch.distributed.pipelining (GPipe, 1F1B, interleaved); the DeepSpeed pipeline engine; Megatron-LM.

Ψ is the parameter count, N the data-parallel degree, t the tensor-parallel degree, p the number of stages and m the micro-batches per step.

Choosing: a decision flow

  1. Does the model fit on one GPU, with its model state and the activations of one micro-batch? Use DP with DDP, and add GPUs until the global batch reaches the largest size that still trains well. If it fits only barely, ZeRO stage 1 or 2 frees room for larger micro-batches at no extra traffic.
  2. Does it fit once the model state is divided by the GPU count? Use FSDP or ZeRO-3. It keeps the data-parallel programming model and needs no model changes. The 1.5 times traffic hides behind compute as long as each GPU has enough work per layer.
  3. Are single layers too large, activations too big, or the per-GPU batch too small to hide FSDP's gathers? Add TP inside the node, typically 2, 4 or 8 ways.
  4. Is the model still too large for one node's worth of TP? Add PP across nodes, with at least four micro-batches per stage.
  5. Is it one of the largest models, on hundreds or thousands of GPUs? Use 3D parallelism: TP × PP just large enough to hold one replica, and DP for everything beyond.

Steps 3 to 5 follow the takeaways of Narayanan et al. (2021): use tensor parallelism up to the number of GPUs in a server and pipeline parallelism across servers, keep the model-parallel size t · p only as large as memory requires, and scale out with data parallelism.

Combining them: 3D parallelism

In 3D parallelism the GPU count is the product DP × TP × PP, and every GPU belongs to one group along each dimension. In the figure's 3D layout, 8 = 2 × 2 × 2. GPUs 0 and 1 form a TP pair, GPUs 0 and 4 form a two-stage pipeline, and GPUs 0 and 2 hold the same shard in the two replicas and all-reduce its gradients.

Place the dimensions by how often they talk and whether the next operation waits on them:

  • TP talks every layer and blocks, so it gets the fastest links, inside a node or NVLink domain.
  • PP talks once per micro-batch per boundary and can overlap with other micro-batches, so its boundaries go where the network is slowest, between nodes.
  • DP talks once per step and overlaps with backward, so it goes outermost and spans the cluster.

The dimensions are coupled through the batch. The global batch is DP × m × micro-batch size, so a larger DP degree at a fixed global batch leaves fewer micro-batches per pipeline and a bigger bubble. The figure shows the trade at small scale. 3D gives each pipeline only 8 of the 16 micro-batches, but with 2 stages instead of 8 it idles 11% of the step, while pure PP, with all 16 micro-batches, idles 30%.

With this recipe, Narayanan et al. (2021) ran training iterations of a 1-trillion-parameter model on 3,072 A100 GPUs at 502 petaFLOP/s, 52% of the per-GPU peak. Mixture-of-experts models add expert parallelism, and long sequences add context parallelism, as further dimensions of the same grid. Rack-scale NVLink domains such as the 72-GPU GB200 NVL72 move the fast-link boundary from the node to the rack, which lets a TP group span more than eight GPUs at NVLink speed (interconnect details).

Production pitfalls

The pipeline bubble

With 8 stages and 8 micro-batches, GPUs idle for 7/15, or 47%, of every step; 32 micro-batches bring that to 18%. More micro-batches mean a larger global batch or smaller, less efficient micro-batches. Stage balance matters as much: the slowest stage sets the pace, and the first and last stages also carry the embedding and the output layer with the loss, so equal layer counts are not equal work.

TP all-reduce bandwidth

Four all-reduces per layer per micro-batch, each blocking the next operation, make TP the most sensitive of the three to link bandwidth. An H100's NVLink carries 900 GB/s counting both directions, about 450 GB/s each way, while an NDR InfiniBand port carries 50 GB/s each way, roughly a ninth. Keep TP groups inside the NVLink domain, and check the rank-to-GPU mapping at startup: a launcher that numbers ranks differently can place a TP group across nodes without any error.

DP gradient sync at scale

The ring moves about 2 times the gradient bytes per GPU at any scale, but other costs grow with N. A ring takes 2(N − 1) steps, so its latency grows with the GPU count, one reason NCCL also has tree algorithms. At a fixed global batch, each GPU gets less backward work to hide the all-reduce behind; at a fixed per-GPU batch, the global batch grows until it hurts convergence. And a synchronous all-reduce waits for the slowest GPU every step. Tuning DDP's bucket size helps, and so does moving some of the scale into TP and PP.

FSDP wrapping granularity

FSDP gathers one wrapped unit at a time. Units that are too small produce many small, latency-bound collectives; units that are too large raise the memory peak, since each is materialized in full. Wrapping each transformer block is the usual middle ground.

Primary sources

GPU & High-Performance Computing
Distributed Parallelism in Deep Learning

GPU distributed parallelism: Data Parallel (DDP), Tensor Parallel, Pipeline Parallel, and ZeRO optimization for training large AI models.

GPU & High-Performance Computing
Multi-GPU Communication: NVLink, PCIe, and NCCL

How GPUs talk: the bandwidth cliff from HBM to Ethernet, NVLink 5 and GB200 NVL72 topologies, ring AllReduce step by step, and choosing between NCCL, Gloo, and MPI.

GPU & High-Performance Computing
NCCL: How NVIDIA Collective Communication Works

A deep dive into NCCL internals: communicators and channels, how it picks ring/tree/NVLS algorithms and LL/LL128/Simple protocols, reading NCCL_DEBUG logs, and tuning and debugging distributed training.

GPU & High-Performance Computing
Slurm GPU Allocation for Distributed Training

Complete guide to GPU allocation on Slurm — --gres flags, CUDA_VISIBLE_DEVICES remapping, GPU topology and NVLink binding, MIG partitioning, production job scripts, and debugging common GPU errors.

GPU & High-Performance Computing
CUDA Contexts: Ownership, Current Stack, Isolation

Deep dive into the CUDA context object: control vs data plane, inventory (memory, modules, streams, events, graphs), push/pop/setCurrent stacks, primary retain/release, flags and limits, isolation, cost, and traps.

GPU & High-Performance Computing
CUDA Context vs Streams vs MPS: Which Layer Fixes What

Decision map for CUDA: a context is per-process GPU state, a stream is an in-order queue inside a context, and MPS shares one context across processes. Pick the layer that matches the problem.

If you found this explanation helpful, consider sharing it with others.

Mastodon