Data parallelism is a parallel computing strategy in which the same operation is applied simultaneously to different partitions of a dataset distributed across multiple processing units. In machine learning it is the dominant approach to scaling model training: the model is replicated on every worker, each worker computes gradients on a distinct mini-batch shard, and gradients are aggregated through collective communication such as all-reduce before the synchronised parameter update. This contrasts with model parallelism, which splits a single model across devices, and underpins large-scale training on GPU and TPU clusters.

Semantic Classification

Content

Compositional Relationships (Components)

SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:hasPart ai:GradientAggregation))
SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:hasPart ai:AllReduce))
SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:hasPart ai:MiniBatch))
SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:hasPart ai:RingAllReduce))
SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:hasPart ai:GradientCompression))
SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:hasPart ai:ModelReplica))
SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:hasPart ai:ParameterUpdate))

Dependency Relationships

SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:requires ai:CollectiveCommunication))
SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:requires ai:GPUCluster))
SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:requires ai:Synchronisation))
SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:requires ai:NetworkBandwidth))
SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:requires ai:NCCL))
SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:dependsOn ai:Backpropagation))
SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:dependsOn ai:LossFunction))

Capability Relationships

SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:enables ai:LargeScaleTraining))
SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:enables ai:ThroughputScaling))
SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:enables ai:LLMPretraining))
SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:supports ai:DeepLearning))
SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:supports ai:TransformerArchitecture))
SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:supports ai:NeuralNetworkTraining))

Implementation Relationships

SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:implements ai:StochasticGradientDescent))
SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:implements ai:GradientDescent))
SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:uses ai:ParameterServer))
SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:uses ai:MixedPrecisionTraining))
SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:uses ai:Checkpoint))
SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:uses ai:GradientCheckpointing))

Reduction Relationships

SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:reducesTo ai:ParallelProcessing))
SubClassOf(ai:SynchronousDataParallelism
  ObjectSomeValuesFrom(ai:reducesTo ai:DataParallelism))
SubClassOf(ai:FullyShardedDataParallel
  ObjectSomeValuesFrom(ai:reducesTo ai:DataParallelism))
SubClassOf(ai:AsynchronousParameterServer
  ObjectSomeValuesFrom(ai:reducesTo ai:DataParallelism))
SubClassOf(ai:DiLoCo
  ObjectSomeValuesFrom(ai:reducesTo ai:DataParallelism))
SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:reducesTo ai:DistributedComputing))
SubClassOf(ai:FederatedLearning
  ObjectSomeValuesFrom(ai:reducesTo ai:DataParallelism))
SubClassOf(ai:ZeROOptimizer
  ObjectSomeValuesFrom(ai:reducesTo ai:DataParallelism))
SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:enables ai:FaultTolerance))
SubClassOf(ai:DataParallelism
  ObjectSomeValuesFrom(ai:requires ai:Backpropagation))

About

  • Data Parallelism is the most widely deployed strategy for scaling the training of Deep Learning models beyond the capacity of a single compute device. The core idea dates to early work on shared-memory parallel computing but was operationalised for neural network training primarily through the adoption of GPU clusters and high-bandwidth interconnects in the 2010s. In data-parallel training, the model is replicated in full on every worker node; a global mini-batch is partitioned across all workers so that each processes a non-overlapping shard; the forward and backward passes execute independently and simultaneously on each replica; and gradients are synchronised via an all-reduce collective operation before every parameter update. Because each worker applies the same update to its local replica, all replicas remain mathematically identical throughout training. This equivalence to single-device training at the aggregate batch size is the defining mathematical property that makes data parallelism so attractive: it inherits all the convergence guarantees of mini-batch Stochastic Gradient Descent while delivering near-linear throughput scaling as long as communication overhead remains small relative to computation.
  • The central performance constraint of synchronous data parallelism is the communication-to-computation ratio. For small models or very large batches, compute dominates and data parallelism scales near-linearly. For large models with small per-GPU batch sizes — as occurs when fitting a 70B-parameter model on 80 GB A100s — the ratio reverses: each all-reduce step aggregates a gradient vector containing billions of float values, and even at 400 Gbps InfiniBand bandwidth this communication takes hundreds of milliseconds, comparable to or exceeding the compute time per step. This motivates gradient compression (top-k sparsification, PowerSGD low-rank approximation, BF16 quantisation), compute-communication overlap (DDP bucket scheduling, FSDP prefetch), and hybrid parallelism strategies that reduce the number of ranks participating in each all-reduce.
  • The synchronous variant — embodied in PyTorch’s DistributedDataParallel (DDP) and JAX’s pmap — guarantees that gradients are always fully averaged before each step, producing results statistically equivalent to single-device training at the aggregate batch size. However, synchronous all-reduce introduces a communication round-trip latency at every step, making the approach sensitive to stragglers (slow workers) and requiring high-throughput, low-latency interconnects such as NVLink (within nodes) and InfiniBand or RoCE (across nodes). The asynchronous Parameter Server paradigm relaxes this constraint by allowing workers to push and pull gradients without global synchronisation, improving throughput at the cost of gradient staleness and potentially inconsistent updates. The trade-off between consistency and throughput has been extensively studied: asynchronous methods can tolerate staleness of up to τ steps while retaining convergence under appropriate learning-rate schedules, but they require careful engineering of the parameter server to avoid lock contention and communication hotspots.
  • Memory pressure is the primary limitation of classical data parallelism: the full model, optimiser state, and activation buffers must reside on each worker’s device. For models with billions of parameters this exceeds the capacity of individual GPUs. Consider a 7-billion-parameter model in BF16: parameters alone require 14 GB; Adam optimiser states (first and second moment vectors in FP32) require a further 56 GB; gradients require another 14 GB — totalling 84 GB, already beyond the 80 GB A100 capacity, and before accounting for activations. Sharded data parallelism — realised through Microsoft’s ZeRO Optimizer in DeepSpeed (ZeRO Stages 1, 2, and 3) and PyTorch’s Fully Sharded Data Parallel (FSDP, with FSDP2 introduced in PyTorch 2.4 in 2024) — partitions these memory-resident tensors across workers, reducing per-GPU memory by a factor of N while preserving the data-parallel programming model. ZeRO-3 shards all parameter, gradient, and optimiser-state tensors: each rank holds only 1/N of each tensor, issuing all-gather communications before each forward and backward computation unit and reduce-scatter after each backward unit. FSDP2 (PyTorch ≥ 2.4, 2024) refactored this logic using the DTensor abstraction for cleaner composability with Tensor Parallelism and improved prefetch scheduling that hides the all-gather latency behind the previous layer’s backward pass. As of 2025–2026, the state-of-the-art production infrastructure (Megatron-LM, DeepSpeed, Llama training stacks, xAI Colossus) combines data parallelism with Tensor Parallelism and Pipeline Parallelism in a three-dimensional parallelism hierarchy that exploits NVLink bandwidth for tight tensor-parallel communication within nodes, high-bandwidth interconnects for pipeline-parallel stages, and data-parallel replicas across node groups.
  • The scalability of data parallelism is governed by Amdahl’s Law and its stochastic generalisation. The communication fraction f = t_comm / (t_comm + t_comp) determines the maximum speedup: lim_{N→∞} 1 / f as N → ∞. In practice, for transformer architectures at scale, compute efficiency of 35–55% of theoretical peak FLOPs is considered excellent; the Megatron-LM infrastructure achieves ~38% hardware FLOPs utilisation on A100 clusters with 3D parallelism. The linear scaling rule for batch size (Goyal et al., 2017) establishes that with N-fold data parallelism, the learning rate should be scaled by N with a linear warm-up phase to maintain equivalent convergence; this rule breaks down at very large batch sizes (> 32k–64k for ImageNet) where gradient noise is insufficient for exploration and convergence degrades.

Formal Analysis

  • Gradient averaging equivalence: Let w_t denote the model parameters at step t, and let g_i^t = ∇ L_i(w_t) denote the gradient computed on shard i. Synchronous data-parallel SGD updates via w_{t+1} = w_t − η · (1/N) Σ_{i=1}^N g_i^t. This is identical to mini-batch SGD on the full batch ∪_i shard_i, with the same learning rate η, as long as the loss decomposes as L(w) = (1/N) Σ_i L_i(w). The gradient variance seen by each update is reduced by a factor of N relative to single-sample SGD, which is why large-batch training requires a larger learning rate (or equivalently, fewer noisy steps) to explore the loss landscape similarly.
  • Convergence under staleness (asynchronous case): For a τ-bounded-delay parameter server model, convergence to a stationary point of non-convex L at rate O(1/√T) is recoverable under standard smoothness and bounded-gradient assumptions, but with a constant factor degradation proportional to τ. In practice, with τ ≤ 8–16 on large clusters, asynchronous SGD often matches synchronous quality on tasks where gradient correlation across mini-batches is low (large heterogeneous datasets).
  • Communication complexity: Ring all-reduce over N workers for a gradient vector of size G bytes achieves communication volume 2G(N-1)/N bytes per worker (converging to 2G as N → ∞). The ring algorithm achieves this in 2(N-1) communication steps, with each step transferring G/N bytes per link, yielding step time G/(N · B) where B is link bandwidth. The total time is thus 2G/B, independent of N — demonstrating that ring all-reduce is bandwidth-optimal and enables near-linear throughput scaling as N increases (communication time is constant while compute time per step halves with each doubling of N).
  • ZeRO memory analysis: Classical DDP requires 16Ψ bytes per GPU (Ψ parameters: 2Ψ BF16 params + 2Ψ BF16 grads + 4Ψ FP32 params + 8Ψ FP32 Adam states). ZeRO-3 reduces this to 16Ψ/N + activation memory, enabling training of arbitrarily large models across sufficiently many GPUs. For N=64 ranks with a 7B-parameter model, per-GPU memory drops from ~112 GB to ~1.75 GB for model states, with activation memory dominating.

Components and Architecture

  • Model Replica: Each worker holds a complete, synchronised copy of the model parameters (or a shard in FSDP/ZeRO configurations). In classical DDP, each GPU stores a full copy: for a 7B-parameter model in BF16 this occupies 14 GB. FSDP2 and ZeRO-3 replace this with partial shards, requiring all-gather operations before each computation unit to reconstruct the full parameter tensor locally.
  • Data Shard: The global mini-batch is divided into equal-size sub-batches, one per worker, ensuring each training sample is seen exactly once per epoch. Sub-batch size = global batch size / number of workers. With 1,024 workers and a global batch size of 4M tokens, each worker processes approximately 4,096 tokens per step. Shuffling and seeding strategies must ensure that the same samples are not assigned to the same worker repeatedly across epochs.
  • Forward Pass: Each worker computes a forward pass on its local shard, producing intermediate activations at every layer and a per-shard scalar loss. Activations must be stored in memory for use during the backward pass (or recomputed via gradient checkpointing at the cost of 30–40% extra compute). In FP16/BF16 mixed-precision training, activations are stored in half-precision while the master parameters remain in FP32.
  • Backward Pass / Backpropagation: Gradients with respect to all parameters are computed locally on each worker via automatic differentiation (autograd). The backward pass traverses the computational graph in reverse, accumulating gradients via the chain rule. In DDP, bucket-wise all-reduce is initiated as soon as all gradients within a bucket are ready, overlapping communication with the remaining backward computation for subsequent layers.
  • Gradient Aggregation via All-Reduce: Gradients are averaged across all workers using a collective communication primitive. The dominant algorithm is Ring All-Reduce (via NCCL), which divides the gradient tensor into N chunks and performs a reduce-scatter (each rank receives 1/N of the fully reduced gradient) followed by an all-gather (each rank broadcasts its shard to all others), achieving near-bandwidth-optimal communication with total data transfer of 2G(N-1)/N bytes per worker — converging to 2G as N grows. As of 2025, NCCL 2.23 introduced the PAT (Parallel-All-Transfer) algorithm for deployments exceeding 100,000 GPUs, where the ring algorithm’s O(N) step count becomes a latency bottleneck.
  • Synchronised Parameter Update: All workers apply the averaged gradient to update their local replica via an Optimiser (Adam, SGD with momentum, Adafactor, etc.), remaining perfectly synchronised. The update rule for Adam: m_t = β₁ m_{t-1} + (1−β₁) ḡ_t; v_t = β₂ v_{t-1} + (1−β₂) ḡ_t²; θ_{t+1} = θ_t − η · m̂_t / (√v̂_t + ε), where ḡ_t is the all-reduced gradient. All state vectors (m, v, θ) are synchronised after each step, ensuring all replicas are identical.
  • Gradient Compression: Optional techniques reduce communication volume at the cost of approximation error. Top-k sparsification transmits only the k largest-magnitude gradient entries (typically k = 0.1% of total parameters) with coordinates, achieving 1000-fold compression with less than 2% accuracy degradation on image classification tasks. PowerSGD (Vogels et al., 2019) uses low-rank matrix approximation of gradient tensors, achieving high compression ratios with controllable error. 1-bit Adam (Tang et al., 2021) quantises gradients to a single bit with variance correction, achieving 5-fold speedup over baseline on GPT-2 training.
  • Mixed Precision Training: FP16 or BF16 computation with FP32 master weights (used in FSDP/ZeRO optimiser states) reduces memory and bandwidth requirements by 2-fold relative to FP32 training. BF16 (bfloat16) is preferred over FP16 for large-scale LLM training because its wider dynamic range (same as FP32) avoids the gradient underflow and overflow issues that require loss scaling in FP16 training. Mixed precision halves the memory for activations and parameter copies, halves the bandwidth for all-reduce communication (since gradients are communicated in FP16/BF16), and approximately doubles hardware throughput on Tensor Cores and Matrix Multiply units.
  • Checkpoint and fault tolerance: At the scale of thousands of GPUs running for weeks, hardware failures (GPU memory errors, network link failures, power outages) are inevitable. Asynchronous checkpointing — saving model and optimiser state to distributed storage (AWS S3, Google Cloud Storage, HDFS) without stalling training — is essential. PyTorch Distributed Checkpoint (DCP) saves sharded FSDP/ZeRO state efficiently by having each rank save its local shard without requiring all-gather. Elastic training frameworks (PyTorch Elastic, Ray) enable adding or removing workers mid-run, re-sharding the data-parallel groups dynamically.
  • Topology-aware collective scheduling: The mapping of data-parallel ranks to physical network topology (NVLink domains, InfiniBand fat-tree switches, dragonfly groups) determines all-reduce performance. Topology-aware collective libraries (NCCL’s TOPO_DETECT, RCCL for AMD GPUs) automatically construct ring and tree topologies that minimise cross-switch traffic by routing intra-node communication over NVLink and inter-node communication over InfiniBand, achieving 3–5× better bandwidth utilisation than topology-oblivious round-robin rank assignment.

Parallelism Variants and Major Families

  • Synchronous DDP (PyTorch DistributedDataParallel): The standard baseline for data parallelism when the model fits within a single GPU’s memory. Hooks into the PyTorch autograd engine to register gradient-ready hooks that initiate bucket-wise all-reduce operations as gradients accumulate during the backward pass, overlapping communication with subsequent backward computation steps. Achieves near-perfect utilisation when the model is large enough that each backward layer’s compute duration exceeds the time to communicate its gradient bucket. The all-reduce is synchronous — all ranks block at the barrier until the collective completes — ensuring mathematical equivalence to single-device training. Limitations: (1) requires the full model to fit in a single GPU’s memory (typically constraining models to ≤7B parameters on 80 GB A100s); (2) sensitive to straggler workers since all ranks wait for the slowest; (3) the communication and computation overlap effectiveness degrades for very small models where compute times are too short to hide communication latency.
  • Fully Sharded Data Parallel (FSDP / FSDP2): PyTorch’s production implementation of ZeRO-3-style parameter, gradient, and optimiser-state sharding. FSDP wraps model units (transformer blocks, attention modules, MLP layers) and for each unit: issues an all-gather to reconstruct full parameters before the forward pass; executes the forward pass; discards the gathered parameters to free memory; re-gathers before the backward pass; executes the backward pass; and issues a reduce-scatter of the computed gradients to distribute shards back to owning ranks. FSDP2 (PyTorch ≥ 2.4, released August 2024) replaced the wrapping-based API with DTensor-based sharding semantics, exposing a device_mesh abstraction that cleanly composes data-parallel and tensor-parallel dimensions without conflicting collective scheduling. FSDP2 improved prefetch scheduling — initiating the all-gather for the next module while the current module’s computation is in progress — achieving better communication-computation overlap and 5–15% throughput improvements over FSDP1 on standard LLM benchmarks.
  • ZeRO Optimizer (DeepSpeed ZeRO-1/2/3): Microsoft DeepSpeed’s Zero Redundancy Optimizer series. ZeRO-1 shards only the optimiser states (e.g., Adam’s first and second moment vectors) across N ranks, reducing optimiser memory by N-fold while keeping parameters and gradients replicated. ZeRO-2 additionally shards gradients after the reduce-scatter, eliminating gradient redundancy. ZeRO-3 shards all three: parameters, gradients, and optimiser states, requiring all-gather before forward/backward and reduce-scatter after backward for each parameter group. ZeRO-Infinity extends offloading of each tier to CPU memory or NVMe SSD storage via an efficient asynchronous prefetch pipeline, enabling billion-parameter training on a single GPU. ZeRO++ (2023) added hierarchical all-reduce and quantised weights to reduce cross-node communication volume by 4-fold.
  • Parameter Server (Asynchronous SGD): The predecessor paradigm in which a set of parameter server nodes hold the canonical model parameters and worker nodes pull parameters, compute gradients on local data, and push gradient updates asynchronously. Workers do not wait for other workers’ updates before pushing their own; this eliminates the all-reduce barrier and makes the approach robust to straggler workers. The trade-off is gradient staleness: a gradient computed on parameters k steps old introduces a bias proportional to the learning rate and staleness k. DistBelief (Dean et al., 2012) used this architecture at Google. In 2026, parameter servers remain relevant for recommendation systems with very large, sparse embedding tables where the embedding gradient is sparse and asynchronous updates have limited staleness relative to the embedding diversity.
  • Sharded Data Parallel / Hybrid topology-aware: Combines parameter/gradient sharding with hierarchical topology awareness. A common pattern is intra-node FSDP (sharding across the 8 GPUs within a node using NVLink, which runs at 900 GB/s bidirectional) combined with inter-node DDP-style all-reduce (over InfiniBand at 400 Gbps). This exploits the order-of-magnitude bandwidth differential between intra-node and inter-node interconnects: the expensive full parameter all-gather happens within the node on fast NVLink, while the cheaper reduce-scatter of gradient updates travels across slower inter-node InfiniBand. This hybrid approach achieves the memory savings of FSDP with significantly lower cross-node communication overhead than naive FSDP across all ranks.
  • DiLoCo (Distributed Low-Communication training): A research paradigm from Google DeepMind (Douillard et al., 2023) that operates two-level optimisation: an inner loop of H (e.g., 500) standard SGD steps on each worker’s local data shard, followed by an outer synchronisation that averages or aggregates worker model states using a powerful outer optimiser (Nesterov momentum). The inner loop accumulates H steps worth of gradient information before a single cross-worker communication round, reducing communication frequency by H-fold. Experiments showed that H=500 inner steps with cosine-annealed inner learning rate achieves within 1–2% of fully synchronised DDP on LM benchmarks while reducing communication volume by 500×. This makes DiLoCo suitable for geographically distributed training across data centres connected by commodity internet links (latency: 50–200 ms) rather than high-performance InfiniBand (latency: 1–5 μs).
  • Federated Learning: A privacy-preserving distributed learning paradigm where data-parallel training occurs on client devices (hospitals, phones, edge servers) that cannot share raw data. Each round, a subset of clients trains locally on their private data and sends only model updates (gradients or weight deltas) to a central aggregator. FedAvg (McMahan et al., 2017) averages client updates proportional to dataset size. Secure aggregation protocols (using secret sharing or homomorphic encryption) prevent the server from inspecting individual client updates. Differential Privacy (DP-FedAvg) adds Gaussian noise to gradients before aggregation to provide formal privacy guarantees. Federated Learning trades communication efficiency and privacy for challenges of statistical heterogeneity (non-IID data across clients) and system heterogeneity (variable client compute and connectivity).
  • 3D Parallelism: The production standard for frontier LLM training, combining data parallelism (D), Tensor Parallelism (T), and Pipeline Parallelism (P) in a three-dimensional grid where each dimension maps to a distinct communication topology. In Megatron-LM’s canonical configuration, T workers are connected via NVLink (within node) for all-reduce over attention head and MLP column-row splits; P stages communicate activations and gradients via point-to-point GPU Direct RDMA across InfiniBand; D replica groups use ring all-reduce across InfiniBand for gradient averaging. The optimal (D, T, P) configuration is model- and hardware-dependent: for GPT-3 (175B parameters) on 1,024 A100s, Megatron-LM uses T=8, P=16, D=8. Expert Parallelism (EP) adds a fourth dimension for mixture-of-experts models, routing each token to one of E experts distributed across EP ranks via all-to-all collectives.

Use Cases

  • Large Language Models pre-training: GPT-4, Llama 3, Gemini, Claude, Mistral, Qwen, and equivalent frontier models are trained using data-parallel replicas distributed across tens of thousands of GPUs over weeks to months of continuous training. FSDP2 or ZeRO-3 sharding is mandatory to eliminate per-GPU memory constraints at 7B–405B parameter scales. Effective batch sizes of 4M–16M tokens per step are achieved by combining data parallelism (thousands of replicas) with gradient accumulation. The data-parallel all-reduce at these scales aggregates 14–800 GB of gradient data per step, motivating BF16 gradient precision and hierarchical all-reduce (intra-node ring + inter-node tree) to manage bandwidth.
  • Convolutional Neural Network image classification and vision models: The classic use case that originally motivated synchronous DDP. ResNet-50 training on ImageNet with 256 GPUs in one hour (Goyal et al., 2017) demonstrated the viability of the approach. Today, ViT (Vision Transformer) and EfficientNet training at the scale of JFT-3B (Google’s 3-billion-image internal dataset) or LAION-5B requires thousands of TPU or GPU devices with data-parallel batch sizes in the hundreds of thousands of images.
  • Recommendation and ranking systems: Production-scale recommender systems at Meta (DLRM), Google (YouTube), and TikTok process hundreds of billions of user interactions. Embedding tables (representing items and users) are typically too large for a single device and are sharded across many servers using embedding parallelism (a form of model parallelism), while the dense MLP layers that process embedding lookups are replicated in a data-parallel fashion. The gradient synchronisation for the dense layers uses all-reduce while embedding gradients use all-to-all or reduce-scatter collectives.
  • Scientific computing, climate modelling, and drug discovery: Distributed ML workflows for physics-informed neural networks exploit data parallelism over large spatial domain decompositions. GraphCast (Google DeepMind, 2023) and FourCastNet (NVIDIA, 2022) are neural weather prediction models trained data-parallelly on ERA5 meteorological data. AlphaFold 2 (DeepMind) and AlphaFold 3 used data-parallel training across TPU pods to train on PDB and protein sequence databases. Molecular dynamics foundation models (ESM-2, ESM-3, EvolutionaryScale) are trained with data parallelism over protein sequence databases containing hundreds of millions of sequences.
  • Reinforcement learning at scale: Population-based training (PBT) and massively parallel environment rollouts (as used in OpenAI Five for Dota 2, AlphaStar, and MuZero) rely on data parallelism where parallel environment actors produce experience trajectories that are aggregated into a shared replay buffer. Policy gradient updates are computed data-parallelly over experience batches. The actor-critic separation in IMPALA (Espeholt et al., 2018) explicitly exploits data parallelism for the critic update while allowing asynchronous actor execution.
  • Foundation model fine-tuning and instruction tuning: Supervised fine-tuning (SFT) and RLHF (Reinforcement Learning from Human Feedback) of large pre-trained models uses data-parallel training over instruction datasets. DeepSpeed’s ZeRO-3 with LoRA (Low-Rank Adaptation) sharding enables fine-tuning of 70B-parameter models on 8 × A100 GPUs, making frontier model customisation accessible to academic research groups.
  • Continual and lifelong learning: Incrementally updating models on streaming data streams requires maintaining data-parallel training infrastructure that can efficiently process new data shards while preserving previously learned representations, with replay buffers for catastrophic forgetting mitigation distributed data-parallelly.

Academic Context

  • The theoretical foundations of data-parallel gradient descent trace to Robbins and Monro (1951) on stochastic approximation and to Bottou (1998) on online learning and the dynamics of SGD. The first explicit demonstration of distributed gradient averaging for neural network training appeared in Seide et al. (2014) with 1-bit SGD on speech models, which showed that aggressive gradient quantisation could reduce communication volume 64-fold with minimal accuracy loss. Dean et al. (2012) introduced DistBelief, the first industrial-scale parameter server architecture, demonstrating distributed SGD over 1,000 machines on speech recognition tasks at Google. Zinkevich et al. (2010) provided theoretical convergence guarantees for parallel SGD with gradient averaging, establishing the rate O(1/√(NT)) for convex objectives with N workers and T steps. Li et al. (2014) formalised the parameter server framework with more flexible consistency models and demonstrated scalability to billions of parameters.
  • Goyal et al. (2017, “Accurate, Large Minibatch SGD: Training ImageNet in 1 Hour”) established the linear scaling rule: to preserve the statistical behaviour of SGD when multiplying the mini-batch size by k, multiply the learning rate by k and employ a linear warm-up schedule for the first 5 epochs. This insight enabled synchronous DDP training of ResNet-50 on ImageNet in under one hour across 256 NVIDIA P100 GPUs, proving that data parallelism could achieve state-of-the-art accuracy at industrial scale without approximation. Rajbhandari et al. (2020) introduced ZeRO (Zero Redundancy Optimizer) at Supercomputing 2020, demonstrating training of 17B-, 40B-, and 100B-parameter models by eliminating the redundant storage of model states in classical DDP. The three ZeRO stages progressively partition optimiser states, gradients, and parameters across ranks, achieving 16-fold memory reduction at full ZeRO-3 with modest communication overhead. Ren et al. (2021) extended ZeRO to CPU and NVMe storage with ZeRO-Offload, enabling billion-parameter model training on a single GPU by offloading optimiser states to CPU memory. Zhao et al. (2023) introduced PyTorch FSDP as a production-grade open-source implementation, reporting throughput parity with DeepSpeed ZeRO-3 on GPT-3-scale models while providing a more Pythonic API. Douillard et al. (2023) proposed DiLoCo (Distributed Low-Communication training) from Google DeepMind, demonstrating that performing H=500 inner SGD steps between outer gradient synchronisations reduces communication volume by 500-fold with only 1–2% degradation on language modelling benchmarks, enabling cross-datacenter training without the latency sensitivity of synchronous DDP. The May 2025 empirical study (arXiv:2505.12832) provided systematic benchmarking of DDP, FSDP, and parameter server approaches across GPU cluster configurations, finding FSDP2 optimal for models with 1B–100B parameters and standard DDP optimal for models that fit within a single GPU’s memory.
  • Key research groups advancing data parallelism theory and practice include: the Microsoft Research AI team (ZeRO, DeepSpeed), NVIDIA’s Megatron team (3D parallelism, NCCL), the PyTorch distributed team (DDP, FSDP2, DTensor), Google DeepMind (DiLoCo, Pathways), Meta AI Research (FairScale, OPT training), and academic groups at CMU (MLSys), Berkeley (Ray/ICSI), ETH Zürich (scaling theory), and the Edinburgh Parallel Computing Centre (EPCC, distributed scientific ML).

Current Landscape (2026)

  • As of 2026, data parallelism remains the dominant first-order scaling primitive for all frontier AI training runs. PyTorch DDP and FSDP2 are the de facto open-source frameworks, with FSDP2’s device_mesh API enabling composable hybrid parallelism across data-parallel and tensor-parallel dimensions through a unified device mesh abstraction. DeepSpeed ZeRO-3 continues to be widely used in research settings, particularly via HuggingFace Accelerate integration, which provides a unified API over DDP, FSDP, and DeepSpeed backends. NCCL 2.23 (released 2025) introduced the PAT (Parallel-All-Transfer) algorithm alongside the established ring and tree algorithms, and the NCCLX variant supports collectives over more than 100,000 GPUs — driven by xAI’s Colossus cluster (100,000 H100 GPUs) and similar hyperscaler deployments at Meta (Grand Teton), Microsoft (Maia), and Google (TPUv5 pods). These at-scale deployments have revealed new bottlenecks: all-reduce latency at 100k+ ranks begins to dominate even BF16 gradient volumes, motivating hierarchical all-reduce strategies where gradients are first reduced within each node group and then synchronised across groups.
  • Inference-time parallelism increasingly mirrors training strategies: vLLM, TensorRT-LLM, SGLang, and llama.cpp now support tensor-parallel and data-parallel deployment for serving, with continuous batching (iteration-level scheduling) combining with data-parallel replica sets to maximise GPU utilisation under variable request loads. Ring-based all-reduce is being supplemented by tree algorithms and switch-fabric-aware topologies for disaggregated data centre architectures where fat-tree and dragonfly networks have different optimal collective implementations. In-network computing — using Mellanox SHARP (Scalable Hierarchical Aggregation and Reduction Protocol) built into InfiniBand switches — has moved from experimental to production status, offloading all-reduce arithmetic into the network fabric and reducing host CPU/GPU involvement in gradient aggregation by up to 10-fold for small- to medium-scale collectives. DiLoCo-style intermittent synchronisation is gaining traction for cross-datacenter and on-device training scenarios. The EU AI Act’s Annex III provisions on high-risk AI system documentation, combined with the UK government’s AI Opportunities Action Plan (January 2025), are driving interest in reproducible and auditable training configurations, making explicit parallelism strategies, batch sizes, and learning rate schedules part of standardised model cards and training reports.
  • The boundary between data parallelism and model parallelism is blurring at the implementation level. FSDP2 and ZeRO-3 can be viewed as forms of model parallelism (they shard model parameters) that expose a data-parallel API. Megatron-Core’s latest releases (2025) provide a unified parallelism specification language where the programmer specifies a (data, tensor, pipeline, expert, context, sequence) parallelism tuple and the framework automatically handles collective placement and scheduling. Context parallelism — a new dimension introduced in 2024 for long-context Transformer Architecture training — partitions the sequence dimension across devices, complementing data parallelism’s batch-dimension partitioning.

Key Terminology

  • All-Reduce: A collective communication operation in which every worker sends its partial result (gradient tensor) and receives the aggregate (sum or average) such that all workers end with identical values. Foundational to synchronous data parallelism.
  • Ring All-Reduce: The bandwidth-optimal all-reduce algorithm that arranges workers in a logical ring, divides the gradient tensor into N chunks, and performs a reduce-scatter followed by an all-gather in 2(N-1) communication steps. Communication volume per worker converges to 2G bytes as N → ∞.
  • NCCL (NVIDIA Collective Communications Library): The GPU-native library implementing AllReduce, AllGather, ReduceScatter, Broadcast, and other collective operations optimised for NVLink, InfiniBand, and RoCE interconnects. The de facto standard for GPU-based distributed training.
  • ZeRO Optimizer (Zero Redundancy Optimizer): Microsoft DeepSpeed’s memory-efficiency technique that partitions optimiser states (ZeRO-1), gradients (ZeRO-2), and model parameters (ZeRO-3) across data-parallel ranks, eliminating the redundant copies stored in classical DDP.
  • Fully Sharded Data Parallel (FSDP): PyTorch’s production implementation of ZeRO-3-style parameter and gradient sharding. FSDP2 (PyTorch ≥ 2.4, 2024) uses DTensor for sharding semantics, enabling composability with Tensor Parallelism via device_mesh.
  • Gradient Staleness: In asynchronous data parallelism, the number of steps by which a gradient is delayed relative to the current model version when it arrives at the parameter server. Bounded staleness (τ ≤ 8–16) is tolerable under appropriate learning-rate decay.
  • Linear Scaling Rule: The heuristic (Goyal et al., 2017) that with k-fold batch size increase via data parallelism, the learning rate should be multiplied by k, combined with a k-step linear warm-up, to maintain equivalent convergence behaviour.
  • Straggler: A slow worker node in synchronous data parallelism that delays the all-reduce barrier, reducing GPU utilisation across the entire replica group. Hardware heterogeneity, network congestion, and OS jitter are common causes.
  • 3D Parallelism: The combination of data parallelism (batch dimension), Tensor Parallelism (hidden dimension within layers), and Pipeline Parallelism (layer dimension across stages) used in production LLM training stacks (Megatron-LM, DeepSpeed). Data parallelism is the outermost, most bandwidth-frugal dimension.
  • Bucket-wise All-Reduce (DDP gradient bucketing): PyTorch DDP’s optimisation that groups small parameter gradients into fixed-size buckets (default 25 MB) and launches all-reduce operations per bucket as soon as all gradients in the bucket are ready, overlapping communication with the backward pass to hide latency.

UK Context

  • The UK’s academic and industrial AI community engages deeply with data-parallel infrastructure across both research and production settings. The Edinburgh Parallel Computing Centre (EPCC) at the University of Edinburgh, which hosts the ARCHER2 national supercomputer (748,544 CPU cores, Cray Shasta architecture, HPE Slingshot interconnect), develops distributed ML workflows leveraging data-parallel strategies for scientific workloads including protein structure prediction (AlphaFold distributed training), climate modelling, and computational fluid dynamics. EPCC’s EPSRC Centre of Excellence in Exascale Computing actively researches communication-efficient collective algorithms and adaptive scheduling strategies for heterogeneous HPC/ML workloads. Imperial College London’s High Performance Computing Service and Computing Department maintain research groups on distributed systems, runtime scheduling, and scalable ML frameworks, with active collaboration with industry partners including NVIDIA and Arm. The University of Manchester’s NVIDIA AI Technology Centre (NVAITC) partnership accelerates DDP and FSDP-based training on GPU clusters and contributes upstream to PyTorch distributed development.
  • In Northern England, the manufacturing and industrial AI context is significant. Sheffield’s Advanced Manufacturing Research Centre (AMRC) applies distributed ML for predictive maintenance on factory floor sensor networks, using data-parallel training pipelines for anomaly detection at scale. The N8 Research Partnership (Universities of Durham, Lancaster, Leeds, Liverpool, Manchester, Newcastle, Sheffield, and York) operates the Bede HPC cluster (NVIDIA V100 and A100 GPUs, IBM Power9 nodes) providing shared data-parallel compute for regional research. Leeds Institute for Data Analytics (LIDA) uses data-parallel distributed training for health analytics across NHS Yorkshire datasets. Newcastle’s National Innovation Centre for Data (NICD) supports SMEs in deploying data-parallel training pipelines for industrial applications.
  • Industrial AI labs in London operate at significant scale. Google DeepMind’s London offices are a major consumer of large-scale data-parallel training infrastructure on Google’s TPU pods and H100 GPU clusters; key research outputs including AlphaFold 2/3, Gato, Gemini 1.0/1.5, and Lyria (music generation) were trained using data-parallel strategies across thousands of accelerators. The Gemini Ultra training run, for instance, used a combination of model parallelism and data parallelism across TPU v5 pods with custom JAX-based distributed training infrastructure. ARM Holdings (Cambridge) designs CPU and NPU IP used in billions of edge inference devices; ARM’s ML Research Lab investigates on-device federated learning — a form of data parallelism over private, distributed data — as a mechanism for training foundation models without centralising user data. Wayve (autonomous driving) and Stability AI (generative models) operate multi-GPU clusters in London and Cambridge. Graphcore (Bristol), now acquired by SoftBank, developed the Intelligence Processing Unit (IPU) and associated Poplar SDK, which implements a Bulk Synchronous Parallel (BSP) execution model — equivalent to synchronous data parallelism but with all-reduce executed entirely on-chip using Graphcore’s proprietary IPU-Link interconnect at 2.8 TB/s, eliminating the GPU-host-fabric-switch bottleneck of GPU clusters.
  • The UK government’s AI Opportunities Action Plan (January 2025), authored by Matt Clifford, committed £900 million to the AI Research Resource (AIRR) programme, which specifically targets GPU cluster infrastructure for large-scale data-parallel training. The Bristol-based Isambard-AI supercomputer, part of AIRR and powered by 5,448 NVIDIA Grace Hopper Superchips (GH200), provides 200+ petaFLOPS of AI compute. This is complemented by the Edinburgh AI supercomputer (Rapsodi, provisioned through the AI Research Resource tender) and industrial-scale compute through the Innovate UK cloud compute voucher scheme. These initiatives position the UK as a competitive jurisdiction for large-model pre-training despite the current compute disadvantage relative to US hyperscalers.

Future Directions (2026–2030)

  • Intermittent synchronisation at scale: DiLoCo-style and FedAvg-inspired low-communication data parallelism will mature from research prototypes into production training stacks, enabling training across geographically distributed or heterogeneous GPU clusters with high-latency wide-area interconnects. The open question is how to adapt the outer learning rate and aggregation frequency dynamically as model gradients evolve during training.
  • In-network compute (INC): SmartNICs and in-network all-reduce (BlueField DPUs, Mellanox SHARP, Broadcom Triton) will move gradient aggregation arithmetic into the network switch fabric, reducing GPU-to-host DMA and PCIe overhead. SHARP (Scalable Hierarchical Aggregation and Reduction Protocol) in InfiniBand HDR/NDR is already production-ready; new generations will extend to 400G ethernet fabrics.
  • Heterogeneous and edge data parallelism: Mixed-precision, mixed-device (GPU + TPU + IPU + NPU) data-parallel training will require adaptive work partitioning proportional to device throughput, topology-aware gradient routing that minimises cross-domain communication, and dynamic rescheduling when slower devices become bottlenecks.
  • CXL memory disaggregation: Compute Express Link (CXL) memory pooling will allow workers to dynamically share parameter memory pages across a CXL fabric without the explicit all-gather communications required by ZeRO-3, enabling far higher data-parallel efficiency for models whose parameter shards fit within the CXL pool capacity.
  • Privacy-preserving distributed training: Integration of differential privacy mechanisms (DP-SGD with per-sample gradient clipping and Gaussian noise injection) into FSDP/ZeRO stacks will become standard in regulated sectors (healthcare, finance, public sector) under GDPR Article 25 (data protection by design), the UK Data (Use and Access) Act 2025, and the EU AI Act’s requirements for privacy-preserving processing in high-risk AI applications.
  • Compiler-automated parallelism planning: torch.compile (TorchDynamo + Inductor), XLA, and emerging MLIR-based compiler stacks will absorb parallelism specification, automatically partitioning models and schedules across available device topologies without explicit user annotation. GSPMD (Google, 2021) and its successors are early demonstrations; by 2028 this is expected to be the dominant mode of distributed training specification. This will integrate Data Preprocessing pipeline compilation with training computation, enabling end-to-end optimisation of the full ML workflow from raw data ingestion through Gradient Aggregation and parameter update.
  • Model-free adaptive parallelism: Reinforcement Learning agents will optimise the (D, T, P, EP, SP) parallelism configuration, Batch Size, gradient compression level, and Checkpoint frequency jointly, replacing the current manual grid-search over parallelism configurations. This connects data parallelism research to AutoML methodology.
  • Sequence and context parallelism: As context windows for Transformer Architecture models expand to millions of tokens (Gemini 1.5 Pro: 1M tokens; Claude 3 Opus: 200k tokens), sequence-dimension parallelism — partitioning the attention computation across the sequence axis — will become a standard fifth dimension of the (data, tensor, pipeline, expert, sequence) parallelism hierarchy.
  • Sustainable and energy-efficient data parallelism: As training energy costs draw regulatory attention (EU AI Act energy disclosure requirements), communications-efficient data-parallel strategies that reduce all-reduce frequency (à la DiLoCo) or compress gradients (PowerSGD, 1-bit Adam) will be motivated by energy reduction as much as by throughput gains.

Performance and Benchmarks

  • The throughput of data-parallel training is typically measured in samples per second or tokens per second per GPU, and in the Model FLOPs Utilisation (MFU) metric — the ratio of observed FLOPs throughput to the GPU’s theoretical peak FLOPs. State-of-the-art results on A100 80 GB GPUs as of 2025: GPT-3 (175B parameters) with Megatron-LM 3D parallelism achieves approximately 38% MFU; LLaMA 2 (70B) with FSDP2 achieves approximately 42% MFU at 8-GPU scale, degrading to ~35% MFU at 128-GPU scale due to increased communication overhead. These efficiencies compare to the theoretical maximum of 312 TFLOPs/GPU for A100 in BF16 matrix multiply. On H100 GPUs (989 TFLOPs/GPU BF16), similar configurations achieve 40–48% MFU, reflecting the H100’s improved NVLink 4.0 bandwidth (900 GB/s vs A100’s 600 GB/s) reducing the relative cost of all-gather and reduce-scatter operations.
  • Scaling efficiency — the ratio of N-GPU throughput to N × 1-GPU throughput — is the key metric for evaluating data-parallel configurations. DDP (no sharding) with efficient bucket scheduling achieves 90–95% scaling efficiency at 8 GPUs, degrading to 75–85% at 64 GPUs and 60–75% at 512 GPUs on ImageNet-scale workloads with InfiniBand interconnects. FSDP’s additional all-gather overhead typically reduces scaling efficiency by 5–10% relative to DDP at the same scale but enables training of larger models that would otherwise require model parallelism with higher overheads. The May 2025 empirical benchmark (arXiv:2505.12832) across 8, 16, 32, 64, and 128 A100 GPUs found FSDP2 outperforming DeepSpeed ZeRO-2 by 8% throughput at 64+ GPUs for transformer workloads, attributable to FSDP2’s improved prefetch scheduling that better hides all-gather latency.
  • The communication-to-computation ratio ρ = t_comm / t_comp determines when data parallelism becomes communication-bound. For a model of size P parameters (FP32), all-reduce communicates 2P bytes per step. Compute time per step scales with batch size B as t_comp ∝ B × P / peak_FLOPs. Therefore ρ ∝ 1 / B — larger per-GPU batch sizes reduce the communication-to-computation ratio. This is why practitioners typically maximise per-GPU batch size (limited by GPU memory) before scaling to more GPUs: increasing B first improves GPU utilisation and reduces ρ, whereas increasing the number of GPUs N at fixed B keeps ρ constant while adding N-dependent communication latency from the all-reduce barrier. The practical consequence is that data parallelism is most efficient when per-GPU batch size B ≥ B_min, where B_min is the batch size at which ρ ≤ 0.1 (communication occupies at most 10% of step time). For typical transformer workloads on A100 with InfiniBand HDR, B_min ≈ 512–2048 tokens for 7B-parameter models.
  • The Chinchilla scaling law (Hoffmann et al., 2022) established that compute-optimal training pairs model size M with a dataset size of approximately 20M tokens per billion parameters: N_tokens = 20 × M_billion. This implies that doubling training compute should split equally between model size and training tokens — motivating longer training runs with more data rather than simply larger models. Data parallelism enables the throughput needed to process these Chinchilla-optimal token counts in feasible wall-clock time: training a Chinchilla-optimal 70B-parameter model requires 1.4T tokens; at 10,000 tokens/sec/GPU on 1,024 GPUs (10.24 billion tokens/sec total), this takes approximately 38 hours of continuous training. Without data parallelism, a single A100 at 10k tokens/sec would require 1.6 years. This example illustrates why data parallelism at scale is not merely an optimisation but an enabling technology for the frontier AI research agenda.
  • Benchmark frameworks for evaluating data-parallel training stacks include MLPerf Training (MLCommons), which includes reference implementations for image classification (ResNet-50), object detection (RetinaNet), natural language processing (BERT), recommendation (DLRM), speech recognition (RNN-T), reinforcement learning (MiniGo), and large language models (LLaMA 2 70B). MLPerf Training measures time-to-train to a reference quality threshold, enabling cross-vendor comparison of data-parallel infrastructure. In the November 2024 round, NVIDIA’s H100 clusters achieved GPT-3 175B training in under 11 minutes at 10,752 GPUs — a 3.8× improvement over the same cluster size from the 2022 round, reflecting improvements in NCCL, CUDA, and Megatron-LM communication overlap. AMD’s ROCm ecosystem and RCCL (ROCm Collective Communications Library) provides competitive performance to NCCL on MI300X GPUs for data-parallel training, with Frontier (the world’s first exaFLOP supercomputer at Oak Ridge National Laboratory, using AMD MI250X GPUs) serving as a testbed for 100k-GPU scale data parallelism for scientific ML workloads.

Challenges and Limitations

  • Straggler sensitivity: In synchronous data parallelism, the global step time is bounded by the slowest worker. Hardware variability (thermal throttling, memory errors causing retry overhead), OS jitter (garbage collection, CPU scheduling preemption), and network congestion can make any worker a transient straggler. Straggler mitigation techniques include backup (redundant) workers (as in TensorFlow’s BarrierSyncSGD), dropping the slowest fraction of workers per step, and gradient timeout mechanisms that advance the step with whatever gradients have arrived.
  • Effective batch size and convergence: Data parallelism inherently increases the effective batch size proportionally to the number of workers. Very large batches (beyond 16k–32k images or 4M tokens per step) are known to generalise poorly — the gradient noise at large batch sizes loses the regularisation effect that small-batch SGD provides through its implicit annealing of noisy updates. This is the “large-batch training problem” (Keskar et al., 2017), which motivates the use of learning rate warm-up, learning rate decay, layer-wise adaptive rate scaling (LARS, You et al., 2017), and ultimately bounds the linear scalability of data parallelism: beyond a critical batch size, adding more data-parallel workers improves throughput but degrades per-sample efficiency.
  • Memory pressure at extreme scale: Even FSDP2 and ZeRO-3 face memory pressure from activation memory (the intermediate tensors stored during the forward pass for use in backpropagation). For a 70B-parameter transformer with 4096-token sequence length and 96 layers, activation memory can exceed parameter memory. Gradient checkpointing (selective recomputation) trades compute for memory by discarding and recomputing activations, but introduces 30–40% compute overhead. Activation quantisation (storing activations in INT8 during the forward pass) is an active research area that could reduce activation memory by 4-fold without recomputation overhead.
  • Heterogeneous hardware: Cloud GPU clusters are rarely perfectly homogeneous — mixed GPU generations (A100 + H100), mixed networking (InfiniBand + Ethernet), and variable-speed storage create performance heterogeneity that the all-reduce collective treats as stragglers. Adaptive scheduling strategies and topology-aware collective libraries mitigate this but cannot fully eliminate the throughput penalty from the slowest link in the communication ring.

Research and Literature

    1. Robbins, H. & Monro, S. (1951). A stochastic approximation method. Annals of Mathematical Statistics, 22(3), 400–407. [Foundations of stochastic gradient methods]
    1. Dean, J., Corrado, G., Monga, R., et al. (2012). Large scale distributed deep networks. NeurIPS 2012. [DistBelief, first industrial parameter server]
    1. Zinkevich, M., Weimer, M., Smola, A., & Li, L. (2010). Parallelized stochastic gradient descent. NeurIPS 2010. [Convergence guarantees for parallel SGD]
    1. Li, M., Andersen, D.G., Park, J.W., et al. (2014). Scaling distributed machine learning with the parameter server. OSDI 2014. [Parameter server formalisation]
    1. Goyal, P., Dollár, P., Girshick, R., et al. (2017). Accurate, large minibatch SGD: training ImageNet in 1 hour. arXiv:1706.02677. [Linear scaling rule for distributed DDP]
    1. Sergeev, A. & Del Balso, M. (2018). Horovod: fast and easy distributed deep learning in TensorFlow. arXiv:1802.05799. [Ring all-reduce implementation]
    1. Rajbhandari, S., Rasley, J., Ruwase, O., & He, Y. (2020). ZeRO: memory optimizations toward training trillion parameter models. SC 2020. [ZeRO optimizer stages]
    1. Ren, J., Rajbhandari, S., Aminabadi, R.Y., et al. (2021). ZeRO-Offload: democratizing billion-scale model training. ATC 2021. [ZeRO-Infinity / NVMe offloading]
    1. Narayanan, D., Shoeybi, M., Casper, J., et al. (2021). Efficient large-scale language model training on GPU clusters using Megatron-LM. SC 2021. [3D parallelism in production]
    1. Zhao, Y., Gu, A., Varma, R., et al. (2023). PyTorch FSDP: experiences on scaling fully sharded data parallel. VLDB 2023. [FSDP design and production results]
    1. Shoeybi, M., Patwary, M., Puri, R., et al. (2019). Megatron-LM: training multi-billion parameter language models using model parallelism. arXiv:1909.08053. [Foundation for hybrid parallelism]
    1. Lepikhin, D., Lee, H., Xu, Y., et al. (2020). GShard: scaling giant models with conditional computation and automatic sharding. arXiv:2006.16668. [Expert parallelism]
    1. Rasley, J., Rajbhandari, S., Ruwase, O., & He, Y. (2020). DeepSpeed: system optimizations enable training deep learning models with over 100 billion parameters. KDD 2020. [DeepSpeed framework]
    1. Peng, Y., Zhu, Y., Chen, Y., et al. (2019). A generic communication scheduler for distributed DNN training acceleration. SOSP 2019. [Gradient communication scheduling]
    1. Lian, X., Zhang, C., Zhang, H., et al. (2017). Can decentralized algorithms outperform centralized algorithms? NeurIPS 2017. [Decentralised data-parallel SGD]
    1. Assran, M., Loizou, N., Ballas, N., & Rabbat, M. (2019). Stochastic gradient push for distributed deep learning. ICML 2019. [Gossip-based gradient averaging]
    1. Douillard, A., Feng, Q., Rusu, A., et al. (2023). DiLoCo: distributed low-communication training of language models. arXiv:2311.08105. [Intermittent synchronisation for geographically distributed training]
    1. Huang, Y., Cheng, Y., Bapna, A., et al. (2019). GPipe: efficient training of giant neural networks using pipeline parallelism. NeurIPS 2019. [Pipeline parallelism complementing data parallelism]
    1. Lin, Y., Han, S., Mao, H., Wang, Y., & Dally, W.J. (2017). Deep gradient compression. ICLR 2018. [Gradient compression for DDP]
    1. Jain, P., Nrusimha, A., Jain, N., et al. (2020). Checkmate: breaking the memory wall with optimal tensor rematerialization. MLSys 2020. [Memory management in distributed training]
    1. Artetxe, M., Bhosale, S., Goyal, N., et al. (2021). Efficient large scale language modeling with mixtures of experts. arXiv:2112.10684. [Expert parallelism combined with data parallelism]
    1. Wang, S., Li, B.Z., Khabsa, M., Fang, H., & Ma, H. (2020). Linformer: self-attention with linear complexity. arXiv:2006.04768. [Memory-efficient transformer related to scaling]
    1. Dettmers, T., Lewis, M., Shleifer, S., & Zettlemoyer, L. (2022). 8-bit optimizers via block-wise quantization. ICLR 2022. [Quantised optimiser states in sharded DDP]
    1. Arxiv (2025). A study on distributed strategies for deep learning applications in GPU clusters. arXiv:2505.12832. [Empirical evaluation of DDP, FSDP, and parameter server trade-offs]
    1. Arxiv (2025). Collective communication for 100k+ GPUs. arXiv:2510.20171. [NCCLX and PAT algorithm for massive-scale all-reduce]
    1. Arxiv (2024). Efficient training of large language models on distributed infrastructures: a survey. arXiv:2407.20018. [Comprehensive survey of parallelism strategies]
    1. Kaplan, J., McCandlish, S., Henighan, T., et al. (2020). Scaling laws for neural language models. arXiv:2001.08361. [Theoretical basis for data-parallel scaling with batch size]
    1. Micikevicius, P., Narang, S., Alben, J., et al. (2018). Mixed precision training. ICLR 2018. [FP16 training enabling larger batch sizes in DDP]

Provenance