Model Efficiency & Scaling2019advanced12 min read
Megatron-LM: Training Multi-Billion Parameter Language Models Using Model Parallelism
Megatron-LM: تدريب نماذج لغوية بمليارات المعاملات باستخدام التوازي على مستوى النموذج
Shoeybi, M. · Patwary, M. · Puri, R. · LeGresley, P. · Casper, J. · Catanzaro, B. — arXiv
The problem
By 2019, it was clear that larger models yield better NLP performance. GPT-2 had 1.5 billion parameters, and scaling further promised even better results. But a single has limited memory — typically 16 or 32 GB — far too little to hold the weights, gradients, optimizer states, and activations of a multi-billion model. alone cannot solve this: it replicates the entire model on every GPU, so if the model does not fit on one GPU, it does not fit on any. Existing model-parallel solutions like GPipe required custom compilers or framework rewrites, making them hard to adopt.
The contribution
An efficient intra-layer approach for Transformers that requires no custom compiler — just a few communication primitives inserted into native PyTorch. The key insight is to split the matrices within each Transformer layer (MLP and ) using column-then-row partitioning, so that each GPU holds a slice of every layer and only one per sub-block is needed. Using this approach plus mixed-precision , the authors trained an 8.3B-parameter GPT-2 model on 512 GPUs achieving 15.1 PetaFLOPs at 76% . They also showed that rearranging in BERT-style models is critical for scaling beyond BERT-Large, and their 3.9B BERT model achieved state-of-the-art on the RACE dataset.
The impact
Megatron-LM became the standard infrastructure for training large language models at NVIDIA and beyond. Its tensor parallelism approach was adopted by virtually every subsequent large-scale training system, including the training of GPT-3, BLOOM, and many other billion-parameter models. Combined with DeepSpeed's ZeRO optimizer, it enabled 3D parallelism (data + tensor + pipeline) that powers today's largest models. The open-source codebase remains one of the most widely used foundations for distributed LLM training.
Imagine a restaurant kitchen preparing an enormous banquet. Data parallelism is like having ten identical kitchens each cooking the full menu from their own ingredients — fast, but every kitchen needs all the equipment. is an assembly line: one station chops, the next sautés, the next plates — fast, but stations wait idle while others finish.
Megatron-LM's tensor parallelism is different: it takes each cooking station and splits the work within it. Instead of one chef chopping all the vegetables, four chefs stand at the same cutting board, each chopping a quarter of the pile simultaneously. They only need to combine their piles once before passing to the next station. Minimal waiting, maximum parallel chopping — and no special kitchen redesign required.
The memory wall: why data parallelism is not enough
In data parallelism, the full model is copied to every GPU. Each GPU processes a different mini-batch, computes gradients, and then all GPUs synchronize their gradients via an all-reduce operation. This works beautifully — until the model is too large to fit on a single GPU.
Consider a model with 2 billion parameters in FP16. The weights alone consume 4 GB. Add FP32 optimizer states ( stores two moments per parameter), gradients, and activations, and a single GPU needs 40+ GB — exceeding the memory of most GPUs available in 2019.
Pipeline parallelism, used by GPipe, partitions the model into sequential stages placed on different GPUs. But this creates a "pipeline bubble" — GPUs sit idle waiting for activations from previous stages. The deeper the pipeline, the worse the bubble.
What was needed was a way to split each layer across GPUs, so that every GPU participates in every layer simultaneously. This is tensor parallelism — and Megatron-LM made it practical for Transformers.
The core idea: slicing the MLP across GPUs
A Transformer MLP block has two linear layers with a activation in between. The computation is followed by . The weight matrices and are the largest memory consumers. The question is: how do you split them across GPUs?
Option 1: split along its rows. This means each GPU gets a horizontal slice of and a corresponding vertical slice of . The problem is that GeLU is nonlinear: . You would need to synchronize (all-reduce) before the GeLU — an extra communication step right in the middle of the computation.
Option 2 (Megatron's choice): split along its columns. Each GPU gets a vertical slice: . Now each GPU independently computes , because the GeLU can be applied to each partition separately. No synchronization needed until the very end. Then split along its rows so it directly accepts the partitioned output from GeLU. The only synchronization is a single all-reduce after the second matrix multiply, right before .
The communication pattern is elegantly simple. Two custom operators handle everything: is an identity in the and an all-reduce in the . is an all-reduce in the forward pass and an identity in the backward pass. These are conjugates — and together ensure that the forward pass produces correct outputs and the backward pass produces correct gradients, with exactly one all-reduce per direction per MLP block. The entire implementation fits in a few lines of PyTorch.
Attention: a natural fit for parallelism
Multi-head is inherently parallel. Each attention head operates on a different learned projection of the input — the heads are independent computations that only need to be concatenated at the end. Megatron exploits this by partitioning the Q, K, V projection matrices column-wise so that each GPU computes a subset of attention heads.
Concretely, if a model has 16 attention heads and 4 GPUs, each GPU handles 4 heads. The Q, K, V weight matrices are split so each GPU has its portion of heads. After computing attention, each GPU has partial output. The output linear projection is split along its rows (just like in the MLP), so the partial outputs feed directly in without synchronization. A single all-reduce after the output projection combines results — identical to the MLP pattern.
The beauty of this design is that the attention block and the MLP block use the same communication pattern: one operator at the input, one operator at the output. The total for a full Transformer layer is just 2 all-reduces in the forward pass and 2 in the backward pass.
The f and g operators: two lines of code that enable everything
The elegance of Megatron's approach lies in two custom autograd operators. Think of them as traffic signals at the boundary of each parallel block:
Operator : In the forward pass, simply passes the input through unchanged (identity). In the backward pass, performs an all-reduce to sum gradients across GPUs. This is placed before each parallel block — each GPU receives the full input and computes its own partition.
Operator : In the forward pass, performs an all-reduce to sum partial outputs across GPUs. In the backward pass, passes gradients through unchanged (identity). This is placed after each parallel block — the partial results from all GPUs are combined into the full output.
Together, and form a conjugate pair. They are the only communication needed. No custom compiler, no special framework — just two custom PyTorch autograd.Function classes.
Simplified to show the idea — not the real implementation.
import torch
import torch.distributed as dist
class f(torch.autograd.Function):
"""Forward: identity. Backward: all-reduce gradients."""
@staticmethod
def forward(ctx, x):
return x # pass through unchanged
@staticmethod
def backward(ctx, grad):
dist.all_reduce(grad) # sum gradients across GPUs
return grad
class g(torch.autograd.Function):
"""Forward: all-reduce partial outputs. Backward: identity."""
@staticmethod
def forward(ctx, x):
dist.all_reduce(x) # sum partial results from all GPUs
return x
@staticmethod
def backward(ctx, grad):
return grad # pass through unchanged
# Usage in a tensor-parallel MLP:
# x = f.apply(x) # each GPU gets full input
# y = gelu(x @ A_i) # each GPU computes its column partition
# z = y @ B_i # each GPU computes its row partition
# z = g.apply(z) # all-reduce → full outputScaling results: 512 GPUs, 76% efficiency
The authors trained GPT-2 models of increasing size: 1.2B, 2.5B, 4.2B, and 8.3B parameters. All used 8-way tensor within each DGX node (8 GPUs connected by NVLink) and data parallelism across nodes.
The 8.3B model on 512 GPUs achieved 15.1 PetaFLOPs sustained — 76% scaling efficiency compared to a strong single-GPU baseline of 39 TeraFLOPs (30% of V100 peak). This was remarkable for 2019: most existing approaches suffered significant efficiency drops at this scale.
A key engineering decision was keeping tensor-parallel groups within a single NVLink domain (one server). NVLink provides 300 GB/s bandwidth between GPUs in a server, while inter-node connections are typically 10–25× slower. By confining the communication-heavy tensor parallelism to the fast interconnect and using the less-communication-intensive data parallelism across nodes, Megatron achieved near-linear scaling.
Fixing BERT: rearranging layer normalization for scale
Scaling BERT beyond 336 million parameters had proven problematic — training became unstable and performance degraded rather than improved. The Megatron-LM authors diagnosed the issue: in the original BERT architecture, layer normalization is applied after the (post-norm). This means the values flowing through the residual path keep growing with depth, eventually causing instability.
The fix was simple but critical: move layer normalization to before the self-attention and MLP blocks (pre-norm), following the pattern used in GPT-2. With this change, the residual path remains clean — it carries unnormalized values that add linearly — and normalization is applied to the inputs of each sub-layer, keeping activations bounded.
With this rearrangement, the authors successfully trained BERT models up to 3.9 billion parameters — 12× larger than the original BERT-Large — and achieved state-of-the-art results on the RACE reading comprehension benchmark with 90.9% accuracy.
Benchmark results: bigger models, better performance
The GPT-2 models showed a clear trend: larger models converge faster and to lower perplexity. The 8.3B model achieved state-of-the-art results across multiple benchmarks:
- WikiText-103: Perplexity of 10.8, down from the previous best of 15.8 — a 32% relative improvement.
- LAMBADA: 66.5% accuracy, up from 63.2% — demonstrating improved long-range language understanding.
- RACE (reading comprehension): 90.9% accuracy with the 3.9B BERT model, surpassing the previous best of 89.4%.
These results confirmed a key hypothesis: given enough compute and an efficient parallelism strategy, simply making Transformer models larger yields significant quality improvements — a theme that would define the next era of NLP with GPT-3 and beyond.
Legacy: the distributed training playbook
2018
GPipe (Google) — Pipeline parallelism
Split the model into sequential stages across GPUs. Effective but limited by pipeline bubbles and required a custom framework.
2019
Megatron-LM — Tensor parallelism
Intra-layer model parallelism that splits weight matrices within each Transformer layer. Native PyTorch, 8.3B parameters, 76% scaling efficiency on 512 GPUs.
2019
ZeRO (DeepSpeed) — Optimizer state sharding
Partitions optimizer states, gradients, and parameters across data-parallel ranks. Orthogonal to tensor parallelism — the two combine powerfully.
2020
GPT-3 (175B parameters)
Trained using Megatron-style tensor parallelism combined with data parallelism. Proved that scaling to 175 billion parameters unlocks few-shot learning.
2021
3D parallelism era
Megatron-DeepSpeed combined tensor + pipeline + data parallelism. This "3D parallelism" became the standard recipe for training models at the hundred-billion scale and beyond.
2022
BLOOM (BigScience, 176B)
The largest open-source language model at the time, trained using Megatron-DeepSpeed on 384 A100 GPUs. Direct descendant of Megatron-LM's techniques.
Megatron-LM's impact goes beyond any single model. It established the engineering playbook for distributed LLM training: use tensor parallelism within a node (fast interconnect), pipeline parallelism across a small number of nodes, and data parallelism across the full cluster. Every major training system today — from Meta's LLaMA infrastructure to Google's PaLM training setup — follows this template. The insight that a few all-reduce operations in native PyTorch could replace custom compilers democratized billion-parameter training.
CitationShoeybi, Patwary, Puri, LeGresley, Casper, Catanzaro. Megatron-LM: Training Multi-Billion Parameter Language Models Using Model Parallelism. arXiv, 2019.
Terms in this paper
- Tensor Parallelismتوازي مصفوفات الموتّرات
- Model Parallelismتوازي النموذج
- Data Parallelismتوازي البيانات
- Pipeline Parallelismالتوازي المتسلسل للطبقات
- Allreduceالاختزال الشامل
- Scaling Efficiencyكفاءة التوسّع
- Mixed Precision Trainingالتدريب بالدقة المختلطة
- Throughputمعدل التدفق والإنتاجية
- FLOPsالعمليات الحسابية العائمة
- GPUوحدة معالجة الرسوميات
- Gradient Accumulationتجميع التدرجات الحسابية
- Layer Normalizationالتسوية الطبقية
- Residual Connectionالوصلة التجاوزية
- GELUوحدة الخطأ الخطية الغاوسية (GELU)
- Transformerالمحوِّل