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.
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 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}")
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
| Option | Strength | Cost | Choose it when |
|---|---|---|---|
| DeepSpeed | ZeRO, offload, pipeline and MoE engines | Owns backward and step; ops to build | You need offload, its pipeline or MoE engine |
| PyTorch FSDP | Sharding native to PyTorch | Fewer built-in extras | You want to stay on core PyTorch APIs |
| Accelerate | One script, several backends | An extra layer to debug | You switch backends or use Hugging Face trainers |
| Megatron-style tensor parallel | Splits each matrix multiply | Model code written for it | Single 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
| Symptom | Likely cause | First response |
|---|---|---|
| Job hangs at init on multi-node | Wrong network interface or blocked port | Set NCCL_SOCKET_IFNAME in ~/.deepspeed_env; check master port |
| Error building extension on some nodes | JIT build with mismatched nvcc or no compiler | Prebuild ops in the image; compare ds_report across nodes |
| Batch size assertion at startup | Triangle inconsistent with data-parallel size | Set two of the three values and let DeepSpeed derive the third |
| Pipeline engine rejects config | ZeRO stage 2 or 3 with pipeline | Use stage 0 or 1 with the pipeline engine |
| MoE quality poor, some experts idle | Auxiliary loss not added, routing collapsed | Add l_aux to the loss; watch exp_counts |
What to do next
- Run
ds_reporton every node image and prebuild the ops your config uses. - Write the batch triangle down for your parallel layout and check it against the config.
- Convert the training loop to
engine.backwardandengine.stepand remove every direct optimizer call. - 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.
- For multi-node runs, put NCCL settings in
~/.deepspeed_envor use--no_sshwith explicit ranks. - Add the pipeline or MoE engine only after a plain ZeRO run works, and re-derive the batch size when you do.