Why Model Sharding Is Necessary at 70B+ Parameters
A 70B parameter model at BF16 precision requires 140GB of GPU memory just for the model weights. An additional 140GB is needed for optimizer states (Adam optimizer stores two momentum terms per parameter), and roughly 140GB for gradients during the backward pass. Total memory requirement: 420GB for a standard training configuration. The H100 SXM5 provides 80GB of HBM per GPU. The math is simple: no single GPU can hold a 70B model for training. Model sharding - distributing the model state across multiple GPUs - is not an optimization choice; it is a requirement for any model above approximately 7B parameters at BF16.
The landscape of model sharding strategies has evolved rapidly over the past three years. The first generation was naive data parallelism (DDP), where each GPU held a complete copy of the model and processed different data batches. DDP works up to roughly 7B parameters on H100 but becomes infeasible for larger models because each GPU's 80GB HBM cannot hold the full model state. The second generation was ZeRO (Zero Redundancy Optimizer), which sharded optimizer states, gradients, and parameters across GPUs, eliminating the redundant copies that DDP maintained. The third generation, currently dominant in 2026, is FSDP (Fully Sharded Data Parallelism) and DeepSpeed ZeRO Stage 3, which shard all three model state components and have largely converged in capability.
The choice between sharding strategies depends on model size, GPU count, and interconnect bandwidth. For models up to 30B parameters on 8-16 GPUs with NVLink interconnect, FSDP with ZeRO-2 (sharding optimizer states and gradients but not parameters) provides the best balance of memory savings and communication efficiency. For models above 30B on 16+ GPUs, FSDP ZeRO-3 (sharding all three: optimizer states, gradients, and parameters) is necessary, with the tradeoff of increased all-gather communication for parameter unsharding during the forward and backward passes.
FSDP vs DeepSpeed ZeRO: Which Sharding Framework to Use
FSDP (Fully Sharded Data Parallelism) is PyTorch's native implementation of the ZeRO-3 algorithm. It is integrated directly into the PyTorch distributed package and supports the widest ecosystem of model architectures without requiring model code changes. The key advantage: FSDP works with any PyTorch model by wrapping layers or submodules with the FSDP wrapper, which automatically handles sharding, unsharding, and gradient synchronization. FSDP's auto-wrap policy can recursively shard submodules, optimizing the tradeoff between communication efficiency and memory savings. FSDP is the default choice for teams building custom training pipelines in PyTorch.
DeepSpeed ZeRO is Microsoft's distributed training framework that implements all three ZeRO stages plus additional memory optimizations not available in FSDP. The key differentiators: ZeRO-Infinity extends ZeRO to offload parameters and optimizer states to CPU or NVMe storage when GPU memory is exhausted, enabling training of models up to 1 trillion parameters on a single DGX node. DeepSpeed's ZeRO-offload feature moves optimizer states and gradients to CPU memory, reducing GPU memory usage by 50-70% at the cost of increased training step time. DeepSpeed also includes the ZeRO-3 throughput optimizer that overlaps communication with computation to hide all-gather latency.
The framework decision impacts training throughput by 5-15% depending on model architecture and cluster configuration. Benchmark data as of mid-2026 shows FSDP and DeepSpeed ZeRO-3 within 3% of each other in throughput for standard Transformer architectures on 64 H100 GPUs with NVLink interconnect. The gap widens for non-standard architectures (MoE, Mamba, RWKV) where DeepSpeed's specialized kernel implementations provide more consistent performance. The practical recommendation: use FSDP for standard Transformer training where PyTorch integration and ecosystem compatibility are priorities. Use DeepSpeed for models above 100B parameters where ZeRO-Infinity's CPU/NVMe offload is needed, or for MoE models where DeepSpeed's MoE kernel optimizations provide measurable throughput gains.
Tensor Parallelism and Pipeline Parallelism: When to Layer on Top
FSDP and ZeRO are data-parallel strategies - they shard model state but each GPU processes different data. Tensor parallelism (TP) and pipeline parallelism (PP) are model-parallel strategies - they shard the model architecture itself, with GPUs processing different parts of the model for the same data. TP splits each layer's operations across GPUs. For a Transformer layer's self-attention, TP splits the QKV projection and attention computation across GPUs, reducing per-GPU memory for activation storage at the cost of all-to-all communication at each layer boundary.
When to add TP on top of FSDP: when the model's memory per layer exceeds a single GPU's HBM for a single batch, even with ZeRO-3 sharding. This occurs at around 200B+ parameters for Transformers with standard layer sizes. TP reduces the per-GPU memory for activations and parameters further by splitting each layer across GPUs. The communication cost of TP is significant: each layer's all-to-all communication adds roughly 2x the layer's activation size in data transfer. TP-8 (splitting across 8 GPUs) on H100 with NVLink adds 5-15% to per-step time versus no TP but enables training of models that would not fit otherwise.
Pipeline parallelism (PP) splits the model vertically by layers, with each GPU owning a contiguous set of layers. PP reduces memory by limiting the number of layers per GPU but introduces a pipeline bubble idle time. For a 4-stage PP configuration, the pipeline bubble is roughly 25% of training throughput. PP is most useful when TP cannot scale further due to communication overhead, which typically happens at 8 GPUs of TP (TP-8 is the practical maximum for most Transformer architectures). The common hybrid approach for 100B+ models: TP within a node (NVLink-connected GPUs) with PP between nodes (InfiniBand-connected), combined with FSDP across nodes for additional data parallelism.
Communication Tuning: NCCL Configuration for Sharded Training
NCCL configuration is the most impactful performance variable for sharded training. The default NCCL settings prioritize compatibility over performance, leaving 20-40% throughput on the table for well-tuned configurations. The critical NCCL environment variables: NCCL_ALGO selects the communication algorithm (Ring, Tree, or CollnetDirect). Ring is best for sharded training with uniform message sizes because it achieves the highest bandwidth utilization for long-duration transfers. Tree is better for latency-sensitive operations with small message sizes like gradient all-reduce.
NCCL_PROTO selects the transport protocol (Simple, LL, LL128). Simple is the standard protocol and works for all message sizes. LL (Low Latency) uses GPU shared memory for small messages under 256KB and reduces latency by 30-50% for gradient reduction shards. LL128 is optimized for 128-byte aligned messages and provides the best performance for message sizes between 256 bytes and 8KB. For ZeRO-3 training where per-layer parameter shards range from 1MB to 100MB, NCCL_PROTO=Simple with NCCL_ALGO=Ring provides the best overall throughput.
NCCL_MIN_NCHANNELS controls the number of communication channels used for parallel transfers. The default of 1 underutilizes NVLink fabric. Setting NCCL_MIN_NCHANNELS to match the number of NVLink links per GPU (18 for H100 SXM5) maximizes bandwidth utilization for intra-node communication. For inter-node communication over InfiniBand, NCCL_NET_GDR_LEVEL should be set to 3 to enable GPU Direct RDMA, bypassing CPU memory for data transfers between GPUs across nodes. GPU Direct RDMA reduces inter-node communication latency by 30-50% and is essential for ZeRO-3 training across multiple nodes.
Memory vs Performance Tradeoffs: Gradient Checkpointing, Activation Recomputation
Gradient checkpointing (activation recomputation) trades compute for memory by recomputing intermediate activations during the backward pass instead of storing them from the forward pass. The tradeoff: 15-25% additional compute time for 50-70% reduction in activation memory. For FSDP ZeRO-3, activation memory is sharded across GPUs, so the absolute memory savings from gradient checkpointing are reduced but still meaningful. On a 64-GPU cluster training a 70B model, gradient checkpointing reduces per-GPU activation memory from roughly 4GB to 1.5GB, enabling larger batch sizes or reducing the need for ZeRO-3 parameter sharding.
Mixed precision training (BF16 compute, FP32 master weights) is standard practice that provides a free memory reduction: the model's forward and backward passes execute at BF16 (half the memory of FP32), while optimizer states are maintained in FP32 for numerical stability. The total memory without mixed precision at FP32 for a 70B model is 420GB (weights + optimizer + gradients). With BF16 mixed precision, the total drops to 210GB (140GB at BF16 for weights and gradients, plus 70GB for FP32 optimizer states). This is the single most impactful memory-saving configuration and is enabled by default in modern training frameworks.
CPU and NVMe offload from DeepSpeed ZeRO-Infinity pushes memory savings to the extreme by moving optimizer states and gradients to CPU DRAM or NVMe SSD when GPU memory is exhausted. The configuration: set offload_optimizer and offload_params in the DeepSpeed config to cpu or nvme. The performance penalty is significant: CPU offload adds 20-40% to training step time due to PCIe transfer latency. NVMe offload adds 50-80% due to SSD access latency. Offload should be used only when GPU and CPU memory are both exhausted, typically for models above 200B parameters on small GPU clusters. For most 70B-100B training runs on modern GPU clusters with 8-16 GPUs, offload is unnecessary and counterproductive.
Cluster Sizing for Sharded Training: GPU Count and Interconnect Requirements
The GPU cluster size for sharded training depends on the total model state memory and the per-GPU HBM capacity. The formula: minimum GPU count = total_model_state_memory / (per_GPU_HBM - activation_memory_per_GPU - KV_cache_memory). For a 70B model at BF16 mixed precision with standard sequence lengths, the minimum is 4-8 H100 GPUs (ZeRO-3 with gradient checkpointing). For a 400B MoE model, the minimum is 16-32 H100 GPUs. Adding more GPUs beyond the minimum does not linearly increase throughput due to communication overhead scaling.
Interconnect requirements for sharded training scale with the sharding strategy. FSDP ZeRO-3 requires high-bandwidth GPU-GPU communication for the all-gather operations per layer. Intra-node: NVLink 4 at 900 GB/s per GPU is sufficient for up to 8-GPU ZeRO-3 training. Inter-node: InfiniBand NDR400 at 400 Gbps per node is the minimum for multi-node ZeRO-3; solutions below this bandwidth (200 Gbps Ethernet) will become communication-bound for models above 30B parameters. For tensor parallelism, intra-node NVLink is required - TP over InfiniBand adds prohibitive latency.
The cluster guidance for 2026 GPU purchases: for teams training models up to 70B, 8x H100 SXM5 nodes with NVLink and InfiniBand interconnect provide the best cost-performance. 4 nodes (32 GPUs) handle most foundation model training runs with comfortable margins. For teams training models above 70B, B200 nodes with NVLink 5 fabric memory provide architectural advantages that reduce the sharding complexity - single-node training of models up to 400B becomes feasible with the 192GB per GPU on B200. The GPU choice should consider the sharding strategy: H100 requires aggressive sharding for large models, while B200 reduces sharding overhead. The cost comparison must factor in the sharding efficiency difference, not just the per-GPU hourly rate.
