Skip to content
Road to Intelligence

Concept · Chapter 9: How an LLM Is Actually Built

Tensor and Pipeline Parallelism

Must knowUnderstand13 minDifficulty

Model parallelism splits the model itself: tensor parallelism divides each layer's matrices across GPUs, pipeline parallelism gives each GPU a block of consecutive layers, and large runs combine both with data parallelism.

The problem

When a single layer, or a single sequence's activations, is too big for one GPU, sharding the optimizer states is not enough; the computation itself must be divided.

The solution

Split matrix multiplications column- and row-wise across the GPUs of one server (tensor parallelism), chain servers into a pipeline of layer blocks fed with many micro-batches (pipeline parallelism), and replicate the whole arrangement for data parallelism.

The consequence

Trillion-parameter training became possible, at the price of heavy communication inside servers, idle 'bubbles' in pipelines, and complex system design.

Tensor parallelism: split the matrices

A layer's work is mostly matrix multiplication. Split a weight matrix A by columns into A₁ and A₂ on two GPUs, and each computes half of the output, XA₁ and XA₂, from the same input X. Megatron-LM arranged the splits inside the attention heads and the feed-forward block so that each layer needs only a couple of all-reduce operations, and trained an 8.3-billion-parameter model on 512 GPUs at 76% scaling efficiency Established.

The communication happens inside every layer, many times per step, so it needs the fastest links available. Narayanan and colleagues' guideline: use tensor parallelism up to the number of GPUs in one server, and pipeline parallelism across servers Established.

Pipeline parallelism: an assembly line

Give GPU 1 the first quarter of the layers, GPU 2 the next quarter, and so on. Activations travel forward from GPU to GPU, gradients travel back. Only neighbours talk, and the messages are small.

The problem is idleness: GPU 4 waits while the first input passes through GPUs 1–3. The fix is to split each batch into many micro-batches so the stages overlap. With K stages and M micro-batches, each GPU sits idle for a fraction

bubble=K−1M+K−1\text{bubble} = \frac{K-1}{M+K-1}

of the step. Tiny example. With 4 stages and 4 micro-batches: 3 ÷ 7 ≈ 43% idle. With 16 micro-batches: 3 ÷ 19 ≈ 16%. GPipe found the bubble overhead negligible once there were at least four times as many micro-batches as stages Established.

Four GPUs, one four-layer model · what each GPU holds
GPU 1batch slice 1L1L2L3L4GPU 2batch slice 2L1L2L3L4GPU 3batch slice 3L1L2L3L4GPU 4batch slice 4L1L2L3L4all-reduce gradients once per step
Each GPU holds
A full copy of the model, and a different slice of the batch.
What they exchange
Once per step, every GPU averages its gradients with all the others (an all-reduce), so all copies stay identical.
Used
Always, as the outermost layer of parallelism: it is how a run uses thousands of GPUs.
The catch
The whole model, its gradients and optimizer states must fit on every GPU.

Real runs combine them. Llama 3 405B used tensor parallelism across the 8 NVLink-connected GPUs of each server, 16 pipeline stages, and sharded data parallelism (FSDP) across the rest, plus a fourth kind, context parallelism, for very long sequences: up to 16,384 GPUs.

Putting it together

Real runs nest the three: tensor parallelism inside each server, pipeline stages across servers, and data parallelism across copies of the whole pipeline. Narayanan and colleagues trained a 1-trillion-parameter model this way at 502 petaFLOP/s on 3,072 GPUs, 52% of the hardware's theoretical peak Established. Llama 3 405B used tensor parallelism of 8 within each NVLink-connected server, 16 pipeline stages and 128-way sharded data parallelism on 16,384 GPUs, plus context parallelism for very long sequences Established.

What to remember

  • Tensor parallel: split each weight matrix; all GPUs work on the same tokens; combine results inside every layer.
  • Keep tensor parallelism inside one server, where GPUs share very fast links.
  • Pipeline parallel: each GPU holds consecutive layers; micro-batches flow through like an assembly line.
  • Bubble: (K − 1) / (M + K − 1) idle time for K stages and M micro-batches.
  • Llama 3 405B: tensor 8 × pipeline 16 × data 128 = 16,384 GPUs.

Key papers

Important

Megatron-LM: Training Multi-Billion Parameter Language Models Using Model Parallelism

Mohammad Shoeybi, Mostofa Patwary et al. · 2019

Introduced the tensor-parallel layout for Transformers that most large training systems still use: split each layer's matrices across GPUs.

How to read it: Figure 3 (how an MLP and an attention block are split) is the part to understand.

~25 min readarXiv:1909.08053✓ verified 2026-10-04
Important

Efficient Large-Scale Language Model Training on GPU Clusters Using Megatron-LM

Deepak Narayanan, Mohammad Shoeybi et al. · 2021

Showed how to compose tensor, pipeline and data parallelism to train trillion-parameter models efficiently on thousands of GPUs.

How to read it: Read the takeaways in Section 3; they summarise how to choose the parallel sizes.

~45 min readarXiv:2104.04473✓ verified 2026-10-04
Essential

The Llama 3 Herd of Models

Aaron Grattafiori, Abhimanyu Dubey et al. · 2024

The most complete public account of building a frontier-scale model end to end: data pipeline, scaling-law experiments, 16,384-GPU training, failures and all.

How to read it: It is 90+ pages. For this chapter read Section 3 (pre-training) only: data, scaling laws, infrastructure and the training recipe.

~2 h readarXiv:2407.21783✓ verified 2026-10-04

Watch