Megatron-LM is NVIDIA's open source codebase for training large transformer models across hundreds or thousands of GPUs. Its lasting contribution is a set of parallelism techniques that now appear in almost every large training stack: tensor parallelism that splits each layer's matrices across GPUs, pipeline parallelism with interleaved schedules, sequence parallelism for activation memory, and a distributed optimizer that shards optimizer state across data-parallel ranks. The reusable parts live in a library called Megatron-Core; the repository also ships training scripts such as pretrain_gpt.py that drive it from command-line flags.

This article treats Megatron as a system: how the parallel dimensions map onto ranks, sizing a 70B run on 512 GPUs, a launch command checked against the current repository, and the overlap, checkpoint and failure-mode details of practice. The layer-split algebra is in tensor parallelism in depth and the schedules in pipeline parallelism.

Where Megatron came from

Three papers define what Megatron does. Shoeybi and colleagues (2019) split the first MLP matrix by columns and the second by rows, and attention by heads, so each transformer block needs two all-reduces forward and two backward; they trained an 8.3 billion parameter model on 512 GPUs.

Narayanan and colleagues (2021) combined tensor, pipeline and data parallelism, which they called PTD-P, added the interleaved pipeline schedule, and reported training a 1 trillion parameter model on 3,072 A100 GPUs at 502 petaFLOP/s, 52 percent of theoretical peak. Korthikanti and colleagues (2022) added sequence parallelism and selective activation recomputation, which together cut activation memory about fivefold and removed most of the cost of full recomputation. Later additions such as context and expert parallelism extend that framework rather than replace it.

The parallel dimensions and what each costs

Each parallel dimension divides different state and pays for it with different communication. Megatron lets you set each one independently; the product of the model-parallel sizes and the data-parallel size must equal the world size.

DimensionFlagWhat is splitCommunicationWhere it belongs
Tensor (TP)--tensor-model-parallel-sizeweight matrices inside each layerall-reduce, or all-gather plus reduce-scatter with SP, inside every layerwithin one NVLink domain
Sequence (SP)--sequence-parallelactivations of norm and dropout along the sequence, over the TP groupreplaces each TP all-reduce with the same bytes as gather plus scatteralways on when TP > 1
Pipeline (PP)--pipeline-model-parallel-sizelayers into stagespoint-to-point activation sends per micro-batchacross nodes
Context (CP)--context-parallel-sizethe sequence for attentionexchange of key and value blocks between CP rankslong sequences only
Expert (EP)--expert-model-parallel-sizeMoE expertsall-to-all token dispatch and combineMoE models
Data (DP)derivedthe batchgradient reduce-scatter or all-reduce once per stepeverything left over

TP talks inside every layer, so it stays on the fastest links: at most the GPUs per node. PP sends one activation tensor per micro-batch per stage boundary, so it tolerates inter-node links. DP communicates once per step, overlapped with backward, and absorbs the remaining GPUs. The collectives are described in NCCL collectives.

How ranks map onto GPUs

512 GPUs as TP=8 x DP=16 x PP=4 (Megatron default order tp-cp-ep-dp-pp, cp=ep=1)stage 0layers 0-19 + embeddingnode 0ranks 0-7 (TP group)node 1ranks 8-15...node 15ranks 120-127stage 1layers 20-39node 16ranks 128-135 (TP group)node 17ranks 136-143...node 31ranks 248-255stage 2layers 40-59node 32ranks 256-263 (TP group)node 33ranks 264-271...node 47ranks 376-383stage 3layers 60-79 + output, lossnode 48ranks 384-391 (TP group)node 49ranks 392-399...node 63ranks 504-511Horizontal arrows: pipeline send/recv of activations between stages (point-to-point, inter-node)Inside each box: 8 GPUs on NVLink run tensor parallelism with all-reduce or reduce-scatter/all-gatherDown each column: 16 data-parallel replicas of that stage reduce gradients after backwardone training step64 micro-batches per replica flow through 4 stages, then DP gradient reduction
A 512-GPU job as 4 pipeline stages, each replicated 16 times; every replica of a stage is one 8-GPU node running tensor parallelism.

Megatron-Core builds its process groups in parallel_state.initialize_model_parallel, whose order argument defaults to tp-cp-ep-dp-pp. The first dimension in that string varies fastest across global ranks, so consecutive ranks form a TP group. Because torchrun numbers the GPUs of one node consecutively, a TP size of 8 on 8-GPU nodes puts each TP group on one node's NVLink, which is exactly what you want. With context and expert parallelism both at 1, a rank's coordinates reduce to simple arithmetic:

def coords(rank, tp=8, pp=4, world=512):
    """Megatron default order tp-cp-ep-dp-pp with cp = ep = 1."""
    dp = world // (tp * pp)
    return {"tp": rank % tp, "dp": (rank // tp) % dp, "pp": rank // (tp * dp)}

def groups_of(rank, tp=8, pp=4, world=512):
    dp = world // (tp * pp)
    c_ = coords(rank, tp, pp, world)
    base = c_["pp"] * tp * dp
    return {
        "tp_group": [base + c_["dp"] * tp + i for i in range(tp)],
        "dp_group": [base + d * tp + c_["tp"] for d in range(dp)],
        "pp_group": [s * tp * dp + c_["dp"] * tp + c_["tp"] for s in range(pp)],
    }

g = groups_of(0)
print(g["tp_group"])       # [0, 1, ..., 7]       one node
print(g["dp_group"][:3])   # [0, 8, 16]           same stage, other nodes
print(g["pp_group"])       # [0, 128, 256, 384]   the four stages

When context or expert parallelism is enabled the layout changes: expert groups are built over a separate decomposition, so do not reuse this formula for MoE runs. Ask Megatron instead; parallel_state.get_tensor_model_parallel_rank(), get_pipeline_model_parallel_rank() and get_data_parallel_rank() return the truth for the running job, and logging them with the hostname at start-up settles most placement arguments.

Worked example: 70B parameters on 512 GPUs

Take a dense decoder with about 70 billion parameters, 80 layers, hidden size 8,192 and sequence length 4,096, trained on 64 nodes of 8 GPUs with 80 GB each. Choose TP=8 to fill each node, PP=4, which leaves DP = 512 / (8 x 4) = 16.

Model state. Each GPU holds 1/32 of the weights: about 70e9 / 32 = 2.19 billion parameters, ignoring the embedding imbalance discussed below. Megatron's distributed optimizer documentation gives the bytes per parameter for bf16 parameters with fp32 main gradients as 18 without the distributed optimizer and 6 + 12/d with it, where d is the data-parallel size. That is 2.19e9 x 18 = 39.4 GB per GPU without it, and 2.19e9 x (6 + 12/16) = 14.8 GB with it. The saving is the ZeRO stage 1 and 2 idea applied inside Megatron; see the ZeRO optimizer for the derivation.

Activations. Korthikanti and colleagues estimate activation memory per layer per micro-batch as s.b.h.34 / t bytes when tensor and sequence parallelism are on and the attention score matrix is recomputed, with s the sequence length, b the micro-batch size, h the hidden size and t the TP size. Here that is 4,096 x 1 x 8,192 x 34 / 8 = 143 MB per layer. A stage has 20 layers, and under the 1F1B schedule the first stage keeps activations for up to p = 4 micro-batches in flight, so about 20 x 4 x 143 MB = 11.4 GB. The interleaved schedule raises that somewhat on the first stage.

Bubble. With a global batch of 1,024 sequences (about 4 million tokens) and micro-batch 1, each data-parallel replica processes m = 1,024 / 16 = 64 micro-batches per step. The 1F1B pipeline bubble is (p - 1) / m = 3 / 64, or 4.7 percent of step time idle. With the interleaved schedule and v virtual stages per GPU it falls to (p - 1) / (v.m); four chunks of 5 layers give 3 / 256, about 1.2 percent, at the price of v times as many pipeline sends.

About 26 GB of state and activations leaves headroom on 80 GB GPUs. Spend it deliberately: drop PP to 2, raise the micro-batch size, or turn off recomputation, measuring each over a few hundred steps.

A launch command

The flags below appear in the repository's own example scripts at the time of writing. Megatron now generates many arguments from configuration dataclasses, so run pretrain_gpt.py --help on your checkout and fail loudly on unknown flags rather than trusting an old blog post, including this one.

export CUDA_DEVICE_MAX_CONNECTIONS=1     # keeps comm and compute kernel order predictable for overlap

torchrun --nproc_per_node 8 --nnodes 64 --node_rank $NODE_RANK \
  --master_addr $MASTER_ADDR --master_port 29500 \
  pretrain_gpt.py \
  --num-layers 80 --hidden-size 8192 --num-attention-heads 64 \
  --seq-length 4096 --max-position-embeddings 4096 \
  --tensor-model-parallel-size 8 --sequence-parallel \
  --pipeline-model-parallel-size 4 --num-layers-per-virtual-pipeline-stage 5 \
  --micro-batch-size 1 --global-batch-size 1024 \
  --bf16 --use-distributed-optimizer \
  --overlap-grad-reduce --overlap-param-gather \
  --recompute-activations \
  --attention-backend auto \
  --lr 1.5e-4 --min-lr 1.5e-5 --lr-decay-style cosine --lr-warmup-fraction 0.01 \
  --clip-grad 1.0 --weight-decay 0.1 --train-iters 250000 \
  --data-path $DATA_PREFIX --split 949,50,1 \
  --save $CKPT --load $CKPT --save-interval 1000 --ckpt-format torch_dist \
  --log-interval 10 --eval-interval 1000 --eval-iters 10

Megatron checks at start-up that the global batch divides by micro-batch size times DP, the layers by PP and the heads by TP. --recompute-activations selects selective recomputation. Tokenizer flags are omitted; they depend on how you preprocessed the corpus.

Overlap and memory settings

Five settings decide whether a correctly sized run is also a fast one.

  • Distributed optimizer. --use-distributed-optimizer replaces the gradient all-reduce with a reduce-scatter, steps only the local shard of fp32 master weights and Adam moments, then all-gathers the updated bf16 parameters. Same bytes on the wire, far less memory.
  • Gradient reduce overlap. --overlap-grad-reduce launches the reduction for a bucket of gradients as soon as backward has produced it, instead of after the last layer.
  • Parameter gather overlap. --overlap-param-gather hides the all-gather of updated parameters behind the next step's forward pass.
  • Sequence parallelism. Without it every TP rank stores full copies of the norm and dropout activations. With it those are split along the sequence, at no extra communication volume.
  • Selective recompute. Recomputing only the attention score and softmax path costs a few percent of compute, compared with roughly a third for full recomputation of every layer.

Profiling overlap is covered in collective overlap. The quickest check compares step time with a FLOPs estimate: at 4 million tokens per step, a 70B dense model needs about 6 x 70e9 x 4.2e6 = 1.76e18 FLOPs, so on 512 GPUs at a sustained 400 TFLOP/s each the step takes about 8.6 seconds. If yours takes 14, a third of the machine is idle.

Checkpoints and resuming

Megatron checkpoints are sharded by construction, because no rank holds the whole model. The older torch format writes one file per model-parallel rank and is tied to the TP and PP sizes that wrote it. The torch_dist format uses Megatron-Core's distributed checkpointing, which records each tensor's global shape and the slice each rank owns, so a checkpoint saved at TP=8, PP=4 can be loaded into a different layout. Use it for anything you expect to fine-tune, convert, or resume on a different cluster.

Two rules matter more than the format. At 512 GPUs hardware failures are routine, so save at least hourly. And prove resume works early: save at step 100, kill the job, resume, and compare the loss at step 101 with an uninterrupted run, which tests data position, schedule state and random generators.

Failure modes

Most Megatron incidents fall into a short list.

  • Hang at start-up. One node launched with different flags or a different code revision builds different process groups and waits forever in a collective. Hash the full argument list and the git revision on every rank and compare before the first step.
  • Out of memory only on stage 0 or the last stage. The first stage holds the input embedding and the most in-flight activations; the last holds the output layer and the loss over the whole vocabulary. A 128,000-token vocabulary at hidden size 8,192 is about a billion parameters per matrix, roughly 130 million per TP rank, on those stages alone; Megatron also rounds the vocabulary up to a multiple of --make-vocab-size-divisible-by (default 128) times TP. Uneven layer placement, or moving layers off the end stages, fixes it.
  • Loss differs after a parallelism change. Reduction order changes, so small drift is expected; a jump at the first step after resume means a sharding or loading bug.
  • Slow steps with no error. One slow GPU or link stalls every collective. Track per-rank step times and NCCL timings and drain the outlier node rather than tuning the whole job.

Trade-offs against FSDP

Megatron fits when sharding parameters alone is not enough: one layer strains a GPU, or DP-only scaling has run out of batch size. Its costs: the model must use Megatron's parallel layers, checkpoint conversion to Hugging Face format is a separate step, and flags interact.

PyTorch FSDP takes the opposite approach: keep the model code unchanged, shard parameters, gradients and optimizer state across data-parallel ranks, and gather each layer's weights just in time. Below a few tens of billions of parameters on a modest cluster it is usually simpler and nearly as fast; see PyTorch FSDP in depth. Above that, or at thousands of GPUs, the per-layer weight gathers of pure sharding grow expensive and the 3D layout wins.

What to do next

Use this list before your first multi-node Megatron run.

  1. Write down the model shape, GPU count, GPUs per node and memory per GPU, then choose TP as the smallest power of two up to GPUs per node that fits a layer, PP as small as fits, and DP as the rest.
  2. Compute bytes per GPU with 6 + 12/d for model state and s.b.h.34/t per layer for activations; add the embedding and output layer to the end stages.
  3. Compute micro-batches per replica and the bubble (p - 1) / m; aim for under 5 percent, or use the interleaved schedule.
  4. Run pretrain_gpt.py --help on your checkout and confirm every flag you plan to use.
  5. Turn on the distributed optimizer, sequence parallelism, gradient and parameter overlap, and selective recomputation, and set CUDA_DEVICE_MAX_CONNECTIONS=1.
  6. Log each rank's TP, PP and DP coordinates with its hostname, and hash flags and code revision across ranks.
  7. Compute expected step time from 6 x parameters x tokens, compare it with measured time, and profile the gap before scaling up.
  8. Save in torch_dist format, and test kill-and-resume at step 100.
Key takeaway: Megatron-LM composes tensor parallelism inside an NVLink domain, pipeline parallelism across nodes and data parallelism over the rest, with sequence parallelism, selective recomputation and a distributed optimizer to fit memory. Size a run with 6 + 12/d bytes per parameter, s.b.h.34/t bytes per layer of activations and a bubble of (p - 1)/m, verify every flag against your checkout, and test resume before the run matters.