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.

Open in Lab
Compare three parallelism strategies. Data parallelism replicates the model. Pipeline parallelism slices by layer. Tensor parallelism slices within each layer.
The demo wakes as you arrive…

The core idea: slicing the MLP across GPUs

A Transformer MLP block has two linear layers with a activation in between. The computation is Y=GeLU(XA)Y = \text{GeLU}(XA) followed by Z=Dropout(YB)Z = \text{Dropout}(YB). The weight matrices AA and BB are the largest memory consumers. The question is: how do you split them across GPUs?

Option 1: split AA along its rows. This means each GPU gets a horizontal slice of AA and a corresponding vertical slice of XX. The problem is that GeLU is nonlinear: GeLU(X1A1+X2A2)≠GeLU(X1A1)+GeLU(X2A2)\text{GeLU}(X_1 A_1 + X_2 A_2) \neq \text{GeLU}(X_1 A_1) + \text{GeLU}(X_2 A_2). 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 AA along its columns. Each GPU gets a vertical slice: A=[A1,A2]A = [A_1, A_2]. Now each GPU independently computes GeLU(XAi)\text{GeLU}(X A_i), because the GeLU can be applied to each partition separately. No synchronization needed until the very end. Then split BB 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 .

[Y1,Y2]=[GeLU(XA1),  GeLU(XA2)][Y_1, Y_2] = [\text{GeLU}(XA_1), \;\text{GeLU}(XA_2)]
Column-parallel MLP — GeLU applied independently per GPU partition — By splitting AA along columns, each GPU computes its own GeLU(XAi)\text{GeLU}(XA_i) independently. This avoids a synchronization point before the nonlinearity — the key insight that makes Megatron's approach efficient.
Open in Lab
Watch how the MLP weight matrices are split across GPUs. Column-parallel on A avoids synchronization before GeLU; row-parallel on B requires only one all-reduce at the end.
The demo wakes as you arrive…

The communication pattern is elegantly simple. Two custom operators handle everything: ff is an identity in the and an all-reduce in the . gg is an all-reduce in the forward pass and an identity in the backward pass. These are conjugates — ff and gg 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 OO is split along its rows (just like BB 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 ff operator at the input, one gg 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.

Open in Lab
See how attention heads are distributed across GPUs. Each GPU computes its own subset of heads independently, with only one all-reduce needed after the output projection.
The demo wakes as you arrive…

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 ff: In the forward pass, ff simply passes the input through unchanged (identity). In the backward pass, ff 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 gg: In the forward pass, gg performs an all-reduce to sum partial outputs across GPUs. In the backward pass, gg 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, ff and gg form a conjugate pair. They are the only communication needed. No custom compiler, no special framework — just two custom PyTorch autograd.Function classes.

The f and g operators — Megatron's entire communication layer in PyTorchpython

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 output

Scaling 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.

Open in Lab
Explore how throughput scales with GPU count. 8-way tensor parallelism within a node combined with data parallelism across nodes achieves 76% efficiency at 512 GPUs.
The demo wakes as you arrive…

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.

Open in Lab
Compare post-norm (original BERT) vs pre-norm (Megatron-LM fix). Watch how activation magnitudes behave across 24 layers.
The demo wakes as you arrive…

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

  1. 2018

    GPipe (Google) — Pipeline parallelism

    Split the model into sequential stages across GPUs. Effective but limited by pipeline bubbles and required a custom framework.

  2. 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.

  3. 2019

    ZeRO (DeepSpeed) — Optimizer state sharding

    Partitions optimizer states, gradients, and parameters across data-parallel ranks. Orthogonal to tensor parallelism — the two combine powerfully.

  4. 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.

  5. 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.

  6. 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