DeepSpeed is Microsoft's open-source library for training large PyTorch models across many GPUs. Most people meet it through the ZeRO optimizer and treat the rest as a black box. That works until a job hangs at startup, a CUDA extension fails to compile on the tenth node, or the pipeline engine refuses a config that ran fine without it. Debugging those needs a picture of the whole library: what it wraps, what it takes over from PyTorch, and which parts are optional.

This article is that picture: the engine and what its backward and step really do, the config and its batch-size rule, the launcher, the compiled operators, the pipeline engine and the mixture-of-experts layer. ZeRO gets one section and a pointer to its own page. It ends with a worked plan for a 13B model, a comparison with FSDP, Accelerate and Megatron, and the failures you will actually see.

Advertisement

The component map

DeepSpeed is several subsystems behind one entry point. The runtime engine wraps your model and optimizer and owns mixed precision, loss scaling, gradient accumulation, clipping and the ZeRO partitioner. The launcher starts one process per GPU on every node. The op builder compiles fused kernels such as FusedAdam and the CPU Adam used for offload. The pipeline engine handles models split into sequential stages, the MoE layer adds expert parallelism, and a checkpoint engine saves the sharded state.

Everything is driven by one JSON config, so the same script can run as plain data parallel on one GPU or as ZeRO stage 3 with offload on 64 GPUs by changing the file, not the code. That is also the main source of confusion: behaviour you expect to see in Python lives in a config key.

What deepspeed.initialize builds around your nn.Module (one process per GPU)deepspeed launcherhostfile, env, ranksds_config.jsonbatch, precision, ZeRODeepSpeedEnginewraps model + optimizerspawnsconfiguresprecision + loss scalingfp16 / bf16ZeRO partitionerstage 0-3, offloadgradient accumulationboundary-aware stepop builderFusedAdam, CPUAdam, AIOprebuilt or JITNCCL collectivestorch.distributedPipelineEnginePipelineModule, 1F1BMoE layergate + experts, ep_sizecheckpoint enginesharded save / loadThe engine is a subclass of nn.Module: forward() still runs your model, but backward() and step() now belong to DeepSpeed.
DeepSpeed's pieces. The launcher and config produce an engine per GPU; precision, ZeRO and accumulation live inside it, compiled ops and NCCL sit underneath, and the pipeline, MoE and checkpoint subsystems plug in.

What deepspeed.initialize gives you

deepspeed.initialize takes your module, its parameters and a config, and returns four things: the engine, the optimizer, a data loader if you passed a dataset, and a learning-rate scheduler if the config defined one. The engine is itself a torch.nn.Module, so calling it runs your model's forward. What changes is ownership of everything after the forward.

engine.backward(loss) replaces loss.backward(). It applies the fp16 loss scale, divides the loss by the accumulation steps so gradients average, and reduces or reduce-scatters gradients across ranks at the accumulation boundary. engine.step() replaces optimizer.step(), scheduler.step() and zero_grad() together and only counts until the boundary. Call both on every micro-batch and let the engine decide when real work happens; calling the raw optimizer or loss.backward() bypasses scaling and partitioning and silently corrupts gradients.

import deepspeed

deepspeed.init_distributed()          # NCCL process group from RANK / WORLD_SIZE / MASTER_ADDR
model = build_model()                 # plain torch.nn.Module

engine, optimizer, train_loader, scheduler = deepspeed.initialize(
    model=model,
    model_parameters=[p for p in model.parameters() if p.requires_grad],
    training_data=train_dataset,      # optional: DeepSpeed builds a distributed loader
    config="ds_config.json",
)

for step, batch in enumerate(train_loader):
    batch = {k: v.to(engine.device) for k, v in batch.items()}
    loss = engine(**batch).loss       # forward through the wrapped module
    engine.backward(loss)             # scales loss, accumulates, reduces at the boundary
    engine.step()                     # optimizer + scheduler + zero_grad only at the boundary
    if engine.is_gradient_accumulation_boundary() and engine.global_rank == 0:
        print(f"optimizer step {engine.global_steps}: loss {loss.item():.4f}")

engine.save_checkpoint("ckpt/", tag=f"step{engine.global_steps}")
Advertisement

The config file and the batch-size triangle

The config is where precision, the optimizer, the scheduler, ZeRO and logging are set. The one rule it enforces that trips everyone is the batch triangle: train_batch_size must equal train_micro_batch_size_per_gpu times gradient_accumulation_steps times the number of data-parallel ranks. Set any two and DeepSpeed derives the third; set all three inconsistently and it refuses to start. In the config below, 4 x 8 x 16 GPUs = 512.

The data-parallel rank count is not always the GPU count. With a 4-stage pipeline on 16 GPUs there are only 4 replicas, so the same settings give a global batch of 128, not 512. Recompute the triangle whenever you add pipeline or tensor parallelism.

{
  "train_batch_size": 512,
  "train_micro_batch_size_per_gpu": 4,
  "gradient_accumulation_steps": 8,
  "gradient_clipping": 1.0,
  "bf16": { "enabled": true },
  "optimizer": { "type": "AdamW", "params": { "lr": 3e-4, "weight_decay": 0.1 } },
  "scheduler": { "type": "WarmupLR", "params": { "warmup_num_steps": 2000 } },
  "zero_optimization": { "stage": 2, "overlap_comm": true, "contiguous_gradients": true },
  "steps_per_print": 50,
  "wall_clock_breakdown": false
}

You can pass your own PyTorch optimizer or name one in the config, which lets DeepSpeed use its fused or CPU implementation. Pick one path for the optimizer and scheduler, or two schedulers will step the learning rate.

The launcher, hostfiles and environment

The deepspeed command is a launcher, not a trainer. Across nodes it reads a hostfile, one line per host with a slots count, SSHes to each host and starts one process per GPU with rank, world size and master address in its environment. --num_nodes and --num_gpus limit how much of the hostfile is used; --include and --exclude pick specific hosts and GPUs.

Clusters without passwordless SSH, such as most Kubernetes setups, use --no_ssh: run the launcher on every node with its --node_rank and a shared --master_addr and --master_port. Shell variables do not travel over SSH, so NCCL settings go in ~/.deepspeed_env, which the launcher exports to every rank. deepspeed.init_distributed() also works under torchrun or Slurm when the usual rank variables are set.

# hostfile: one line per node, "slots" = GPUs on that node
node-a slots=8
node-b slots=8

# SSH-based launch from node-a (DeepSpeed ssh-es into every host in the file)
deepspeed --hostfile=hostfile train.py

# no SSH (Kubernetes, Slurm): run this on every node with its own rank
deepspeed --hostfile=hostfile --no_ssh --node_rank=$NODE_RANK \
  --master_addr=node-a --master_port=29500 train.py

# ~/.deepspeed_env on the launching node: KEY=VALUE lines exported to every rank
NCCL_SOCKET_IFNAME=eth0
NCCL_DEBUG=WARN

Custom ops: prebuilt or just-in-time

Several features depend on DeepSpeed's own compiled operators: FusedAdam for GPU optimizer steps, CPU Adam for offloaded optimizer state, and the asynchronous I/O op for NVMe offload. The pip package ships kernel sources; an op not built at install time is compiled just in time when a job first needs it, with whatever CUDA toolkit and compiler the machine has.

JIT compilation is convenient on a workstation and fragile on a cluster: every node compiles on first use, nodes with a different nvcc or no compiler fail, and the toolkit must match the CUDA version PyTorch was built with. Prebuild the ops you use into the container image with the DS_BUILD_* variables and confirm with ds_report, which also prints the torch, CUDA and DeepSpeed versions for bug reports.

# what can this machine build, and what is already built?
ds_report

# prebuild the ops you need at install time instead of JIT-compiling on first use
DS_BUILD_CPU_ADAM=1 DS_BUILD_FUSED_ADAM=1 pip install deepspeed --no-build-isolation

# or try to build every op compatible with this machine (needs nvcc and a C++ compiler)
DS_BUILD_OPS=1 pip install deepspeed --no-build-isolation

ZeRO, briefly

ZeRO shards training state across data-parallel ranks: stage 1 shards optimizer state, stage 2 also gradients, stage 3 also the parameters, gathering each layer just before use. Offload moves state to CPU memory or NVMe. Stages 1 and 2 cost about the traffic of ordinary data parallel; stage 3 adds roughly half again, because parameters are all-gathered in both forward and backward.

The config keys, bucket sizes, zero.Init for models too big to build on one GPU, and export of sharded checkpoints are covered in the ZeRO optimizer in depth; the partitioning design is in ZeRO sharding architecture, and offload trade-offs in optimizer state offloading.

The pipeline engine

Pipeline parallelism splits the layers into stages on different GPUs and streams micro-batches through them. DeepSpeed uses a separate model class, PipelineModule, and engine. You describe the model as a flat list of layers, and DeepSpeed cuts it into num_stages contiguous pieces, balanced by parameter count by default, by layer count with "uniform", or by layers whose class matches a regex with "type:[regex]".

Wrapping layers in LayerSpec defers construction until after partitioning, so each rank allocates only its own stage. TiedLayerSpec shares one module between stages, typically input embeddings reused as the output projection; DeepSpeed all-reduces the tied gradients after backward.

The loop changes shape: you call engine.train_batch(data_iter=...) once per global batch, and the engine pulls the micro-batches it needs and runs the schedule. The pipeline engine asserts that the ZeRO stage is 0 or 1 and rejects stages 2 and 3. For scheduling and the bubble cost see pipeline parallelism architecture.

import torch.nn as nn
from deepspeed.pipe import PipelineModule, LayerSpec, TiedLayerSpec

layers = (
    [TiedLayerSpec("embed", EmbeddingLayer, vocab, d_model)]           # first stage
    + [LayerSpec(TransformerBlock, d_model, n_heads) for _ in range(48)]
    + [LayerSpec(nn.LayerNorm, d_model),
       TiedLayerSpec("embed", EmbeddingLayer, vocab, d_model,
                     forward_fn=lambda layer, x: layer.project_to_vocab(x))]  # last stage
)

net = PipelineModule(layers=layers, num_stages=4,
                     loss_fn=cross_entropy_loss, partition_method="parameters")
engine, _, _, _ = deepspeed.initialize(model=net, model_parameters=net.parameters(),
                                       config="ds_config.json")   # ZeRO stage 0 or 1 only

data_iter = iter(train_loader)       # yields (inputs, labels) micro-batches
for step in range(total_steps):
    loss = engine.train_batch(data_iter=data_iter)   # all micro-batches of one global batch

The MoE layer

A mixture-of-experts layer replaces one feed-forward block with many experts and a gate that routes each token to its top-k. DeepSpeed's MoE module clones an expert module num_experts times and spreads them over ep_size ranks. With 16 experts and ep_size=8 each rank holds two, and tokens move between ranks by all-to-all before and after the experts run.

The forward returns the output, an auxiliary load-balancing loss, and per-expert token counts. Add the auxiliary loss to your training loss, or the gate learns to send everything to a few experts. capacity_factor bounds how many tokens an expert accepts; overflow tokens are dropped by default and pass through only the residual. A persistently idle expert means routing has collapsed.

import torch.nn as nn
from deepspeed.moe.layer import MoE

class MoEBlock(nn.Module):
    def __init__(self, d_model, d_ff, num_experts=16, ep_size=8):
        super().__init__()
        expert = nn.Sequential(nn.Linear(d_model, d_ff), nn.GELU(), nn.Linear(d_ff, d_model))
        self.moe = MoE(hidden_size=d_model, expert=expert, num_experts=num_experts,
                       ep_size=ep_size, k=2, capacity_factor=1.25)

    def forward(self, x):
        out, l_aux, exp_counts = self.moe(x)   # l_aux: load-balancing loss for this layer
        return x + out, l_aux

# in the training step: add every layer's auxiliary loss to the task loss
loss = task_loss + 0.01 * sum(aux_losses)

Worked example: a 13B model on two 8-GPU nodes

Take a 13-billion-parameter decoder, bf16 training with AdamW, on two nodes of eight 80 GB GPUs with fast intra-node links and a slower network between nodes. Assuming 16-bit parameters and gradients and fp32 master weights and Adam moments, the training state is 2 + 2 + 12 = 16 bytes per parameter, or 208 GB in total, before any activations.

Plain data parallel needs all 208 GB on every GPU, so it is out. ZeRO stage 1 over 16 ranks keeps 26 GB of parameters and 26 GB of gradients per GPU plus about 9.8 GB of optimizer shard, about 62 GB, leaving little room for activations. Stage 2 brings it to 26 + 182 / 16, about 37 GB, leaving around 40 GB for activations. Stage 3 reaches 13 GB but all-gathers every layer twice per step, half of it across the slow inter-node link.

The plan is stage 2 with activation checkpointing, a micro-batch that fits in the remaining 40 GB, and accumulation set through the triangle. Move to stage 3 only if longer sequences push activations past the budget; a two-stage pipeline, one stage per node, helps only if the inter-node link is the bottleneck, and forces ZeRO back to stage 1.

DeepSpeed, FSDP, Accelerate or Megatron

OptionStrengthCostChoose it when
DeepSpeedZeRO, offload, pipeline and MoE enginesOwns backward and step; ops to buildYou need offload, its pipeline or MoE engine
PyTorch FSDPSharding native to PyTorchFewer built-in extrasYou want to stay on core PyTorch APIs
AccelerateOne script, several backendsAn extra layer to debugYou switch backends or use Hugging Face trainers
Megatron-style tensor parallelSplits each matrix multiplyModel code written for itSingle layers are too big

They combine: Accelerate can drive DeepSpeed, and large runs use tensor parallelism inside a node with ZeRO across nodes. See PyTorch FSDP in depth and tensor parallelism in depth.

Failure modes

SymptomLikely causeFirst response
Job hangs at init on multi-nodeWrong network interface or blocked portSet NCCL_SOCKET_IFNAME in ~/.deepspeed_env; check master port
Error building extension on some nodesJIT build with mismatched nvcc or no compilerPrebuild ops in the image; compare ds_report across nodes
Batch size assertion at startupTriangle inconsistent with data-parallel sizeSet two of the three values and let DeepSpeed derive the third
Pipeline engine rejects configZeRO stage 2 or 3 with pipelineUse stage 0 or 1 with the pipeline engine
MoE quality poor, some experts idleAuxiliary loss not added, routing collapsedAdd l_aux to the loss; watch exp_counts

What to do next

  1. Run ds_report on every node image and prebuild the ops your config uses.
  2. Write the batch triangle down for your parallel layout and check it against the config.
  3. Convert the training loop to engine.backward and engine.step and remove every direct optimizer call.
  4. Estimate 16 bytes per parameter for mixed-precision Adam, divide by the ZeRO stage's sharding, and pick the lowest stage that leaves room for activations.
  5. For multi-node runs, put NCCL settings in ~/.deepspeed_env or use --no_ssh with explicit ranks.
  6. Add the pipeline or MoE engine only after a plain ZeRO run works, and re-derive the batch size when you do.
Key takeaway: DeepSpeed is an engine that takes over backward, step, precision and memory partitioning from PyTorch, driven by one config, plus a launcher, compiled ops, and optional pipeline and MoE subsystems. Most production problems come from its edges: the batch-size triangle, JIT-compiled ops that differ between nodes, environment that does not travel over SSH, and the pipeline engine's ZeRO stage 1 limit. Build the ops into the image, let the engine own the optimizer, size memory before choosing a ZeRO stage, and add parallelism one dimension at a time.