Conceptio › Archive › arXiv CS
arXiv CSopen access

The Life of a Token: from Words to Bits on the Wire

· arxiv_cs
arXiv CS · Papers · License: Open Access
Open Source ↗Direct PDF ↓
distributed-systemsinternetnetworkingprotocols
networking, internet, protocols, distributed systems

The Life of a Token: From Words to Bits on the Wire Davide Avesani1,∗, Pengwenlong Gu1 , Sotiris Skaperas1 , Stefano Secci1,1

arXiv:2609.19924v1 [cs.DC] 17 Sep 2026

Abstract Large Language Models (LLMs) transform vast collections of unstructured text into semantic patterns used for language generation and reasoning tasks. Behind their ease of use lies a complex process: words become tokens, tokens become vectors, and vectors ultimately give rise to streams of bits that flow through HighPerformance Computing (HPC) systems. As modern LLMs grow to billions or trillions of parameters, this path increasingly unfolds across thousands of interconnected accelerators, making the underlying communication fabric a critical and often opaque component of model training. This tutorial follows the journey from words to network traffic and explains how language is translated into communication flows within HPC training systems. Using concrete examples from Dante’s Divine Comedy, we illustrate how model architecture, tokenization, embeddings, and parallelization strategies shape the volume, structure, and timing of data exchanged across the network. We combine architectural analysis, analytical traffic models, and numerical examples to characterize the communication requirements of LLM training. Our goal is to demystify how words travel across the network and provide practical insights into the network capabilities required to support the journey from text to a trained model. Keywords: large language models, distributed training, collective communication, communication traffic modeling, high-performance computing 1. Introduction: The journey begins In the opening lines of La Divina Commedia, Dante Alighieri finds himself in a dark forest [1], poised at the threshold of transformation. Large Language Model (LLM) systems begin in a similarly opaque realm: a vast, unstructured forest of text. Through training, this text is transformed into statistical and semantic representations that enable language generation, reasoning, and generalization across unseen tasks [2]. Beneath the natural language ability of LLMs lies a concrete systems process. Words are transformed into tokens, tokens into vector representations, vectors into activations and gradients, and these ‘tensors’ ultimately into streams of bits exchanged across large scale distributed infrastructures. This tutorial examines this transformation from a communication oriented perspective. Rather than focusing on model accuracy or architectural optimization, we focus on how the operations performed during LLM training translate into data movement among accelerators. Selected passages from La Divina Commedia are used as a running example to ground the discussion. By following specific fragments of the poem through tokenization, distribution, and processing, we connect abstract training concepts to concrete data flows within the infrastructure. This narrative anchors technical concepts to a tangible reference point, enabling an intuitive and structured explanation of the computational and communication mechanisms involved. Tracing the life of a token has become increasingly challenging as it requires the understanding of both the model that processes it, and the infrastructure that makes such processing possible at scale. Since the introduction of the Transformer architecture [3], improvements in LLM performance have been largely driven by increases in model size, training data, and computational budget. This trend is reflected in the progression from models with hundreds of millions of parameters, such as GPT-2, to hundreds of billions of parameters, as in GPT-3, and more recently to models approaching the trillion parameter scale. Scaling law studies further show that model performance follows predictable trends as model size, dataset size, and compute increase [4]. This scaling pushes training far beyond the computational and memory capacity of a single machine. Modern LLMs are therefore trained on large clusters composed of hundreds or thousands of specialized accelerators, typically Graphics or Tensor Processing Units (GPUs or TPUs). For instance, GPT-3 was reportedly trained using more than 10,000 (NVIDIA V100) GPUs [5], while GPT-4 is estimated to have used on the order of tens of thousands of (NVIDIA A100) GPUs [6]. ∗ Corresponding author

Email addresses: [email protected] (Davide Avesani), [email protected] (Pengwenlong Gu), [email protected] (Sotiris Skaperas), [email protected] (Stefano Secci)

Life of a Token: from text to network traffic in distributed LLM training A tutorial roadmap from token creation to communication and network implications Dataset and Tokenization

1

Section III — determining T

2

Transformer and Training Tensors

Distributed Training

3

Section IV — determining d, L, P, and mini-batch size

Raw text corpus D

Section V — where traffic originates

Node 1

Node 2

⚙

518

1527

Preprocessing and tokenization

Main parallelization strategies Data Parallelism (DP) – gradients

311

2690

Mini-batch / Local batch processing

•••

Input tensor:

Token IDs Fixed-length sequence (length T)

Key output: Token volume, sequence count, sequence length T

Pipeline Parallelism (PP) – activations between stages

Activation payload

Sact = BTdsact

Gradients payload

Sgrad = Psgrad

AllReduce

◌

ReduceScatter / AllGather

Inter-node network

•••••

↗↙

↔

Point-to-Point

✣

All-to-All

┴

Broadcast / Scatter

NVLink / NVSwitch / PCIe

×

GPU

↔

↔

Tensor payload

Collective algorithm

RoCE / InfiniBand / Ethernet

↔

↔

•••

↔

↔

↔

Performance factors Per-rank traffic

• Bandwidth • Latency • Congestion

Key output: Estimate size, frequency, and destination of each communication object

Key output: Which tensors are exposed to communication

↗↙

Network topology (e.g., leaf–spine / Clos)

Mapping flow

MoE / Expert Parallelism (EP) – routed tokens

Key output: Tensor dimensions and candidate communication objects

Intra-node interconnects

↕↔

Tensor / Sequence Parallelism (TP / SP) – intra-layer activation exchange

X ∈ ℝB × T × d

Section VII — from payloads to bits on the wire

Communication primitives

Node N

•••

Stack of L Transformer layers •••

Network Infrastructure and Implications

5

Section VI — from tensors to communication payloads

Accelerators / GPUs across nodes

••• Token embeddings (dim. d)

“Nel mezzo del cammin di nostra vita...” (La Divina Commedia)

Collective Communication and Traffic Estimation

4

• Placement • Topology

Final outcome: Network traffic on links, packets, and bytes

The journey of a token

Becomes network traffic

How to use the framework 1

Identify model and training dimensions (d, L, P, T, s, V)

2

Derive tensor sizes (Sact, Sparam, Sgrad)

3

Identify the parallellism strategy (derive B and M)

4

Map tensors to communication primitives (AllReduce, P2P)

5

Estimate traffic per collective (Collective Operations)

6

Interpret the result through topology and transport

Figure 1: Illustration of the Life of a Token framework. The figure traces the transformation of a token across the distributed LLM training stack, from raw text and tokenization, through model processing and parallel execution, to the network communication generated during large scale training.

The distribution of training across such a large system considerably increases the complexity of tracing this journey. Moreover, in such large distributed training scenario, computation is no longer the only limiting factor. Model parameters, gradients, optimizer states, and intermediate activations must be exchanged across devices, making communication a central component of the training process. Several studies have shown that communication can account for a significant fraction of LLM training time and, depending on the specific scenario, may even become the main performance bottleneck [7, 8, 9, 10, 11, 12]. As a result, improving large scale LLM training requires not only faster accelerators, but also communication infrastructures that are optimized to support the traffic generated during training. Accurately characterizing LLM training traffic is therefore essential for designing infrastructures capable of supporting future training workloads, including scenarios in which accelerator clusters are provided as shared or on demand G/TPU as a Service platforms. A complete account of this communication behavior is challenging because it is jointly shaped by multiple, tightly coupled factors: (i) dataset properties and partitioning, (ii) model architecture and scale, (iii) parallelization strategies and workload distribution, (iv) Collective Communication Operations (CCOs), and (v) network topology, interconnect technologies, and communication protocols. These factors are typically studied across separate communities, including machine learning, distributed systems, and networking. As a result, their combined effect on communication behavior remains difficult to capture, making accurate prediction and system-level optimization a nontrivial problem. To bridge these perspectives, we introduce the ‘Life of a Token‘ framework, illustrated in Fig. 1, as the organizing abstraction of the tutorial: it guides the reader through data processing during distributed LLM training, from raw text to tokens, Transformer tensors, and network traffic. Linking these stages provides a unified view of how decisions made at the data, model, and execution levels imply changes at the network-level. 1.1. Tutorial objective and contributions The literature on LLMs has expanded rapidly in recent years, driven by significant industrial investment, continuous architectural innovation, and the growing scale of deployed models. At the same time, training these models across increasingly large and distributed accelerator clusters has made the network a critical component of the overall training system. Optimizing network performance is therefore essential, but doing so requires first characterizing the communication traffic generated during training. 2

Although prior work provides detailed insights into specific aspects of LLMs and distributed training, the resulting body of knowledge remains fragmented across machine learning, distributed systems, and networking. As a result, translating high level model descriptions into concrete communication requirements for real world deployments remains challenging. This tutorial addresses this gap by providing a communication oriented framework for reasoning about how LLM training generates network traffic. We identify four limitations that motivate the need for such a tutorial. 1. Limited connection between model design and infrastructure design. The literature provides extensive resources on the end-to-end development of LLMs, including widely recognized works such as [13] and [14]. However, these contributions primarily focus on model architecture, training procedures, and implementation details, while providing limited insight into how model-level choices translate into infrastructure requirements. In particular, the impact of AI model parameters on communication volume and network requirements remains insufficiently exposed. 2. Lack of systematic tools for communication modeling. Since the introduction of GPT-3 [15], LLM architectures have rapidly evolved, incorporating numerous optimizations and design variations. This diversity makes it difficult to consistently evaluate the communication implications of different configurations. A step by step methodology for estimating communication payloads and mapping them to network-level requirements is still largely missing. 3. Fragmented analysis of communication mechanisms. Prior work has analyzed and demystified important components of the distributed training stack, including communication patterns in Transformer models [16], collective communication protocols and algorithms [17], large scale parallel training systems [11, 12], and production network fabrics for AI workloads [18]. However, these contributions are typically organized around individual layers of the stack, rather than around the end-to-end path from data to network traffic. A unified perspective that connects dataset processing, model architecture, tensor dimensions, parallelization choices, CCOs, and network infrastructure is still needed. 4. Incomplete reporting of communication relevant training details. Technical reports on state of the art (SOTA) LLMs often emphasize architecture, dataset scale, context length, and benchmark performance [19, 20, 21, 22, 23, 24]. In contrast, details such as parallelization strategy, accelerator placement, interconnect technology, and network topology are often omitted or reported inconsistently. This makes it difficult to translate model level descriptions into communication payloads and network requirements. Closest works take complementary but distinct perspectives. Duan et al. [25] survey distributed LLM training across infrastructure, parallelization, system optimization, and reliability, prioritizing comprehensive coverage rather than a continuous derivation of communication traffic. Liang et al. [26] classify communication efficient techniques across algorithms, frameworks, and infrastructure, focusing primarily on mechanisms for reducing communication overhead rather than on how model and training choices generate that communication. Tazi et al. [27] provide a practical tutorial with code and extensive scaling experiments, emphasizing how distributed training techniques are implemented and combined. Song et al. [28] instead follow collective operations through planning, execution, adaptation, and computation-communication coordination, after those operations have already been exposed by the workload. The Life of a Token complements these works by beginning upstream, with raw text and a specified model and training configuration, and following a consistent derivation through tokenization, tensor construction, parallelization, collective invocations, and network traffic. Its contribution is not broader coverage or a new parallelization algorithm, but a reusable cross-layer derivation method. Fig. 1 summarizes the framework, formalized in Section 2. The tutorial makes four contributions: 1. We present a token-centered, five-stage framework connecting dataset processing, Transformer tensors, distributed execution, collective communication, and network infrastructure. 2. We provide a step-by-step methodology for estimating the principal tensor payloads, communication operations, invocation frequencies, and per-accelerator and aggregate traffic associated with representative DP, PP, and TP configurations. 3. We connect logical payload estimates to representative collective algorithms and discuss how interconnect bandwidth, latency, topology, and accelerator placement influence communication behavior and training scalability. 4. We combine examples based on Dante’s Divine Comedy with a consistent GPT-2-like reference configuration to turn abstract model and training quantities into concrete payload and traffic estimates. Table 1 summarizes these distinctions. The remainder follows the framework’s five stages, while Appendix A provides additional derivations and extends the analysis to Mixture-of-Experts (MoE) models.

3

Table 1: Positioning of this tutorial with respect to representative adjacent bodies of work. Adjacent line of work LLM foundations, architectures, and scaling

Datasets and tokenization

Distributed LLM training and parallelization

Collective communication and communication libraries

AI cluster and datacenter networking

Communication-oriented surveys, characterization, simulation, and traffic modeling

Recent large-model technical reports

Representative references Transformer and GPT-style models, scaling laws, and general LLM surveys [2, 3, 4, 15, 29, 30]

Main focus Explain the architectural principles of LLMs, including attention, Transformer blocks, decoder-only models, and the scaling of model size, data, and compute. Large-scale corpora, dataset conDescribe how raw text is colstruction, tokenization algorithms, lected, cleaned, tokenized, and and token-counting rules of thumb organized into token sequences for [31, 32, 33, 34, 35, 36, 37, 38, 39] LLM training.

Data, tensor, pipeline, sequence, and hybrid parallelism; memory optimization; distributed-training surveys, systems, and practical tutorials [25, 27, 40, 12, 41, 42, 43, 44, 45, 46, 47, 48, 49, 50, 51, 52, 53] Collective primitives, CCL implementations, collective synthesis, topology-aware algorithms, and collective-centric tutorials [28, 54, 55, 56, 17, 57, 58, 59, 60] GPU interconnects, RDMA, InfiniBand, RoCE, NVLink/NVSwitch, TPU-scale systems, congestion control, and AI-oriented data-center fabrics [61, 62, 63, 64, 65, 66, 18, 67, 68, 69] Communication surveys, empirical characterization, bandwidth studies, training simulators, and topology-aware analyses [26, 16, 10, 6, 70, 71, 72, 73]

State-of-the-art LLM reports and large-scale model descriptions [24, 22, 74, 75, 23, 76, 77, 21, 20]

Explain, implement, or optimize techniques for partitioning data, parameters, layers, activations, and computation across accelerators.

Analyze operations such as AllReduce, ReduceScatter, AllGather, Broadcast, point-to-point exchange, and All-to-All, together with their planning and implementation over GPU clusters. Study the infrastructure that transports training traffic, including intra-node and inter-node interconnects, transport protocols, congestion mechanisms, and data-center topologies.

Difference from this tutorial Provide limited guidance on how architectural quantities translate into communication payloads and network requirements. Focus on data preparation but rarely follow how tokenization and sequence construction influence tensor dimensions and distributed communication payloads. Organize the problem around parallelization techniques, training systems, or implementation guidance rather than a continuous derivation from tokenized input to communication payloads and network requirements. Generally begin with alreadydefined messages or collective operations, leaving the upstream connection to tokenization, model tensors, and parallelization choices implicit. Typically treat the communication workload as an input to network analysis rather than deriving it from model dimensions, training choices, and parallelization strategies.

Survey, measure, model, or simulate communication behavior in distributed Transformer and LLM training systems.

Primarily produce taxonomies, platform-specific measurements, or simulator predictions rather than a single pedagogical procedure connecting raw text, tensor construction, collective invocations, and network traffic. Report model architectures, train- Often provide limited visibility ing data, context lengths, scaling into the communication patchoices, and performance improve- terns, CCOs, network topology, ments for modern LLMs. and infrastructure assumptions required to train such models at scale.

2. The Life of a Token framework This section formalizes the Life of a Token framework that structures the rest of the tutorial. We first present it as a conceptual abstraction for following data across the training stack, and then describe how the same abstraction is used to reason about communication payloads. We then turn this abstraction into a step-bystep methodology for estimating communication payloads and conclude by defining the notation and modeling scope used in the rest of the paper. The notation, assumptions, and key quantities underlying the framework are introduced here and further developed in Section 4, where we examine in detail how they arise from the LLM training pipeline. 2.1. Framework overview At a high level, the framework provides the organizing abstraction used throughout the tutorial to connect the logical steps of LLM training with the communication traffic generated during distributed execution. In this view, the token acts as a conceptual thread: it starts as a discrete identifier produced by the tokenizer, becomes part of a dense activation tensor after the embedding layer, propagates through the Transformer stack, and eventually contributes to the tensors exchanged across accelerators during training. Fig. 1 summarizes the tutorial roadmap in five stages, which are formally introduced in Section 2.2. We introduce the framework using a dense decoder-only Transformer as the reference path, since it cleanly exposes the main data transformations and communication mechanisms. More specialized architectures, especially MoE models, extend this baseline by adding communication stages through token routing. We use the dense case to build intuition and revisit MoE models later as a natural extension. 2.2. Communication estimation methodology We now apply the five stage roadmap in Fig. 1 as an operational procedure for moving from the training input to the communication objects eventually observed by the network. 4

1. Dataset and tokenization: the first stage determines how raw text is transformed into the token sequences processed during training. Starting from the training dataset, the text is cleaned, tokenized, and organized into fixed length sequences of length T . This stage determines the total number of tokens, the number of training sequences, and the sequence length used by the model. 2. Transformer and training tensors: the second stage follows how token identifiers are mapped into dense embeddings and processed by the Transformer stack. For the dense decoder-only Transformer used as a reference, the main quantities are the vocabulary size V , the sequence length T , the hidden dimension d, the number of Transformer layers L, the number of trainable parameters P , the mini-batch size B, and the numerical precision s, expressed in bytes per value. After the embedding layer, the input tensor has shape X ∈ RB×T ×d . (1) The corresponding activation payload is Sact = BT dsact .

(2)

Similarly, the parameter and gradient payloads can be approximated as Sθ = P sθ (3a) Sgrad = P sgrad (3b) These quantities do not yet represent network traffic. They identify the candidate payloads that may become communication objects once training is distributed across multiple accelerators. 3. Distributed training: the third stage considers how the training workload is partitioned across accelerators. Different parallelization strategies expose different tensors to communication, leading to distinct communication patterns and traffic characteristics. This stage therefore determines which tensors become communication objects. 4. Collective communication and traffic estimation: the fourth stage maps the identified communication objects to abstract patterns and estimates their traffic. Depending on the parallelization strategy, tensors may be exchanged between replicas, passed between pipeline stages, or redistributed across devices within a layer. At this stage, the tensor determines the logical payload, while the communication pattern determines how that payload is exchanged across devices. 5. Network infrastructure and implications: the fifth stage maps the estimated payloads onto the underlying network infrastructure. Communication may remain within a node, traverse inter-node links, or cross multiple tiers of the data center network. Actual communication behavior depends on interconnect bandwidth, latency, routing, contention, congestion control, and topology aware placement. By identifying which tensors dominate communication, how often they are exchanged, and where they travel through the system, the framework eases reasoning about parallelization strategies, collective algorithms, accelerator placement, and network requirements. Overall, this methodology moves beyond a token-level view and frames the analysis as a cross layer reasoning tool. It links modeling and training decisions to the main communication objects, the ways they are exchanged, and the network resources they place under pressure. 2.3. Notation and scope Table 2 summarizes the tutorial’s notation. We distinguish tensor payloads from network traffic: quantities such as Sact , Sθ , and Sgrad identify the size of candidate communication objects, while the actual traffic observed on the network depends on the distributed training strategy, the collective communication algorithm, and the underlying topology. Low level implementation details, kernel level optimizations, memory allocation behavior, and framework specific scheduling choices are outside the main scope unless they directly affect communication patterns. Hence we want to emphasize the cross layer relationship between model configuration, tensor dimensions, communication primitives, and network-level behavior. 3. Dataset and tokenization: determining T The training dataset plays a central role in LLM development, as its scale, quality, and diversity shape model behavior and performance. Modern LLM datasets aggregate heterogeneous sources, including web crawls (e.g., CommonCrawl [32], C4 [78]), books, public academic repositories (e.g., arXiv), code repositories (e.g., GitHub), Wikipedia, and large curated corpora such as The Pile [33]. A comprehensive overview of such datasets is provided by Liu et al. [31]. Once collected, the dataset undergoes preprocessing, where raw data is cleaned, normalized, and structured for training [79]. The resulting text is then tokenized, i.e., converted into a sequence of discrete units called 5

Table 2: Main notation used throughout the tutorial.

Model and data

Communication and storage

Notation

Definition

Notation

Definition

B bµ d

Local mini-batch size. Pipeline micro-batch size. Model hidden size.

A Cdev (M, pcc ) A Csys (M, pcc ) M

D L m P T V V = |V|

Tokenized training dataset. Number of Transformer layers. Number of pipeline micro-batches. Number of trainable parameters. Sequence length in tokens. Tokenizer vocabulary. Vocabulary size.

pcc sact sgrad Sθ = P sθ

Per-accelerator communication volume. Aggregate communication volume. Logical communication payload in bytes. Number of communication participants. Bytes per activation value. Bytes per gradient value. Full model parameter payload.

tokens, which provide the numerical input representation used by the LLM. Various techniques have been developed to perform tokenization, with subword based approaches such as Byte Pair Encoding (BPE) [34], WordPiece [36], and SentencePiece/Unigram LM [35] being the most widely adopted in modern LLMs. These methods balance vocabulary size and representation efficiency by decomposing text into frequent subword units. Recent work has also begun to formalize the theoretical role of tokenization in shaping model efficiency and behavior [37]. The tokenizer is typically trained separately prior to model training and remains fixed during optimization. Its design directly affects the number of tokens generated from raw text and, consequently, the amount of data that must be processed. The tokenized dataset is then divided into fixed length sequences of T ordered token identifiers, with each sequence serving as one training sample. The sequence length T is kept fixed during training to ensure a consistent input representation across samples, enabling efficient tensor based computation and stable optimization. The value of T is kept constant within a training phase to obtain tensors with consistent dimensions. In practice, however, its value may be increased across training stages to extend the supported context length [80, 81]. The sequence length T is therefore an important model and training hyperparameter. Fig. 2 summarizes the transition from raw text to fixed length token sequences. 3.1. Notation and dimensionality Dataset size typically scales with model capacity, ranging from tens of gigabytes (GB) for smaller models to multi-terabyte collections for large-scale systems, as illustrated by datasets such as Dolma and FineWeb [82, 83]. During tokenization, the raw text is transformed into a sequence of discrete tokens drawn from a fixed vocabulary V, whose size is V = |V|. Let D denote the tokenized training dataset. It can be represented as D = (x1 , x2 , . . . , x|D| ), xi ∈ V, where |D| is the total number of tokens in the dataset. The vocabulary size V is a key design parameter in LLM engineering. It is typically fixed during model construction and influences both the expressiveness of the token representation and the number of tokens required to encode a given text. In practice, vocabulary sizes commonly range from tens of thousands to a few hundred thousand tokens, depending on the tokenizer design and the languages covered by the model [84, 85]. The exact mapping from raw text to tokens depends on the tokenization algorithm and its implementation. As a commonly used rule of thumb for English text, one token corresponds on average to approximately four characters, or three quarters of a word [38, 39]. Because most characters in predominantly English text occupy one byte when encoded in UTF-8, this corresponds to an order-of-magnitude estimate of approximately four bytes of raw text per token. If Stext denotes the size of the cleaned text in bytes and βtok the average number of , where βtok ≈ 4 for English-dominated corpora. The raw-text bytes represented by one token, then |D| ≈ Sβtext tok tokenized dataset is subsequently partitioned into N fixed-length training sequences, Si = [ti,1 , ti,2 , . . . , ti,T ], i ∈ {1, . . . , N }, where T is the sequence length and N ≈ |D|/T . Depending on the preprocessing procedure, an incomplete final sequence may be padded, packed with tokens from another document, or discarded. In practice, typical values of T range from 1024 to 4096 tokens for many LLMs, while long-context models may use substantially larger values. Example 3.1 illustrates the scale of tokenized datasets and the resulting sequence counts encountered during LLM training. Example 3.1 (PB Corpus Token and Sequence Scale). Consider a cleaned textual dataset of size Stext = 1 PB, where 1 PB = 1015 bytes. Using the approximation βtok ≈ 4 bytes per token, the total number of tokens 6

1 RAW TRAINING DATASET

2

Unstructured data from sources like books, webpages, code, and articles.

TOKENIZED DATASET

EXAMPLE

TOKENIZATION SIZE: TERABYTES TO PETABYTES (1012 to 1015+ bytes)

3 FIXED-LENGTH SEQUENCES

Dataset → long stream of discrete tokens.

45 N

Tokenized data → fixed length-T sequences

45

“Nel mezzo del cammin di nostra vita” 301 el

757 46748 1624 me zzo del

Example: 7 words 10 tokens In mean: 1 token = 3/4 of a word

6730 cam

1083 min

1891 91127 55576 di nostra vita

SIZE: TERABYTES TO PETABYTES (1012 to 1015+ bytes)

SEQUENCING

t1 t1

301 t2 t2

757

…

1891 91127 55576

t3

Length T … tT−2

tT−1

tT

t3

Length T … tT−2

tT−1

tT

Length T

Figure 2: From dataset size to tokenization to sequence organization

15

is |D| ≈ Sβtext = 104 = 2.5 × 1014 . Thus, a PB scale textual corpus corresponds to approximately 250 trillion tok tokens. If the corpus is divided into fixed length sequences, the number of resulting sequences is approximately 2.5×1014 ≈ 2.44×1011 , corresponding to about 244 billion sequences. Nseq ≈ |D| T . For T = 1024, this gives Nseq ≈ 1024 14 For T = 4096, the number of sequences becomes Nseq ≈ 2.5×10 ≈ 6.10 × 1010 , corresponding to about 61 4096 billion sequences. To provide an intuition of such scale, consider La Divina Commedia, which contains approximately 100,000 words [86]. Using the approximation that one word corresponds to about 1.25 tokens, the poem contains approximately |D| ≈ 100,000 × 1.25 = 125,000 tokens. If the text is divided into sequences of length T = 1024, the number of sequences is approximately 123. Compared with the petabyte scale corpus discussed above, which contains approximately 2.5 × 1014 tokens, the tokenized Divina Commedia represents only 5 × 10−8 % of the total dataset token volume. 4. From tokens to Transformer tensors The previous section described how raw text is converted into fixed-length token sequences of length T . In the Life of a Token framework, we next follow these sequences into the LLM where token-ID tensors are mapped to dense embeddings. After briefly introducing the Transformer architecture and the decoder-only reference model, this section traces how tokens are embedded and processed through the Transformer layers, establishing the main quantities used in the subsequent communication analysis, such as V , T , d, L, P , B, sact and sgrad . 4.1. From the original Transformer to decoder-only LLMs The Transformer architecture forms the core Neural Network (NN) design on which most modern LLMs are built. It was originally introduced as an encoder-decoder model for sequence-to-sequence modeling [3]. It is composed of two main components: an encoder, which processes the input sequence, and a decoder, which generates the output sequence. Over time, this design has evolved into three main variants, each tailored to different types of tasks: • Encoder-decoder: the original formulation, combining an encoder that processes the input and a decoder that generates the output. This structure is commonly used for sequence to sequence tasks such as machine translation and summarization (e.g., T5 [78], BART [87]). • Encoder-only: a simplified version that retains only the encoder component. These models are primarily designed to understand and represent text, and are widely used for tasks such as classification, question answering, and information extraction (e.g., BERT [88], RoBERTa [89]). • Decoder-only: a variant that uses only the decoder component to generate text, one token at a time. This architecture gained significant popularity with the introduction of the GPT-2 model [2] and gained widespread adoption with GPT-3 [15]. It has now become the dominant choice for modern generative LLMs, including the GPT, Gemini, LLaMA, and Mistral model families [15, 76, 81, 90]. 4.2. Reference decoder-only Transformer architecture The decoder-only variant of the Transformer architecture [3], popularized by GPT-style models and in particular by GPT-2 [2], has become the foundation on which modern LLMs are built. As a reference, we adopt a simplified decoder-only Transformer inspired by the GPT-2 architecture [2]. Illustrations of the original encoder-decoder architecture and the decoder-only reference architecture are provided in the Appendix A.1. After GPT-2, many architectural variations have been proposed, but the decoder-only Transformer remains the dominant backbone of modern LLMs. A concise overview of these variants is provided in [91]. In this tutorial, however, we retain the proposed simplified structure throughout the paper for the following reasons: 7

• It captures the core building blocks that remain common across most contemporary LLMs. • It provides a clear and interpretable representation of how data propagates through the model, which is essential for analyzing communication patterns. • Although modern architectures introduce optimizations, the fundamental communication behavior remains largely governed by the same objects, making this abstraction sufficient for our analysis. To maintain clarity and avoid unnecessary complexity, we do not consider advanced architectural extensions here such as MoE, as they introduce specialized communication patterns. MoE architectures and their communication implications are examined in Appendix A.4. The reference architecture is composed of the following components: • Input embedding and positional encoding: this layer transforms the processed tokens into vector representations. Each token is mapped to a dense vector of dimension d through an embedding matrix of size V × d, where V = |V| is the vocabulary size. The embedding matrix is part of the model parameters and is therefore learned jointly with the other model parameters. Positional information can also be added to the embedding representation through mechanisms such as rotary position embeddings [92]. • L stacked Transformer blocks: each block takes a sequence of hidden representations as input and produces updated representations with the same overall tensor shape. The number of Transformer blocks L is a design parameter defined during model construction and varies with the scale of the LLM. In practice, L can range from around a dozen to more than one hundred Transformer layers. For example, the OPT family ranges from 12 layers in OPT-125M to 96 layers in OPT-175B, while MT-NLG 530B uses 105 layers [93, 94]. • Masked multi-head self-attention: each of the L Transformer layers contains a Masked multi-head selfattention block. This component represents the key innovation of the Transformer architecture; it allows each token to selectively relate to previous tokens in the sequence, building a context aware representation based on the surrounding elements. At a high level, this process is implemented through the scaled dot   ⊤ QK √ Vatt, where Q, K, product attention mechanism, defined as: Attention(Q, K, Vatt) = softmax dh and Vatt are the query, key, and value tensors, respectively (the notation Vatt is used to avoid confusion with the vocabulary size V ). To increase modeling flexibility, this operation is performed in parallel across multiple attention heads, allowing the model to capture different types of relationships within the sequence. • Feed Forward Network (FFN): this is the second component inside each Transformer layer, which applies a transformation to each position separately, refining the representation obtained from the attention layer. At a high level, the FFN consists of two linear transformations separated by a nonlinear activation function, and can be expressed as: FFN(h) = W2 σ(W1 h + b1 ) + b2 , where h ∈ Rd is the representation of one token, W1 projects it from the hidden dimension d to an intermediate dimension dff , and W2 projects it back from dff to d. The vectors b1 and b2 are the corresponding bias terms. The function σ(·) denotes a nonlinear activation function, such as the Gaussian Error Linear Unit (GeLU)1 used in the proposed GPT-2 architecture, although the specific activation function depends on the model architecture. This operation temporarily expands the dimensionality of the representation (typically to a multiple of d) before projecting it back to the original size, allowing the model to capture more complex patterns. It is important to note that the FFN preserves the overall dimensionality of the data while further transforming each token representation before passing it to the next layer. • Output layer (linear projection and softmax): the final stage of the model transforms the output representations into predictions over the vocabulary. Each token representation, of dimension d, is projected onto the vocabulary space through a linear transformation defined by a matrix of size V × d, where V = |V|. This operation produces a vector of values known as logits, which represent unnormalized scores for each possible token in the vocabulary. A softmax function2 is then applied to convert these scores into a probability distribution over the vocabulary which are used to compute the cross entropy loss against the target token at each position. Due to the large vocabulary size, this projection is computationally and memory intensive, as it requires producing a vector of V logits for each token position. 1 GeLU is a smooth nonlinear activation function defined as GELU(x) = xΦ(x), where Φ(x) is the cumulative distribution function of the standard normal distribution. Other Transformer architectures may use different activation functions or gated variants, such as ReLU, SiLU, or SwiGLU.  PV 2 For a logit vector z ∈ RV , softmax is defined as softmax(z) = exp(z ) i i j=1 exp(zj ), producing nonnegative values that sum to one.

8

EXAMPLE

Divina Commedia

Nel mezzo del cammin di nostra vita ... TOKENIZATION

N

45

el

301

me

757

zzo

46784

del

1624

cam

6730

min

1083 ⋮

[EOS]

50256

Split into T-length sequences. T tokens

S1 S2 S3 SNM

⋮

… … …

⋮

…

Group into a mini-batch

Token Embedding Layer

Embedding X ∈ ℝB × T × d (d = embedding dimension) Input to Transformer chain

Masked Multi-Head Self-Attention

Mini-batch of size B S1 S2 S3 SNM

… … … ⋮

×L layers

Layer Norm

Feed Forward

⋮

…

Tensor shape: (B, T) (token IDs)

Layer Norm

Input: Token IDs

Embedding Layer

Output: Embedding Vectors

N

45

0.12

-0.44

...

0.87

el

301

0.03

0.91

...

-0.11

-0.21

0.07

...

0.56

0.31

-0.15

...

0.42

0.33

...

-0.77

me

757

zzo

46748

del

1624

-0.03

cam

6730

0.28

0.11

...

0.09

min ...

1083

0.17

-0.88

...

0.24

...

... |V| ...

}

d elements

d = 2048

Example (FP16): with d = 2048, each token is represented by 2048 elements × 2 bytes ≈ 4 KB

Transformer block

Figure 3: Illustration of the data flow during LLM training. Raw text from La Divina Commedia is tokenized, divided into fixed length sequences, grouped into a mini batch of size B, and transformed into the input tensor X ∈ RB×T .

Several recent works further optimize the attention mechanism itself. FlashAttention [95] improves memory efficiency through optimized attention kernels while FlashAttention 2 [96] further improves parallelism and work partitioning on the GPU. Linear attention [97] approximates attention to improve scaling with sequence length. Other variants modify how keys and values are shared across heads, including multi-query attention [98] and grouped-query attention [99]. A broader overview of efficient attention mechanisms is provided in [100]. Although these techniques modify the internal execution of the attention block, they leave the input and output tensor dimensions unchanged. Since only internal attention execution changes, we do not distinguish among these variants. The FFN block has also been the target of efficiency improvements. Sparse and MoE variants increase model capacity without increasing the active computation proportionally [101, 102, 103, 104]. As mentioned above, we discuss MoE based architectures separately in the Appendix A.4. 4.3. Inside the input pipeline As introduced in Section 3, the basic training sample processed by the model is a sequence of T token identifiers. In practice, processing a single sequence at a time is generally inefficient, as it may underutilize the parallel computational resources available on modern GPUs. Instead, multiple sequences are grouped together into a tensor (X), commonly referred to as a mini-batch, containing B sequences. The mini-batch size is chosen to maximize the throughput of the training system, i.e., the number of tokens processed per unit time. Intuitively, this corresponds to selecting the largest mini-batch that can fit within the available accelerator memory while avoiding out of memory conditions. Larger mini-batches generally improve hardware utilization, although the optimal value ultimately depends on both the model architecture and the available hardware resources. To summarize, the input tensor X provided to the Transformer during training can be represented as X ∈ {1, . . . , V }B×T , where B denotes the mini-batch size, and T the sequence length. Fig. 3 summarizes how raw text is transformed into tokenized sequences and subsequently aggregated into a mini-batch, which serves as the input to the LLM during training. 4.4. The journey of the input tensor We now follow the journey of the input tensor X through the Transformer model, focusing on the transformations it undergoes within the network. Understanding these transformations is essential for characterizing the communication and network traffic patterns discussed later in this tutorial. The first operation applied to the input tensor X is the embedding layer. This layer maps each token identifier to a dense vector representation of fixed dimension (d). In SOTA LLMs, the embedding dimension typically ranges from a few hundred to several thousand elements (e.g., d ∈ [768, 16,384]), with each element stored using numerical formats such as FP16, BF16, or FP32. As a result, each token is transformed from a single integer identifier into a high dimensional vector representation. Consequently, after the embedding layer, the input tensor assumes the form defined in (1) where each token is represented by a d dimensional embedding vector. Fig. 3 illustrates this process, showing how token identifiers are mapped to their corresponding embedding vectors. A positional encoding is then typically added to the embedding representation to encode token ordering within the sequence, while preserving the same tensor dimensions. After this transformation, the tensor X is commonly referred to as the input activation or hidden state. The input activations are then processed by the first Transformer layer. In the decoder-only architecture considered here, each layer combines masked multi-head attention with a FFN. The tensor X first enters the attention block, where token representations are updated using information from the rest of the sequence. 9

Figure 4: Tensor shape preservation: X ∈ RB×T ×d retains its shape across attention, FFN, and the Transformer chain.

The resulting tensor is then processed by the FFN, which further updates the embedding values. Throughout these operations, the structure of the output tensor remains unchanged: the number of sequences, the sequence length, and the embedding dimension are preserved. Consequently, the output of the Transformer layer is still represented as the tensor defined in Eq. (1) although the values contained in the embedding vectors have been updated. The output tensor is then passed to the next Transformer layer, where the same process is repeated. Layer after layer, the embedding vectors are progressively refined as increasingly rich contextual information is incorporated into their representations. This process continues until the final Transformer layer is reached. Throughout the entire network, the tensor retains the same shape, with only the values of its embeddings being updated. Fig. 4 summarizes this process, illustrating how the embedding values are progressively updated while the tensor shape remains unchanged throughout the Transformer chain. The first significant dimensionality change occurs only at the final projection layer. Specifically, the tensor X is projected onto the vocabulary space, producing a tensor of logits Z ∈ RB×T ×V . Each vector of length (V ) contains a score for every possible next token in the vocabulary. These scores are subsequently converted into probabilities through a softmax operation, yielding an output tensor of the same dimensions. 4.5. The LLM training step We have followed the journey of an input mini-batch through the Transformer model and observed how its dimensionality remains unchanged throughout the network. However, this journey represents only one component of the overall training process. To correctly introduce the distributed training techniques discussed in the following sections, we first review the fundamental concepts and terminology associated with LLM training. A training step can be conceptually divided into four main phases: 1. Forward pass: this phase corresponds to the process described in the previous section. The input tensor is first transformed into input activations by the embedding layer, and is subsequently processed by the sequence of Transformer layers. The final tensor representation is then projected onto the vocabulary space and used to predict the next token at each position of the sequence. 2. BW pass: Once the training loss has been computed, the model must determine how each parameter contributed to the prediction error. This information is captured by quantities known as gradients, which indicate how the loss would change if a given parameter were modified. Backpropagation computes these gradients and propagates them from the output layer to the input layer. From a data-flow perspective, the backward (BW) pass mirrors the forward (FW) pass in reverse, propagating gradients rather than activations. These gradient tensors retain the overall dimensions of the corresponding activations as they traverse the Transformer layers. 10

3. Gradient accumulation: a mini-batch represents the amount of data that can be processed by the model at a given time and is typically constrained by the available accelerator memory. However, model parameters are not necessarily updated after every mini-batch. Instead, multiple mini-batches are often processed before completing a training step. The set of tokens processed before a parameter update is commonly referred to as the batch (or global batch): LLM training runs typically use global batch sizes of roughly 106 to 107 tokens per training step [105]. Since processing an entire global batch simultaneously is often infeasible, it is divided into multiple mini-batches, which are processed sequentially. The gradients computed during the BW pass of each mini-batch are accumulated, typically through summation, and are only used to update the model parameters once all mini-batches belonging to the global batch have been processed. 4. Weight update: After computing the gradients, the model parameters are updated using an optimization algorithm. The new parameters are obtained by adjusting the current weights based on the computed gradients, moving them in a direction that reduces the loss and improves the model’s predictions. This update completes a training iteration, commonly referred to as a training step. 4.6. Batch Scale and Single GPU Training To illustrate the scale of batches used in LLM training, consider La Divina Commedia, whose approximately 100,000 words [86] correspond to roughly (125,000) tokens. For comparison, a model such as GPT-2 is commonly trained using batch sizes on the order of (5×105 ) tokens. In this case, the entire Divina Commedia would occupy only about one quarter of a single training batch. Considering instead GPT-3 175B, which uses a batch size of 125,000 approximately (3.2 × 106 ) tokens, the complete Divina Commedia would represent only 3.2×10 6 × 100 ≈ 3.9% of a single batch. Equivalently, approximately 26 copies of the entire poem would be required to fill one GPT-3 training batch. Let us now follow a concrete mini-batch through the training pipeline and examine the dimensions of the input tensor X and the outputs generated at each stage. We base this example on a modified version of the implementation of the GPT-2 model provided in [106], and we report the tensor dimensions and throughput observed in the single GPU setup. Example 4.1 (Processing a mini-batch on a single GPU). Consider running the GPT-2 model implementation with sequence length T = 1024 and hidden dimension d = 768, where activations are stored using FP16 precision (2 bytes per value), running on an NVIDIA L40S accelerator equipped with 48 GB of VRAM. A single input sequence of length 1024 is first transformed by the embedding layer into a tensor containing 1024 × 768 = 786,432 FP16 values. The resulting activation tensor therefore occupies 1024 × 768 × 2 = 1,572,864 bytes ≈ 1.5 MB. With this setup, processing a single sequence achieves a throughput of approximately 65,000 tokens/s. Loading the model and processing the sequence requires about 2 GB of GPU memory. The available GPU memory therefore allows multiple sequences to be processed simultaneously. In this setup, the maximum mini-batch size that could be accommodated was B = 32 sequences. Consequently, the input tensor before the embedding layer contains 32 × 1024 = 32,768 tokens . After the embedding layer, the tensor becomes X ∈ R32×1024×768 with a corresponding activation memory of 32 × 1024 × 768 × 2 = 50,331,648 bytes ≈ 50 MB. On the same setup, this configuration achieves a throughput of approximately 105,000 tokens/s. Assume now a global batch size of 524,288 tokens. Since each mini-batch contains 32 × 1024 = 32,768 tokens, the number of mini-batches that must be processed before performing a weight update is 524,288 32,768 = 16 mini-batches. Therefore, a complete training step consists of processing approximately 16 mini-batches. The gradients generated during the BW pass of each mini-batch are accumulated, and only after all 16 mini-batches have been processed, the model parameters are updated. The reason why these experiments were conducted using GPT-2 is that it is small enough to perform a training step on a single L40S accelerator. Models with a parameter count comparable to GPT-3 175B require substantially more memory than is available on a single GPU, making distributed training a necessity. This challenge is discussed in the next section. 5. Distributed training: where traffic originates So far, we have analyzed LLM training in a single device setting, following how the input tensor X is processed through the model. In this section, we extend this view to distributed training. We first discuss why parallelization is required at the scale of large LLMs, and then analyze how different parallelization strategies partition the training workload across accelerators, generating the communication patterns and network traffic studied throughout this tutorial.

11

5.1. The need for distributed training As shown in Example 4.1, our single GPU setup is able to process approximately 105,000 tokens/s when training GPT-2. Assuming a training dataset containing 3 × 109 tokens, the total training time would be 3×109 ≈ 28,500 seconds, corresponding to roughly 8 hours of training on a single accelerator. approximately 105,000 To understand why distributed training becomes necessary, let us progressively scale this example. GPT-2 contains approximately 124-128 million parameters, whereas larger models such as GPT-3 XL already contain about 1.3 billion parameters [15], more than 10× larger. As a simplified approximation, let us assume that the achievable throughput decreases proportionally with model size. Under this assumption, the processing throughput would decrease from approximately 105 to 104 tokens/s. Now consider training such a model 12 8 on a dataset containing 1012 tokens. The required training time would become 10 104 = 10 seconds, which corresponds to more than 3 years of continuous training on a single GPU. Although highly simplified and based on several assumptions, this example highlights a fundamental challenge of modern LLM training: as model size and dataset size increase, training on a single accelerator rapidly becomes computationally impractical. Computation time, however, is only part of the problem: memory capacity quickly becomes an even more severe limitation. A simple estimate of the memory required to store a model can be obtained by multiplying the number of parameters by the storage size of each parameter. For example, assuming FP16 precision (2 bytes per parameter), a model containing 100 billion parameters requires approximately 200 GB of memory simply to store the parameters. During training, the memory requirements increase substantially due to the need to store gradients, optimizer states, and intermediate activations. In practice, the total memory footprint is often 4× to 10× larger than the model size alone. For instance, the activation memory required during GPT-3 training exceeds the memory occupied by the model parameters by more than 5× [12]. Consequently, training modern LLMs on a single accelerator becomes infeasible for two main reasons: 1. Memory limitations: the total memory required during training, including model parameters, activations, gradients, and optimizer states, exceeds the capacity of a single device. While current high end accelerators provide tens or, in some cases, hundreds of gigabytes of memory (e.g., NVIDIA H100 [107]), these capacities remain insufficient for training the largest models. 2. Computational time: the amount of computation required to train models containing hundreds of billions or trillions of parameters, on datasets composed of trillions of tokens, would result in training times measured in years on a single accelerator. Despite continuous efforts by hardware vendors to increase the memory capacity available per accelerator, memory growth has not kept pace with the rapid scaling of LLMs. As reported in [10], model sizes have grown faster than the memory capacity of individual accelerators, creating an increasingly wide gap. Consequently, modern LLM training relies on distributed techniques that partition computation and memory across multiple accelerators. 5.2. Parallelization Techniques and Communication Primitives The need to overcome the memory and computational limitations has led to the development of several distributed training strategies. These techniques differ in how the training workload is partitioned across multiple accelerators and, consequently, in the communication patterns they generate. Since network traffic is ultimately a consequence of how data, computations, and model parameters are distributed, understanding these parallelization strategies is essential for estimating the volume of data exchanged during training. In this section, we introduce the most widely adopted parallelization techniques and provide the background needed to derive their communication volumes. These techniques can be distinguished by the component of the training workload that they distribute across accelerators. The main strategies considered in this tutorial are: • Data Parallelism (DP): the full model is replicated across multiple devices, and each device processes a different portion of the input data. The results are then combined to update a shared set of model parameters [41, 42, 43]. • Pipeline Parallelism (PP): the model is divided into sequential stages, each assigned to a different device, enabling different parts of the model to be executed concurrently on different portions of the input [44, 45, 46, 47]. • Tensor Parallelism (TP): the computations within individual layers are split across multiple devices, allowing large model components to be distributed and processed in parallel [40]. Sequence Parallelism (SP) extends TP by partitioning the input sequence across devices [108] and is discussed in Appendix A.2. • Hybrid Parallelism: DP, TP, SP, and PP are in practice often combined to enable efficient training of large scale LLMs [12, 41]. 12

Figure 5: DP workflow. Each GPU maintains a replica of the model and processes a different mini-batch locally. After the FW and BW passes, the locally computed gradients are exchanged and aggregated across all GPUs through an AllReduce CCO.

To establish a clear baseline, we consider simplified implementations of these strategies. Appendix A.2 discusses common optimizations and how different implementations alter the resulting traffic patterns. Duan et al. [25] and Rostam et al. [109] review the distribution and optimization of LLM workloads across large scale systems, covering parallelization strategies, associated challenges, performance considerations, and the physical data center infrastructure (hardware connections, interconnects) used for training. MoE models introduce Expert Parallelism (EP) [101, 103], which we discuss in Appendix A.4. Before examining the parallelization strategies, we introduce the required communication terminology. Distributed training uses Collective Communication Operations (CCOs) to aggregate or redistribute tensors across accelerators, typically through Collective Communication Libraries (CCLs). Their algorithms and implementations are discussed in Section 6. The main collective considered here is AllReduce, which aggregates a tensor across a group of accelerators and returns the result to every participant. In this section, we treat it as a logical synchronization operation and focus on the payload, participants, and invocation frequency. We now analyze the main parallelization strategies. Unless stated otherwise, the numerical examples in this section use the GPT-2-like configuration introduced previously, derived from the nanoGPT implementation [106]. The model has P ≈ 128M parameters, L = 12 Transformer layers, hidden dimension d = 768, and sequence length T = 1024. All assumptions used in the examples are based on the performance characteristics of an NVIDIA L40s GPU equipped with 48 GB of memory. The reported execution times were obtained from experiments conducted on the same hardware. Communicated activations and gradients use FP16/BF16 precision, such that sact = sgrad = 2 bytes per value. 5.2.1. Data Parallelism (DP) To simplify the discussion, we begin by assuming that the entire model can fit within the memory of a single accelerator. The idea behind DP is to train several identical copies of the same model in parallel. Each accelerator stores a full replica of the model, but processes a different portion of the training data. The FW and BW passes are performed independently on each accelerator, producing local gradients from the data processed on that device. When a parameter update is required, these gradients are synchronized across all accelerators in the DP group, typically through an AllReduce CCO. After synchronization, each accelerator applies the same weight update, ensuring that all model replicas remain identical. In this way, DP increases the amount of training data processed in parallel while keeping all copies of the model synchronized. Fig. 5 illustrates this execution flow, highlighting the local computation performed on each GPU and the subsequent gradient synchronization phase. From a communication perspective, the total number of gradient values is equal to the number of model parameters. Consequently, whenever gradients are synchronized, each accelerator participates in the synchronization of a gradient tensor whose size is approximately equal to the model size. Given P and sgrad , the logical gradient tensor synchronized by each DP worker has size MDP defined in (3b). For FP16/BF16 gradient communication, sgrad = 2 bytes, and thus MDP = P × 2 bytes. In standard DP, gradient synchronization is usually implemented through an AllReduce operation. The following example applies the reference configuration to eight DP workers. Example 5.1 (Gradient Synchronization in DP). Using the reference configuration, consider DP across pDP = 8 GPUs. Each GPU stores a full model replica and processes a local mini-batch of B = 32 sequences. The target global batch size is Bglobal = 524,288 tokens. The number of gradient accumulation steps per GPU is B 524,288 Nacc = BTglobal pDP = 32·1024·8 = 2. Thus, before each optimizer update, every GPU processes two local mini-batches, 13

Figure 6: PP workflow. Transformer layers are partitioned across multiple GPUs, while the activation tensor X is exchanged between adjacent pipeline stages during the FW and BW passes through point-to-point communication.

corresponding to 2 · 32 · 1024 = 65,536 tokens per GPU. Across all eight GPUs, this gives the target global batch size: 8 · 65,536 = 524,288 tokens. After local gradient accumulation, the DP workers synchronize their gradients. The logical gradient tensor synchronized by each worker has size MDP = P sgrad = 256 × 106 bytes ≈ 256 MB. Therefore, each GPU participates in one AllReduce over an approximately 256 MB logical payload per optimizer update. The communication volume generated by the collective implementation is derived in Section 6. This example highlights the main communication characteristic of DP: communication is relatively infrequent, since it occurs once per optimizer update, but the synchronized object is proportional to the model size. Therefore, DP becomes network-intensive when the model is large, when many DP workers participate in the synchronization, or when optimizer updates occur frequently. 5.2.2. Pipeline Parallelism (PP) PP represents a second widely adopted parallelization strategy for LLM training. In this case, we start from the assumption that the model is too large to fit within the memory of a single accelerator. The idea behind PP is to split the Transformer stack across multiple accelerators, so that each device stores and executes only a consecutive subset of layers. For example, a model with 96 Transformer layers, such as GPT-3 [15], can be distributed across 8 GPUs, with each GPU storing 12 layers instead of the full model. During the FW pass, the input is first processed by the accelerator hosting the initial layers of the model. The resulting activation tensor is then sent to the next accelerator, which processes the following layers. This process continues stage by stage until the final accelerator produces the model output. During the BW pass, gradients follow the same pipeline in the opposite direction, moving from the last stage back toward the first. From a communication perspective, PP differs from DP because data is not synchronized across all accelerators. Instead, communication mainly consists of P2P exchanges between adjacent pipeline stages. For each processed input, the activation tensor is transmitted FW across a pipeline boundary, while the corresponding activation gradient tensor is transmitted BW across the same boundary. These activation and activation gradient exchanges therefore represent the primary source of network traffic in PP training. Fig. 6 illustrates the execution flow of PP, highlighting how the activation tensor propagates across adjacent pipeline stages during the FW and BW passes. For a minibatch of size B, sequence length T , hidden dimension d, and activation precision sact bytes, the activation tensor exchanged between two adjacent pipeline stages has the shape in (1) and the payload MPP defined in (2). Since the activation tensor is transmitted once during the FW pass and its corresponding gradient once during the BW pass, the communication volume across one pipeline boundary is CPP,boundary ≈ 2MPP = 2Sact where Sact is the activation payload defined in (2). This volume is local to each pipeline boundary: unlike DP, the tensor is not synchronized across all accelerators, but exchanged only between neighboring pipeline stages. The following example instantiates this calculation for a GPT-2-like model partitioned across eight pipeline stages. The example deliberately uses a naive non-overlapped pipeline execution in order to expose the size and timing of the activation transfers before introducing bubble-reducing schedules. Example 5.2 (Activation Transfers in PP). Using the reference configuration, consider pPP = 8 pipeline stages processing a mini-batch of B = 32 sequences. We consider a naive pipeline partition where the 12 Transformer 14

Figure 7: TP workflow: Transformer blocks span multiple GPUs, with CCOs exchanging X during FW and BW.

layers are distributed across the 8 GPUs as [1, 1, 2, 2, 2, 2, 1, 1]. Thus, the first two and last two stages contain one Transformer layer each, while the four middle stages contain two layers each. This partition is deliberately unbalanced to account for the additional input and output vocabulary layers assigned to the first and last stages, which can create both computational and memory imbalances [110]. Assume the following per-layer layer execution times: tlayer FW ≈ 20 ms, and tBW ≈ 40 ms. The FW pass execution times of the eight stages are therefore [20, 20, 40, 40, 40, 40, 20, 20] ms, which gives a total naive FW pipeline latency of 20 + 20 + 40 + 40 + 40 + 40 + 20 + 20 = 240 ms. Similarly, the naive BW pipeline latency is approximately 2 × 240 = 480 ms. At each pipeline boundary, adjacent stages exchange the activation tensor defined in (1). Its logical payload is MPP = Sact = 32 · 1024 · 768 · 2 ≈ 50 MB where Sact is defined in Eq. (2). Each pipeline boundary therefore carries approximately 50 MB during the FW pass and the same amount of activation-gradient data during the BW pass, for a total of approximately 100 MB per mini-batch. In this naive non-overlapped execution, a fullmini-batch activation transfer therefore occurs across each pipeline boundary once during the approximately 240 ms FW traversal, and the corresponding activation-gradient transfer occurs once during the approximately 480 ms BW traversal. This example highlights the main communication characteristic of PP: communication is activation-sized and local, as tensors are exchanged only between adjacent pipeline stages. However, these transfers lie on the critical path, and their dependencies may leave stages idle, creating pipeline bubbles. Practical PP implementations divide each mini-batch into m micro-batches of size bµ , such that B = mbµ . Each micro-batch is propagated independently through the pipeline, and the tensor exchanged between adjacent stages is Xi ∈ Rbµ ×T ×d , i ∈ {1, . . . , m}. Assuming equal precision for activations and their gradients, the communication volume per microPm batch and boundary is 2bµ T d sact . Across all micro-batches, the total remains i=1 2bµ T d sact = 2Sact . Microbatching therefore changes the granularity and timing of transfers rather than their total volume. Scheduling techniques that exploit this decomposition are introduced in Appendix A.2. 5.2.3. Tensor Parallelism (TP) As in the case of pipeline parallelism, we assume that the model is too large to fit within the memory of a single accelerator. Unlike PP, however, TP partitions the computation occurring inside individual Transformer layers across multiple GPUs. In particular, the matrix multiplications performed in the self-attention and FFN blocks are distributed across several accelerators, each of which computes only a portion of the operation. The detailed partitioning schemes are described in the original Megatron-LM work [40], while an intuitive pedagogical walkthrough of these operations is also available online [111]. In this tutorial, we focus primarily on the resulting communication patterns. Compared to DP and PP, TP introduces significantly more frequent communication, since synchronization occurs inside every Transformer layer. As discussed in [40], tensor-parallel execution typically requires two AllReduce CCOs per Transformer layer during the FW pass and two AllReduce CCOs during the BW pass. Fig. 7 illustrates the execution flow of TP, highlighting the CCOs performed inside each Transformer block during the FW and BW passes. From a communication perspective, the tensor exchanged by 15

Table 3: Distributed Training Communication Patterns

Technique

Communication Pattern

DP

Model sized gradient synchronization with payload Sgrad defined in (3b) across data parallel replicas, typically through AllReduce. Activation and activation-gradient tensors exchanged through P2P communication between adjacent pipeline stages with activation tensor shape defined in (1). Multiple CCOs inside each Transformer layer during both the FW and BW passes. The exchanged tensor has the shape defined in (1). SP, a TP optimization, follows the same general structure as TP, but with different tensor partitioning and execution scheduling. Combination, placement, and overlap of the communication patterns generated by DP, PP, TP and SP.

PP TP Hybrid

TP is associated with the activation tensor processed by the Transformer layer. The activation tensor involved in TP communication has the shape given in (1) and the payload defined in (2). Since TP typically requires two AllReduce operations in the FW pass and two AllReduce operations in the BW pass for each Transformer layer, the logical amount of activation data synchronized per layer is approximately 4MTP . The following example instantiates this calculation for a GPT-2-like model trained with TP across eight GPUs. Example 5.3 (AllReduce Payload in TP). Using the reference configuration, consider TP across pTP = 8 GPUs processing a mini-batch of B = 32 sequences. The communicated activation tensor has the shape defined in (1) and the payload given in (2). Substituting the above values gives MTP = Sact = 32·1024·768·2 ≈ 50 MB. Assume that the execution of a single Transformer layer is distributed across eight GPUs and requires approximately tFW ≈ 3 ms, and tBW ≈ 6 ms. In the baseline non optimized TP implementation, each Transformer layer performs two AllReduce operations during the FW pass and two during the BW pass, each over a logical payload of approximately 50.3 MB. Under the assumed execution times, this corresponds to one AllReduce every 1.5 ms during the FW pass and every 3 ms during the BW pass. This example highlights the main communication characteristic of TP: the communicated tensor is activationsized, as in PP, but the communication is much more frequent because it occurs multiple times inside every Transformer layer. Therefore, TP is sensitive not only to bandwidth, but also to communication latency and synchronization overhead. The main limitation of TP is that it is highly communication intensive, requiring synchronization operations much more frequently than DP or PP. For this reason, TP strongly depends on extremely high-bandwidth and low-latency interconnects, and becomes significantly less efficient when accelerators are distributed across different nodes without ultra-fast communication links. Moreover, compared to DP and PP, TP is generally more complex to implement and adapt to evolving Transformer architectures. 5.2.4. Hybrid Parallelism and communication summary In practical large-scale LLM training, parallelization strategies are rarely used in isolation. Instead, modern systems combine multiple dimensions of parallelism in order to jointly address model memory constraints, computational scalability, and communication efficiency. This combination of multiple strategies is commonly referred to as hybrid parallelism. In this setting, 3D parallelism usually denotes the joint use of DP, TP, and PP, while 4D parallelism extends this design by introducing an additional dimension, often SP or another context-related partitioning strategy [12, 108, 112]. Hybrid parallelism has become the dominant paradigm for large-scale LLM training. It is employed by many widely adopted systems and frameworks, including MegatronLM [12], DeepSpeed [41], Mesh-TensorFlow [48], and production-scale training systems such as PaLM [75]. Although these systems differ in their implementation strategies, they all rely on combining multiple forms of parallelism across different groups of accelerators. More comprehensive overviews of distributed training strategies and optimization techniques are provided in [25, 113]. From a communication perspective, hybrid parallelism should not be viewed as a completely new communication mechanism. Instead, it combines and overlaps the traffic patterns generated by its constituent strategies. Communication partitioning and hierarchical scheduling can further improve the overlap between these exchanges and model computation [114]. DP introduces model-sized gradient synchronization across model replicas, PP generates point-to-point exchanges of activation tensors between adjacent pipeline stages, and TP requires frequent collective communication inside Transformer layers. Consequently, the simplified formulations introduced in the previous subsections remain the fundamental building blocks for characterizing network traffic also in hybrid training configurations as well. Table 3 summarizes these communication patterns and

16

Execution 1 Backward pass (480 ms)

Forward pass (240 ms)

50 MB 50 MB every 1.5 ms

Execution 2 Backward pass (480 ms)

Forward pass (240 ms)

50 MB every 3 ms

50 MB 50 MB every 1.5 ms

TP

50 MB every 3 ms

TP

0 50 MB 7 × 50 MB

0 50 MB 7 × 50 MB

7 × 50 MB

PP

PP

0 256 MB

0 256 MB

DP

No DP communication 50 0

0

60 120 180 240 0 Time in FW pass (ms)

1 × 256 MB

DP

No DP communication 120 240 360 480 Time in BW pass (ms)

7 × 50 MB

256 MB

No DP communication 50 0

0

60 120 180 240 0 Time in FW pass (ms)

120 240 360 480 Time in BW pass (ms)

Figure 8: Communication schedule for a full batch execution, composed of two consecutive executions and showing TP, PP, and DP transfers during the forward and backward passes.

Fig. 8 shows how the traffic patterns from the previous examples combine in a simple 3D configuration. TP generates frequent 50 MB AllReduce operations, PP exchanges 50 MB activation tensors at stage boundaries via simple P2P operations, and DP performs one 256 MB AllReduce after two gradient accumulation steps. The figure makes it clear that the three parallelization strategies do not create a new traffic pattern when combined. Instead, the traffic generated by DP, PP, and TP coexists within the same training execution, with each strategy contributing its own communication pattern. The timeline is illustrative and compares communication frequency and payload size; it does not represent a measured end-to-end trace of a fully integrated hybrid training configuration. This traffic pattern is consistent with several studies of network behavior during LLM training [10, 16, 18, 67, 115]. Gangidi et al. [18] describe it as predictable and repetitive, with few connections and millisecond scale on off bursts. During each active burst, large tensors are transferred as quickly as the network permits, so the instantaneous data rate can approach the available link bandwidth. However, the average rate over the full training step is lower because communication bursts alternate with periods of computation in which little or no data is transferred. Selecting a specific hybrid parallelization strategy depends on several factors, including model size, sequence length, accelerator memory capacity, and the characteristics of the underlying network topology and infrastructure discussed in Section 7. The input length distribution also matters, since variable length documents can create uneven computation and communication loads across pipeline groups [116]. For this reason, several works have explored automatic approaches for deriving efficient hybrid configurations, including Alpa [49], FlexFlow [117], GSPMD [50], Unity [118], Galvatron [119], and Aceso [120]. Nevertheless, determining the optimal hybrid parallelization strategy for a given training scenario remains an open research problem. Despite this, a commonly adopted design pattern is to map TP groups onto accelerators connected through extremely high-bandwidth and low-latency links, typically within the same node. Conversely, PP and DP are generally more suitable for communication across nodes, where network performance is comparatively lower [11, 121]. Fig. 9 summarizes this placement intuition and the resulting communication patterns. 6. CCOs: from tensor to communication payloads In the previous section, we identified the tensor payloads that are collectively communicated under different distributed training strategies. In the Life of a Token framework, this marks the transition from tensors as computational objects to tensors as communication objects. This section examines how such tensors are processed by CCOs and how the implementation of these operations affects the resulting network traffic. This analysis matters because a tensor size alone is not sufficient to determine communication volume: the same CCO can be implemented through different algorithms, each inducing a different communication schedule, number of transfers, and traffic distribution across links. To illustrate how collective algorithms move, combine, and redistribute data, we focus on AllReduce, one of the most widely used CCOs in our scenario. As discussed in the previous section, it plays a central role in gradient synchronization during large-scale LLM training. It is also a useful tutorial example because, depending on its implementation, it can be decomposed into simpler 17

3D PARALLELISM: COMMUNICATION VIEW Combination of Data Parallelism (DP) × Pipeline Parallelism (PP) × Tensor Parallelism (TP) Pipeline Parallelism (PP) Stages Stage 0 Replica 0

GPU

GPU

Stage 1

…

COMMUNICATION DIMENSIONS Stage P−1

Tensor Parallelism (TP) Collective communication (e.g., all-reduce) inside each layer.

GPU

X ∈ ℝB × T × d

Data Parallelism (DP)

Pipeline Parallelism (PP) Replica 1

Point-to-point communication of activations (forward) and gradients (backward).

Replicas

…

X ∈ ℝB × T × d

…

Data Parallelism (DP) Gradient synchronization across replicas after the backward pass (e.g., all-reduce).

Replica D−1

Tensor Parallel Group

Tensor Parallel Group

X ∈ ℝP × Sgrad

Figure 9: Example of 3D hybrid parallelism combining DP, PP, and TP. GPUs are organized in a 3D grid; different parallel dimensions generate different communication patterns.

operations, such as ReduceScatter, AllGather, Reduce, and Broadcast, allowing us to introduce multiple CCOs through a single representative operation. This section is therefore not intended as a complete survey of CCOs or CCLs. Instead, it provides a conceptual bridge for characterizing how the tensor payloads generated during LLM training translate into the network traffic induced by their collective communication. A broader catalogue of collective operations is available in the NCCL documentation [55]. 6.1. AllReduce implementation strategies We now examine how different implementations of AllReduce translate the same logical tensor payload into different communication volumes. The key question is how the reduction is scheduled, since the schedule determines the number of transfers and the resulting traffic pattern. To make this effect explicit, we compare three representative implementations: a naive AllReduce, a ring-based AllReduce, and a tree-based AllReduce. The naive version provides a simple baseline, while the ring and tree variants illustrate two widely used design principles for organizing collective communication. 6.1.1. Naive AllReduce A simple but inefficient implementation of AllReduce consists of having each GPU broadcast its local tensor to all the GPUs participating in the operation. After receiving the tensors from every participant, each accelerator locally computes the final reduction. Although intuitive, this method generates a very large amount of network traffic and scales poorly as the number of accelerators increases. In this naive implementation, each of the pcc GPUs sends a message of size M to all other pcc − 1 GPUs. Therefore, the traffic generated by each naive naive naive GPU is: Cdev (M, pcc ) = (pcc − 1)M. The total traffic across the system is: Ctotal = pcc Cdev = pcc (pcc − 1)M. 6.1.2. Ring AllReduce Optimized AllReduce implementations avoid the naive broadcast based behavior. The key idea is to arrange the pcc participating GPUs in a logical ring and split the exchanged tensor M into pcc equally sized chunks. Instead of sending the entire tensor to every other GPU, each accelerator exchanges one chunk at a time with its neighbors in the ring. Ring AllReduce is typically implemented in two phases. The first phase is the ReduceScatter phase: tensor chunks circulate around the ring, and partial reductions are accumulated until each GPU holds one fully reduced chunk. The second phase is AllGather : the reduced chunks circulate again so that, at the end of the operation, every GPU reconstructs the complete reduced tensor. This implementation is bandwidth efficient for large tensors because each GPU communicates only with its two ring neighbors and transfers data in a regular chunked pattern. As a result, communication can be pipelined across the ring, and each P2P exchange can use the available bandwidth of the corresponding link, assuming no external contention. This makes ring AllReduce a widely adopted implementation. With this approach, the amount of data sent by ring −1 each GPU is approximately: CGPU = 2 · pcc pcc · M. Consequently, the total amount of data transmitted across ring ring the system is: Ctotal = pcc · CGPU = 2(pcc − 1)M. Fig. 10a summarizes the execution of ring AllReduce across four GPUs.

18

Phase 1: Reduction Phase 1: Reduce-Scatter

Phase 2: All-Gather

GPU 1

GPU 1

1234 GPU 4

1234

GPU 4

GPU 2 1234

GPU 2

1 2 3 4

1 2 3 4

GPU 3

GPU 3

1234

1 2 3 4

GPU 1

(root) 1 2 3 4

1 2 3 4

GPU 2 1234

Phase 2: Broadcast

GPU 1

GPU 4 1234

(a) Example of Ring AllReduce execution across four GPUs

(root) 1 2 3 4

GPU 3 1234

GPU 2 1234

GPU 3 1234

GPU 4 1234

(b) Example of Tree AllReduce execution across four GPUs

Figure 10: Comparison of Ring and Tree AllReduce for four GPUs and their two communication phases.

6.1.3. Tree based AllReduce Tree-based AllReduce provides an alternative implementation that is often better suited for smaller tensors or latency sensitive communication. Instead of arranging GPUs in a ring, the participating accelerators are organized into a logical tree. Communication then proceeds hierarchically along the tree structure. As in ring AllReduce, the operation can be interpreted as two phases. In the first phase, reduce, data moves from the leaves of the tree toward the root. At each internal node, partial tensors received from children are combined with the local tensor, until the root obtains the fully reduced result. In the second phase, broadcast, the reduced tensor is propagated from the root back down the tree so that every GPU receives the final result. For a payload of size M , each edge of the tree carries approximately M bytes during the reduce phase and M bytes during the broadcast phase. Since a tree with pcc GPUs contains pcc − 1 edges, the total amount of data transmitted across tree the system is: Ctotal = 2(pcc − 1)M. From the perspective of a single GPU, the amount of data sent depends on its position in the tree. Leaf nodes send data during the reduce phase and receive data during the broadcast phase, while internal nodes may send and receive data from multiple neighbors. Therefore, tree based AllReduce is often better characterized either by its total system traffic or by its number of communication steps. Averaged C tree −1 tree = 2 pcc = ptotal over all GPUs, the communication volume per GPU is: CGPU,avg pcc M. Assuming a balanced cc binary tree, the number of communication steps is approximately 2 log2 pcc . Tree-based AllReduce is especially useful when the message size M is small or when latency dominates communication time. In this regime, the number of sequential communication steps becomes more important than the total amount of data transferred. A balanced tree completes the reduce and broadcast phases in approximately 2 log2 pcc steps, whereas ring AllReduce requires a number of steps that grows linearly with pcc . For this reason, tree-based algorithms can reduce synchronization latency for small messages, while ring-based algorithms are typically preferred for large tensors where bandwidth efficiency dominates. Fig. 10b illustrates the execution of Tree AllReduce across four GPUs. 6.2. Collective Communication Libraries Parallelization strategies express their required tensor exchanges through CCOs, which define the logical communication among accelerator groups without specifying its hardware implementation. CCLs provide this implementation. A CCL takes a logical CCO and maps it to an executable communication schedule across the available devices. In doing so, it selects or instantiates a collective algorithm, such as a ring, tree, or specific topology aware schemes, divides tensors into chunks, organizes the order of transfers, and coordinates the synchronization among ranks. In practice, the CCL decides the communication schedule followed by the data: which devices talk to each other, which transfers stay inside a node, and which ones cross the network between nodes. The underlying network then carries these transfers over the available interconnects and routes. This distinction is important for traffic analysis. The same tensor payload and the same CCO can lead to different communication patterns depending on how the operation is implemented. For example, an AllReduce based on a ring algorithm results in a sequence of pairwise exchanges, while a tree-based implementation organizes communication in a hierarchical manner. In addition, topology aware implementations can decide whether to keep communication within a node using faster local paths or to route data across nodes when needed, depending on the system layout. These choices affect how data is partitioned, scheduled, and exchanged across devices. Widely used CCLs include NVIDIA NCCL [54], which is commonly used for GPU-based distributed training, as well as RCCL [122] for AMD GPUs, HCCL for Huawei accelerators, PyTorch Gloo [123], and Intel oneCCL [124]. Recent work has analyzed the behavior of these libraries and explored improved collective implementations through topology aware algorithm design, optimized scheduling, and more flexible communication backends [17, 56, 57, 58, 60, 59, 125]. Low level CCL parameters can also be tuned automatically for the hardware configuration and workload, including the interference between communication and computation [126]. In the context of this 19

paper, CCLs are relevant because they form the bridge between payload-level estimates and the traffic eventually observed by the network. They determine how abstract communication primitives are decomposed into concrete transfers, making them essential for connecting LLM training tensors to link-level communication behavior. 6.3. Communication Estimates in Practice The communication volume generated by a CCO depends on both the tensor payload and the algorithm used to implement the operation. Although this tutorial uses AllReduce as a representative example, the same reasoning applies to other CCOs, whose implementations may follow different schedules and therefore generate different traffic patterns [56, 127]. Although Ring and Tree AllReduce move approximately the same total amount of data, they divide this transfer into a different number of communication rounds. Each round introduces a fixed startup cost, including operation launch, synchronization, and message-processing overhead, even when the message itself is very small. Ring AllReduce requires approximately 2(pcc − 1) sequential rounds, because both its ReduceScatter and AllGather phases proceed around the ring. For small tensors, the time required to transmit the data is short, so these repeated startup costs dominate the overall execution time. A tree reduces the number of sequential rounds to approximately O(log pcc ), allowing small messages to complete with fewer synchronization delays. For large tensors, the data-transfer time becomes more important than the startup cost. Ring AllReduce divides the tensor into chunks and pipelines them around the ring, allowing all accelerators to send and receive data simultaneously. Its traffic is also distributed evenly across the participants, which helps sustain high link utilization and approach the available bandwidth. In a tree, communication is concentrated along the tree levels, and links or accelerators closer to the root may become bottlenecks. Treebased implementations may therefore complete fewer rounds but achieve lower sustained bandwidth for large tensors. Consequently, Tree AllReduce is often preferred for small or latency-sensitive messages, whereas Ring AllReduce is usually better suited for large, bandwidth-dominated transfers [127, 56]. Although payload size and collective algorithm determine the baseline communication demand, completion time can also be affected by stragglers and network congestion, particularly in shared environments [128]. 7. Network infrastructure: from payloads to bits on the wire Network infrastructure and interconnect technologies form the fifth element of the proposed framework, linking payload-level communication requirements to the actual movement of data across the training system. The previous stages determine the size of the tensors to be exchanged, the logical communication patterns induced by the adopted parallelization strategy, and the communication volume shaped by CCO algorithms and their implementation inside CCLs. Once these logical exchanges are scheduled, their actual behavior depends on how they are mapped onto the underlying infrastructure. This mapping is influenced by CCL decisions, such as collective algorithm selection, channel scheduling, and topology-aware organization, but is ultimately constrained by the physical and logical arrangement of links among accelerators and nodes. The analysis so far can therefore be summarized as a sequence of connected steps. The model and tokenized input determine quantities such as the parameter count, batch size, sequence length, and activation dimensions. These quantities define the gradient and activation payloads. The parallelization strategy then determines which payloads are exchanged, how often communication occurs, and which accelerators participate. CCO algorithms and CCL implementations further determine how these logical exchanges are divided and scheduled across the network. Because these strategies produce traffic with different sizes, frequencies, and synchronization requirements, they also place different demands on the network. Table 4 summarizes this connection and provides a roadmap for the infrastructure discussion that follows. The table provides the link between the communication analysis of the previous sections and the infrastructure discussion that follows. Because these communication patterns place different demands on the network, this section examines the technologies and design choices available to meet them. We discuss interconnects, protocols, topologies, and routing mechanisms, and relate their bandwidth, latency, and scaling properties to the needs of different distributed training configurations. The goal is to provide practical guidance for matching available network solutions to the communication requirements of DP, PP, TP, and hybrid parallelism, while accounting for the model and input properties that determine the size and frequency of the exchanged tensors. 7.1. Infrastructure hierarchy in AI data centers Large scale LLM training is typically deployed within highly integrated AI data centers, where thousands of accelerators can be coordinated to behave as a single training system [129]. From the perspective of this tutorial, the data center provides the physical substrate over which the communication objects identified in the previous sections are transported. Before discussing specific interconnect technologies and protocols, it is useful to distinguish the main layers of this infrastructure and the types of traffic they carry. 20

Table 4: Relationship between parallelization strategies, communication patterns, and network infrastructure requirements. Strategy

Dominant communication

Traffic timing

Main ment

DP

Model-sized gradient synchronization, typically through AllReduce.

Once per optimizer update, after any gradientaccumulation steps.

High-bandwidth and scalable collective execution for large payloads.

PP

Activation and activationgradient P2P transfers between adjacent pipeline stages. Frequent activation-sized collectives within Transformer layers.

Once per mini/micro-batch at each pipeline boundary in both directions.

Combination of the DP, PP, and TP/SP communication patterns.

The different patterns coexist and may overlap during training.

Low latency and enough bandwidth to transfer activations without delaying the next stage. Very high-bandwidth and low-latency because communication is frequent and tightly synchronized. Topology-aware placement, balanced link utilization, and contention control.

TP/SP

Hybrid

Multiple times per layer during both the FW and BW passes.

network

require-

Typical placement May span multiple nodes, often using hierarchical intranode and inter-node collectives. Adjacent stages are preferably placed close together or connected through highcapacity paths. Usually kept within a node or another high-bandwidth communication domain. TP is generally kept local, PP connects consecutive stages, and DP spans replicated model or pipeline groups.

A first distinction is between frontend and backend network traffic. The frontend network carries auxiliary traffic associated with the operation of the training system, such as dataset ingestion and checkpoint storage [130, 131], together with logging, monitoring, and management access. This traffic is often described as north-south traffic, since it connects the training cluster with storage systems, external services, or control plane components. In contrast, the backend network carries the communication generated by the training process itself. This traffic is commonly referred to as east-west traffic, since it flows among accelerators and nodes participating in the same distributed training job [67, 18]. For large scale LLM training, backend east-west traffic is the dominant concern for performance, as communication time constitute an important percentage of total training time [12]. Fig. 11 outlines the backend infrastructure, distinguishing intra-node from inter-node communication and the main interconnects at each level. Accordingly, backend communication can be viewed at two levels: • Intra-node communication: this occurs among accelerators located within the same node, where a node can be viewed as a single compute unit composed of one or more CPUs and multiple accelerators interconnected through dedicated onboard interconnects that provide direct device-to-device communication within the server. It supports communication between devices that are physically close and often participate in parallelism groups that require frequent and synchronized data exchange. • Inter-node communication: this occurs across different nodes in the training cluster. It enables training jobs to scale beyond a single node capacity and carries communication between nodes over the data center network infrastructure. This hierarchy is important because intra-node and inter-node communication have different performance characteristics and place different requirements on the infrastructure. Intra-node communication typically offers lower latency and higher actual bandwidth, while inter-node communication must scale across many servers and is more exposed to routing, contention, and congestion effects. Because training is synchronized, a small number of slow workers can delay the entire job even when the communicated tensor payload remains unchanged [132]. As a result, the same logical communication operation may behave differently depending on whether it remains inside a node or crosses node boundaries. 7.2. Interconnect technologies and protocol mechanisms The hierarchy introduced above separates backend training communication into intra-node and inter-node exchanges. These two levels are realized through different interconnect technologies and protocol mechanisms, which affect not only the achievable bandwidth and latency, but also the amount of traffic observed below the payload level. In the previous sections, communication volume was derived from tensor sizes and collective operations. At the infrastructure level, this payload is transported through concrete links and protocols, where packetization, headers, flow control, reliability mechanisms, memory-copy paths, and implementation details determine the actual data movement experienced by the system. 7.2.1. Intra-node interconnects Intra-node communication takes place among accelerators within the same server. A baseline technology for this communication is PCI Express (PCIe), which connects CPUs, GPUs, storage devices, network interfaces, and other accelerator cards inside a node. PCIe communication is organized into lanes, commonly grouped 21

Intra-node

Node 1

(within a server)

Intra-node

NVSwitch GPU 0

GPU 1

GPU 2

NIC

NIC

Intra-node

GPU 3

GPU 0

GPU 1

GPU 2

NVSwitch GPU 3

GPU 0

PCIe Switch NIC

NIC

NIC

NIC

Node N

(within a server)

NVSwitch

PCIe Switch NIC

Node 2

(within a server)

GPU 1

GPU 2

GPU 3

PCIe Switch NIC

NIC

NIC

NIC

NIC

Network Switch / Fabric (inter-node) PCIe lines

Inter-node link

NVLink

(Ethernet/InfiniBand/Optical Fiber)

Figure 11: Backend communication hierarchy in distributed LLM training. Accelerators within a node communicate through local high-speed interconnects, while communication between nodes crosses the data-center network fabric.

into configurations such as x4, x8, or x16, with modern GPUs typically using x16 connections. For example, PCIe 6.0 can provide a nominal bandwidth of up to approximately 128 GB/s in an x16 configuration. From a traffic characterization perspective, PCIe transports the training payload together with protocol and control information. Data transfers are packetized into Transaction Layer Packets (TLPs), and additional metadata is introduced for reliability, flow control, acknowledgments, ordering, and integrity checks [133]. Estimating this overhead precisely is difficult because it depends on the PCIe generation, packet size, transfer pattern, and hardware implementation. Direct low-level measurement also typically requires specialized instrumentation. For the purpose of this tutorial, prior studies [134] and vendor or technical documentation [135, 136] suggest that, for large transfers such as those commonly generated during distributed training, PCIe protocol overhead can be approximated as roughly 5-10% of the payload traffic. Although PCIe remains an important intra-node communication substrate, it presents limitations for large-scale GenAI training. First, GPU-to-GPU communication over PCIe may share paths with CPUs, memory, storage, and network interfaces, reducing the actual bandwidth available to training traffic. Second, communication across multiple GPUs can be affected by the internal PCIe layout of the server, including switches and CPU domains. These limitations have motivated the adoption of dedicated accelerator interconnects. In NVIDIA-based systems, NVLink provides direct high-bandwidth GPUto-GPU communication, while NVSwitch extends this design by providing a switching fabric that allows multiple GPUs within a node to communicate with high aggregate bandwidth. These technologies make the node behave more like a tightly coupled accelerator system and are particularly beneficial for communication-intensive parallelization strategies, such as TP. Their detailed low-level overhead is harder to characterize because the implementation is proprietary, but they are designed to provide higher actual bandwidth and lower communication overhead than PCIe-based GPU-to-GPU transfers. The performance characteristics of such intra-node interconnects have been analyzed in prior work, including the study by Li et al. [61]. 7.2.2. Inter-node interconnects and Remote Direct Memory Access (RDMA) Inter-node communication occurs across different nodes in the training cluster and commonly relies on RDMA-capable high performance fabrics. Traditional TCP/IP communication requires data to traverse the operating system networking stack, often involving multiple memory copies and significant CPU intervention. As link speeds increase to hundreds of gigabits per second, this software mediated path can become a bottleneck for distributed AI training [137]. RDMA addresses this limitation by allowing Network Interface Cards (NICs) to transfer data directly between memory regions on different nodes, reducing CPU involvement and communication overhead [64]. In GPU clusters, this idea is extended to accelerator memory. GPUDirect RDMA enables compatible NICs to access GPU memory directly, avoiding intermediate copies through host memory and improving the efficiency of inter-node GPU communication [62]. RDMA is a communication mechanism rather than a single network protocol. One widely deployed implementation is InfiniBand, which provides a native RDMA communication stack and a dedicated high-performance fabric based on InfiniBand adapters, links, and switches [63]. RDMA can also be deployed over Ethernet through RDMA over Converged Ethernet (RoCE). RoCE v1 operates directly at the Ethernet link layer and is therefore limited to layer 2 domains. RoCE v2 encapsulates RDMA traffic within UDP/IP packets, enabling layer 3 routing and making RDMA communication possible across routable Ethernet networks [65]. Another RDMA-over-IP approach is iWARP, 22

Table 5: Interconnect Technologies for Distributed LLM Training.

Level

Technology

Role and traffic-characterization implication

Intra-node

PCIe

Intra-node

NVLink/NVSwitch

Inter-node

InfiniBand

Inter-node

RoCE v2

Memory path

GPUDirect RDMA

General-purpose server interconnect connecting CPUs, GPUs, NICs, storage, and other devices. GPU-to-GPU transfers may share paths with other components and may depend on the server PCIe layout. Traffic includes protocol overhead due to packetization, ordering, flow control, reliability, and integrity mechanisms. Dedicated accelerator interconnect for high-bandwidth GPU-to-GPU communication within a node. It reduces reliance on PCIe paths and is relevant for communication-intensive parallelism dimensions such as TP. Low-level protocol overhead is difficult to expose because implementations are proprietary. Dedicated RDMA-capable fabric widely used in HPC and AI clusters. It provides native RDMA support for low-latency and high-bandwidth node-to-node communication. Traffic characterization must account for message size, congestion, routing, and collective implementation. RDMA over routable Ethernet using UDP/IP encapsulation. It enables RDMA communication over Ethernet data center networks and allows packet-level inspection with standard tools. Header overhead is easier to estimate than for proprietary intra-node fabrics, but actual performance depends on MTU, congestion control, loss behavior, and Ethernet fabric configuration. Mechanism that allows compatible NICs to access GPU memory directly, avoiding intermediate copies through host memory. It reduces CPU involvement and host memory traffic, improving the efficiency of inter-node GPU communication. It affects the memory-copy path rather than defining a separate network fabric.

which runs over TCP/IP, although it is less common in large-scale LLM training deployments [138]. 7.2.3. Observability and protocol overhead in interconnects The choice of interconnect also affects how easily traffic can be measured and how protocol overhead can be estimated. Intra-node technologies such as PCIe, NVLink, and NVSwitch are difficult to inspect directly at the packet or transaction level without specialized hardware support, and proprietary accelerator fabrics expose limited low-level protocol information. In contrast, inter-node traffic carried over Ethernet-based RoCE v2 can often be captured and analyzed using standard packet-inspection tools, since RDMA traffic is encapsulated in UDP/IP packets. This makes header overhead and packet-level behavior easier to study. For the large messages commonly exchanged during distributed LLM training, the protocol overhead of RDMA based inter-node communication is typically small compared with the tensor payload, often limited to a few percent under favorable conditions. However, the actual traffic observed on the network also depends on MTU size, encapsulation, congestion-control behavior, retransmissions, and implementation choices. Therefore, payload-level estimates derived from CCOs should be interpreted as baseline communication requirements, while actual network traffic may include additional protocol and control overheads. Table 5 summarizes the main interconnect technologies discussed in this section and highlights their role in traffic characterization. The key distinction is between payload traffic, determined by the training strategy and collective operation, and actual network traffic, which also depends on the interconnect, protocol stack, and implementation overheads. The table omits fixed bandwidth and latency values because they depend on hardware generation, device model, link configuration, topology, and deployment. Numerical traffic estimates should instead state their assumptions in a separate parameter table, distinguishing nominal link bandwidth from the actual bandwidth seen by the training workload. 7.3. Data center topology and routing effects Once training spans multiple nodes, the data center topology determines how the logical communication schedule is mapped onto physical links. Although topology does not change the original tensor payload, it affects the number of traversed hops, available path diversity, bisection bandwidth, contention, and the actual bandwidth observed by the training workload. Consequently, communication groups should be placed according to their traffic characteristics. TP and SP groups generally benefit from high-bandwidth local communication 23

domains; consecutive PP stages should be placed close together; and DP groups require sufficient inter-node bisection bandwidth for large gradient synchronizations. Modern AI clusters often use multiple network interfaces organized into independent or partially independent rails. Rail-aware placement and routing can distribute synchronized collective traffic across these paths, whereas poor assignments may create hotspots despite sufficient aggregate capacity. Equal-cost multipath, adaptive routing, and rail-aware routing similarly influence link utilization. Congestion control, buffering, and loss recovery further affect communication completion times, particularly in Ethernet-based RDMA deployments. Production systems illustrate these considerations through different designs, including Spectrum-X, Alibaba HPN, large-scale Meta RoCE fabrics, and the optically reconfigurable TPU v4 interconnect [69, 67, 18, 68]. Recent work also treats topology as a design variable that can be optimized for communication patterns, scalability, reliability, and cost [139]. Traffic characterization should therefore distinguish logical communication volume from infrastructure-dependent wire-level behavior, which is determined by placement, routing, and contention. 8. Conclusions and future work The growing scope of LLM training has made communication a central component of system performance. Understanding this communication requires linking concepts that are commonly studied in isolation, from datasets and model architectures to distributed execution and network infrastructure. This tutorial introduced the ‘Life of a Token‘framework to highlight this connection, tracing data from raw text and tokenization through Transformer tensors, parallel execution, collective communication, and finally to the underlying network traffic. The analysis led to several observations that cut across the various levels. Choices regarding the dataset and tokenization determine the number and length of the sequences processed during training, while the model size determines the sizes of the activation, parameter, and gradient tensors. The parallelization strategy determines which tensors become communication objects, which accelerators exchange them, and how frequently these exchanges occur. Collective algorithms then transform the logical payloads of the tensors into communication patterns and into traffic. Finally, the infrastructure determines how these transfers translate into training completion times and scalability. Consequently, no single metric, such as the number of parameters, tensor size, or nominal link bandwidth, is sufficient to characterize the communication behavior of a distributed training configuration. Communication optimization should therefore be considered a cross layer problem. At the software level, collective communication libraries can improve the operations through algorithm selection, hierarchical execution, message partitioning, scheduling, and overlap with computation. At the infrastructure level, accelerator and network designers continue to increase intra-node and inter-node communication capacity and reduce communication latency. Network topology and accelerator placement also determine which transfers use high-capacity local paths and which traverse the broader fabric. Conversely, model partitioning and hybrid parallelization strategies can be adapted to the available topology; these optimization dimensions are closely interrelated: improving one level without considering the others might simply shift the bottleneck elsewhere. Future systems should therefore design the model architecture, parallelization, collective implementation, placement, and network infrastructure in an increasingly integrated manner, optimizing the system as a whole instead of treating each component independently [140, 141]. Rather than prescribing a single optimal configuration, this tutorial provides a framework for identifying where communication originates and where optimization opportunities arise, helping model designers assess infrastructure implications and networking researchers derive traffic requirements from training configurations. Recent work increasingly frames the understanding of LLMs as a multi level scientific problem, spanning learning dynamics, emergent capabilities, internal representations, and the mechanisms that develop during training [142, 143, 144, 145]. These efforts point toward more predictive and falsifiable explanations of deep learning. This tutorial complements them with a systems perspective, showing how large scale training gives rise to tensor exchanges, parallel execution, collective communication, and network traffic. The Life of a Token framework helps make this process more measurable and understandable. Acknowledgements The authors would like to thank Massimo Gallo for his feedback and title idea. We used generative AI tools to assist with figure design and to correct grammar and spelling. This research was conducted as part of the Net4AI project, funded by the French National Research Agency (ANR) under contract ANR-24-CE25-5120.

24

Output Probabilities

Output Probabilities

Output Probabilities

Softmax

Softmax

Softmax Linear

MoE Feed-Forward Layer

Lineax

Linear

Input tokens (from attention)

Add & Norm

Add & Norm Feed Forward Lx

Token 1

Feed Forward

Add & Norm

Add & Norm

Feed Forward

Nx Add & Norm

Multi-Head Attention

Masked Multi-Head Attention

Positional Encoding

Positional Encoding

Positional Encoding

Add & Norm

Output Embedding

Input Embedding

Inputs

Outputs (shifted right)

Inputs

(a) Original encoder–decoder Transformer

Expert 1

Expert 2

…

Expert NE

…

Masked Multi-Head Attention

Positional Encoding

Input Embedding

Token N

MoE Feed-Forward Layer

Lx

Masked Multi-Head Attention

Add & Norm Add & Norm

…

Router / Gating Network

Multi-Head Attention Lx

Token 2

Add & Norm

Output tokens (to next layer)

+

Token 1’

Token 2’

…

Token N’

Input Embedding Routing (each token sent to one expert)

(b) Decoder-only model

Inputs (tokens)

Outputs combined to reconstruct the batch

(c) Token processing in an MoE layer

Figure A.12: Comparison of: (a) the original encoder–decoder Transformer architecture [3]; (b) the GPT-2-inspired decoder-only Transformer adopted as the reference architecture throughout this tutorial [2]; and (c) a representation of a decoder-only Mixtureof-Experts (MoE) architecture.

Appendix Appendix A: Scope and organization The main body establishes simplified communication baselines for Data Parallelism (DP), Pipeline Parallelism (PP), and Tensor Parallelism (TP). This appendix extends that analysis in four directions. Appendix A.1 relates the original encoder-decoder Transformer to the decoder-only architecture used throughout the tutorial. Appendix A.2 examines how state sharding, pipeline scheduling, and sequence partitioning modify the communication baselines. Appendix A.3 connects the analytical DP model to measurements from a controlled small-cluster experiment. Finally, Appendix A.4 follows token representations through a Mixture-of-Experts model and derives the communication introduced by Expert Parallelism (EP). Unless stated otherwise, the notation follows the main tutorial. In particular, B is the local mini-batch size, T the sequence length, d the hidden dimension, and sact the number of bytes used for each activation value. The activation payload is therefore M = BT d sact . (A.1) We use FW and BW for the forward and backward passes, CCOs for collective communication operations, and pcc for the number of accelerators in the relevant communication group. Appendix A.1. From the Original Transformer to Decoder-Only LLMs The original Transformer was introduced as an encoder-decoder architecture for sequence-to-sequence tasks such as machine translation [3]. As shown in Fig. A.12a, the encoder processes the source sequence through repeated self-attention and feed-forward blocks. The decoder generates the target sequence using masked self-attention, cross-attention over the encoder output, and a feed-forward network (FFN). The encoder and decoder therefore play distinct roles: the former constructs a representation of the input sequence, while the latter conditions on that representation and on the previously generated target tokens. Decoder-only models remove the encoder stack and the encoder-decoder cross-attention path. Instead, they apply a single repeated stack of causally masked self-attention and FFN blocks to predict each token from the tokens that precede it. This design was popularized by GPT-style language models [2, 15] and constitutes the reference architecture used throughout the main tutorial. Figure A.12 compares the original encoder–decoder architecture in panel (a) with the decoder-only architecture in panel (b). Panel (c) additionally previews the Mixture-of-Experts (MoE) architecture, in which the dense FFN is replaced by multiple experts and token representations are routed to selected experts for processing. This routing process preserves the input and output tensor shape at the layer boundary but introduces specialized communication patterns. The MoE architecture and its communication implications are examined in detail in Appendix A.4

25

Figure A.13: Logical communication schedules in standard DP and ZeRO.

Appendix A.2. Communication Refinements for Dense-Model Parallelism The communication baselines in the main tutorial describe the dominant behavior of each parallelism dimension: model sized gradient synchronization for DP, activation-sized P2P transfers for PP, and frequent activation sized collectives within Transformer layers for TP. Practical systems refine these baselines to reduce memory redundancy, pipeline idle time, or exposed communication. Such refinements may change the collective sequence and message granularity even when the total communicated volume remains of the same order. State sharding in data parallelism The most widely adopted family of optimizations for DP training focuses on reducing memory redundancy by partitioning optimizer states, gradients, or model parameters across accelerators, instead of fully replicating them on every device. Missing information is reconstructed through CCOs when needed. Two widely adopted examples are ZeRO [42], with its variants ZeRO-1, ZeRO-2, and ZeRO-3, and Fully Sharded Data Parallelism (FSDP) [43]. DeepSpeed [41] is a widely used deep learning optimization library that implements many of these techniques in practice, including the ZeRO family and several of its extensions. From a network perspective, these approaches progressively redistribute communication across the training process. In standard DP, the ReduceScatter and AllGather phases are executed consecutively as part of a single gradient AllReduce. In ZeRO-1 and ZeRO-2, the overall communication volume remains comparable to standard DP, but the communication schedule is reorganized. In particular, ZeRO-2 partitions gradients, so the ReduceScatter and AllGather phases are separated by the optimizer update rather than being executed as one monolithic gradient AllReduce. ZeRO-3 changes the traffic pattern further by also sharding model parameters. Because parameters are no longer fully replicated on each accelerator, the required parameter shards must be gathered before computation. In a simplified view, this introduces an AllGather before the FW pass, another AllGather before the BW pass, and a final ReduceScatter after the BW pass to distribute the reduced gradient partitions. Therefore, unlike ZeRO-1 and ZeRO-2, ZeRO-3 increases the total communication volume because it communicates parameter shards in addition to gradients. The original ZeRO analysis reports this overhead to be approximately 1.5× the communication volume of standard DP [42]. If the standard DP traffic is denoted by CDP , then CZeRO-3 ≈ 1.5CDP . FSDP follows a conceptually similar approach by sharding model states across accelerators, while dynamically reconstructing them during execution [43]. Fig. A.13 summarizes when CCOs are performed in standard DP and in the ZeRO variants. Several refinements and implementation variants have also been proposed, including ZeRO-Offload, ZeRO-Infinity, and ZeRO++ [51, 52, 53], which further optimize memory usage and communication efficiency. Pipeline schedules and communication overlap The main objective of PP optimization is to reduce pipeline bubbles. A common approach is to divide each mini-batch into micro-batches. Instead of processing the entire mini-batch at once, different pipeline stages can simultaneously work on different micro-batches, increasing hardware utilization. Early systems such as 26

Table A.6: First-order communication effects of dense-model refinements.

Refinement

Primary objective

Dominant communication

Effect relative to baseline

ZeRO-1/2

Reduce replicated optimizer or gradient state

ReduceScatter and AllGather around the update

ZeRO-3/FSDP

Shard parameters, gradients, and optimizer state

Parameter AllGather and gradient ReduceScatter

PP schedules

Reduce bubbles and exposed transfer time

SP and long-context variants

Reduce activation redundancy and support long sequences

Point-to-point activation and activation-gradient transfers Collectives over sequence or activation partitions

Comparable first-order volume; schedule is reorganized Additional state reconstruction; overlap and bucketing matter Similar leading-order bytes; smaller messages and different timing Implementation-dependent schedule; activation scale remains central

GPipe [44] reduce bubble overhead by splitting mini-batches into many micro-batches and executing them in a pipelined manner. Later approaches, such as PipeDream’s one-forward-one-backward (1F1B) scheduling strategy [45], further improve utilization by overlapping FW and BW passes across different micro-batches. TeraPipe [46] further improves pipeline efficiency by introducing token-level PP, where each sequence is split into smaller token slices so that stages can exchange smaller tensors more frequently. More recent works, including Zero Bubble scheduling techniques [47], aim to further minimize idle periods and improve pipeline efficiency. Another important research direction focuses on communication-aware pipeline execution, where communication between adjacent stages is overlapped with computation in order to reduce the impact of latency and bandwidth limitations. A representative example is DualPipe [19], which explicitly overlaps communication and computation across pipeline stages. Additional optimizations have also been proposed for reducing activation memory usage, preserving weight consistency, and improving the placement and partitioning of pipeline stages across accelerators [12, 146, 49]. Sequence and long-context partitioning One limitation of the original TP formulation is that only the large matrix multiplications are partitioned across GPUs, while several other operations inside the Transformer layer remain replicated on every accelerator. Although these operations are less computationally intensive, they can still consume a significant amount of memory. To improve memory efficiency, several refinements of TP have been proposed. One widely adopted extension is sequence parallelism (SP) [108], which partitions selected operations along the sequence dimension, reducing memory redundancy while preserving the overall communication structure of TP. From a network perspective, the total communication volume remains broadly comparable to standard TP [108], although communication and computation are scheduled differently. As with the PP optimizations discussed previously, SP primarily reorganizes the timing and scheduling of communication rather than fundamentally changing the communication primitives or the amount of exchanged data. Consequently, the basic TP formulation analyzed above still provides the core intuition required to characterize the resulting network traffic. Once the execution schedule is known, the timing and organization of the communication phases can be derived accordingly. Several additional works have further extended sequence parallelism for long-context LLM training by improving communication efficiency, memory scalability, and execution scheduling across distributed accelerators [147, 148, 149, 150, 151]. Comparative effects of dense-model refinements To a first-order approximation, refinements of a given parallelization technique exchange the same underlying tensors as their baseline formulation. The principal differences concern how these tensors are partitioned, when the corresponding communication operations are scheduled, and how they are interleaved or overlapped with computation. Table A.6 relates the resource bottleneck addressed by each refinement to its principal communication effects. This comparison distinguishes the total volume of transferred data from the portion of communication exposed on the execution critical path. Consequently, an optimization may improve E2E throughput without reducing the total byte volume, for example, by dividing transfers into smaller messages, overlapping them with computation, or redistributing them across the training step.

27

Table A.7: Measured DDP scaling on the seven-node cluster. Speedup and efficiency are computed relative to the single-GPU configuration. The AllReduce share is the fraction of the step time spent in gradient synchronization.

GPUs

1 2 4 7

AllReduce share

(×)

Throughputscaling efficiency (%)

1.00 1.65 3.00 4.69

100 82 75 67

– 15.0 20.2 24.5

Step time

AllReduce

Throughput

Speedup

(ms)

(ms)

(tokens/s)

930 1,130 1,240 1,390

– 170 250 340

1,100 1,810 3,300 5,160

(%)

Appendix A.3. Empirical Illustration: DDP Scaling on a Small GPU Cluster To connect the analytical communication model to E2E behavior, we evaluate a controlled distributed DP workload on a seven-node GPU cluster. Each node contains a dedicated GPU, the nodes are interconnected by 10 Gb/s links, and each GPU hosts one training worker. The workload uses the GPT-2 implementation cited in the main tutorial [106]. Communication activity is recorded with tcpdump. The traces are used to identify gradient-synchronization bursts and estimate achieved transfer rates. Table A.7 reports step time, AllReduce time, throughput, and scaling efficiency for one, two, four, and seven workers. Results and interpretation Because each worker processes approximately the same number of tokens per step, the global workload increases with the number of workers. The experiment therefore characterizes throughput scaling under an approximately constant per-worker workload rather than strong scaling under a fixed global workload. Aggregate throughput increases from 1,100 tokens/s with one GPU to 5,160 tokens/s with seven GPUs, corresponding to a speedup of 4.69× and a throughput-scaling efficiency of 67%. As the number of workers increases, the measured AllReduce time rises from 170 ms with two GPUs to 340 ms with seven GPUs, while its share of the training-step time increases from approximately 15% to 24%. These measurements indicate that collective communication and synchronization impose increasing overhead under the evaluated configuration, thereby limiting the growth in aggregate throughput. Link bandwidth affects the duration of gradient-synchronization bursts, but collective algorithms, latency, software overhead, stragglers, and the ratio of local computation to communication also influence scaling efficiency. Similar communicationrelated scaling limitations have been reported in large-scale LLM training studies [11, 121, 152]. Limitations and scope This experiment provides a controlled empirical illustration of the relationship predicted by the analytical model: as the number of DDP workers increases, gradient synchronization occupies a larger fraction of the training step and reduces scaling efficiency. The seven-node testbed allows this effect to be observed by relating communication time to end-to-end throughput. Although the absolute measurements are specific to the evaluated hardware and 10 Gb/s interconnect and should not be extrapolated directly to modern accelerator clusters, the observed relationship between synchronization cost, communication intensity, and scaling efficiency remains broadly applicable. Appendix A.4. Mixture of Experts The main tutorial focuses on dense Transformer models, where every token is processed by the same set of model parameters. One of the most widely adopted modifications to this architecture concerns the FFN block of the Transformer architecture. MoE architectures [101] replace the single FFN of a Transformer layer with a collection of smaller FFNs, referred to as experts. Let NE denote the number of experts available in a given MoE layer. The key idea is that each token is no longer processed by the same FFN. Instead, after the attention operation, each token is processed only by a specific expert, or by a small subset of experts, selected from the available NE experts. To determine which expert should process a token, MoE layers employ a small routing network, commonly referred to as a router or gating network. Given the token representation produced by the attention mechanism, the router computes a score for each expert. The experts with the highest scores are then selected, typically using a top-K strategy, where only the K highest scoring experts are activated for that token. The outputs of the selected experts are subsequently returned to the model’s computation and forwarded to 28

Figure A.14: FW pass schema of an MoE Transformer layer.

the following layers. Fig. A.12c provides a schematic representation of an MoE layer, showing how tokens are routed to selected experts and subsequently merged to reconstruct the output batch. The main advantage of MoE architectures is that they increase model capacity without proportionally increasing the computation required to process each token. While the total number of model parameters grows with the number of experts, only a subset of these parameters is activated for a given token. Consequently, MoE architectures can achieve parameter counts significantly larger than those of dense models while maintaining comparable computational requirements. To introduce the parallelization strategies commonly adopted for MoE models and understand the communication patterns they generate, it is useful to first examine how a single MoE Transformer layer operates on a single accelerator. Consider an input tensor X ∈ RB×T ×d . As in a dense Transformer, the tensor is first processed by the attention mechanism, which produces an output tensor of the same dimensions. The routing network then computes a score for each token embedding and assigns it to one or more experts according to the selected routing policy. Consequently, the token embeddings contained in X are partitioned into groups, with each group processed by a specific expert. After the expert computations are completed, the resulting embeddings are reassembled according to their original token positions, producing an output tensor with the same dimensions as the input tensor. This tensor is then forwarded to the subsequent Transformer layer. During backpropagation, gradients follow the same routing path in the reverse direction. The gradients associated with each token are first propagated through the expert that processed that token and are then reassembled to reconstruct the gradient tensor corresponding to the original input tensor. Fig. A.14 illustrates how the input tensor is reorganized and processed within a single MoE Transformer layer. In the next subsection, we examine how MoE models are distributed across multiple accelerators and analyze the resulting communication and network traffic patterns. Expert parallelism Although MoE architectures increase model capacity without a proportional increase in per-token computation, they still require distributed execution across multiple accelerators. In particular, during training, the complete set of experts must remain resident across the participating accelerators’ memories, even though each token activates only a small subset of them. Thus, memory usage scales with the total number of experts, whereas per-token computation depends mainly on the number of activated experts. The parallelization strategy most commonly associated with is EP [102, 103]. The key idea is to distribute the experts across multiple accelerators while replicating the remaining components of the Transformer layer, such as the attention mechanism and the routing network. During execution, each accelerator processes its local mini-batch in the same way as a dense Transformer up to the routing stage. After the attention operation, the router assigns each token embedding to a specific expert. Since experts are distributed across accelerators, the selected expert may reside on a different device from the one currently processing the token. In this case, the token embedding is transferred to the accelerator hosting the selected expert, then returned after computation and reassembled into the original output tensor before being forwarded to the next Transformer layer. To maximize hardware utilization, MoE systems typically attempt to balance the number of tokens assigned to each expert. When experts are distributed across multiple accelerators, this routing process generates a communication pattern commonly referred to as All-to-All. This communication pattern is typically implemented 29

Inputs from prior layer MoE model MoE Layer

GPU0

GPU1

GPU2

Gating

Gating

Gating

Dense Layer

All-to-All Dispatch

MoE Layer Dense Layer

Compute0

Compute1

Compute2

All-to-All Combine

MoE Layer

Sum

Sum

Sum

Dense Layer

Outputs to next layer Figure A.15: Example of the All-to-All communication pattern adapted from Fig. 2 of [154].

through the corresponding All-to-All CCO. After the routing stage, each accelerator sends the token embeddings assigned to remote experts and receives the embeddings that must be processed by its local experts. Assuming an ideally load-balanced scenario with pcc accelerators participating in the expert parallel group and a top-K −1 routing policy, each accelerator exchanges approximately: K pcc pcc X, of the local tensor X. In the common −1 top-1 configuration, K = 1, this expression reduces to pcc pcc X. Once the expert computations are completed, a second All-to-All communication phase is performed to return the processed embeddings to their originating accelerators and reconstruct the original tensor organization. Consequently, EP introduces two All-to-All communication phases during the FW pass, with analogous communication occurring during backpropagation. More broadly, MoE training systems can co-schedule intra-node and inter-node communication with expert computation to reduce the exposed All-to-All overhead [153]. Fig. A.15, adapted from Lei et al. [154] provides a schematic representation of the All-to-All communication pattern introduced by EP. Let MMoE = BT d sact denote the byte payload of the local token-representation tensor. Assuming top-1 routing (K = 1), ideal load balance, and an expert-parallel group of pcc = 8 accelerators, the bytes sent by each (1) −1 8−1 accelerator in one All-to-All phase are approximately CA2A ≈ K pcc pcc MMoE = 1 × 8 × 50 MB = 43.75 MB. Under ideal balance, each accelerator receives the same amount. EP uses two such phases in the forward pass and two analogous phases in the backward pass. Using a sent-bytes convention, the resulting volume is therefore 4 × 43.75 = 175 MB per accelerator per MoE layer. The refinements and sparse-model extension lead to four practical conclusions. First, communicated byte volume and communication exposed on the critical path are different quantities: PP schedules and overlapping techniques may improve throughput without reducing bytes. Second, memory-saving methods can increase or redistribute communication, as illustrated by parameter reconstruction in ZeRO-3 and FSDP. Third, the smallcluster DDP result supports the analytical trend but should be interpreted within its experimental boundary rather than as a modern-cluster benchmark. Finally, EP changes the traffic structure from dense-model gradient or activation exchanges to token-dependent All-to-All communication. Its performance therefore depends not only on tensor size and group size, but also on routing balance, expert placement, and the physical network topology. Together, these observations extend the Life of a Token framework beyond the dense baselines considered in the main tutorial. Whether a token representation is sharded, rescheduled, or routed, the analysis

30

continues to identify the tensor being transferred, the participating devices, the invocation frequency of the communication operation, and the point at which the resulting traffic crosses the infrastructure hierarchy. References [1] D. Alighieri, La divina commedia, Le Monnier, Firenze, 1979. [2] A. Radford, J. Wu, R. Child, D. Luan, D. Amodei, I. Sutskever, et al., Language models are unsupervised multitask learners, OpenAI blog 1 (8) (2019) 9. [3] A. Vaswani, N. Shazeer, N. Parmar, J. Uszkoreit, L. Jones, A. N. Gomez, Ł. Kaiser, I. Polosukhin, Attention is all you need, Advances in neural information processing systems 30. [4] J. Kaplan, S. McCandlish, T. Henighan, T. B. Brown, B. Chess, R. Child, S. Gray, A. Radford, J. Wu, D. Amodei, Scaling laws for neural language models, arXiv preprint arXiv:2001.08361. [5] N. Alarcon, Openai presents gpt-3, a 175 billion parameters language model, accessed: 2026-03-05 (Jul. 2020). URL https://developer.nvidia.com/blog/openai-presents-gpt-3-a-175-billion-parameter s-language-model/ [6] J. J. Tithi, H. Wu, A. Abuhatzera, F. Petrini, Scaling intelligence: Designing data centers for next-gen language models, arXiv preprint arXiv:2506.15006. [7] S. Cheng, J.-L. Lin, M. Emani, S. Raskar, S. Foreman, Z. Xie, V. Vishwanath, M. T. Kandemir, Thorough characterization and analysis of large transformer model training at-scale, Proceedings of the ACM on Measurement and Analysis of Computing Systems 8 (1) (2024) 1–25. [8] M. Wang, C. Meng, G. Long, C. Wu, J. Yang, W. Lin, Y. Jia, Characterizing deep learning training workloads on alibaba-pai, in: 2019 IEEE international symposium on workload characterization (IISWC), IEEE, 2019, pp. 189–202. [9] J. Li, Y. Jiang, Y. Zhu, C. Wang, H. Xu, Accelerating distributed {MoE} training and inference with lina, in: 2023 USENIX Annual Technical Conference (USENIX ATC 23), 2023, pp. 945–959. [10] B. Hanindhito, B. Patel, L. K. John, Bandwidth characterization of deepspeed on distributed large language model training, in: 2024 IEEE International Symposium on Performance Analysis of Systems and Software (ISPASS), IEEE, 2024, pp. 241–256. [11] Z. Jiang, H. Lin, Y. Zhong, Q. Huang, Y. Chen, Z. Zhang, Y. Peng, X. Li, C. Xie, S. Nong, et al., {MegaScale}: Scaling large language model training to more than 10,000 {GPUs}, in: 21st USENIX Symposium on Networked Systems Design and Implementation (NSDI 24), 2024, pp. 745–760. [12] D. Narayanan, M. Shoeybi, J. Casper, P. LeGresley, M. Patwary, V. Korthikanti, D. Vainbrand, P. Kashinkunti, J. Bernauer, B. Catanzaro, et al., Efficient large-scale language model training on gpu clusters using megatron-lm, in: Proceedings of the international conference for high performance computing, networking, storage and analysis, 2021, pp. 1–15. [13] S. Raschka, Build a large language model (from scratch), Simon and Schuster, 2024. [14] J. Alammar, M. Grootendorst, Hands-on large language models: language understanding and generation, O’Reilly Media, Inc., Sebastopol, CA, USA, 2024. [15] T. Brown, B. Mann, N. Ryder, M. Subbiah, J. D. Kaplan, P. Dhariwal, A. Neelakantan, P. Shyam, G. Sastry, A. Askell, et al., Language models are few-shot learners, Advances in neural information processing systems 33 (2020) 1877–1901. [16] Q. Anthony, B. Michalowicz, J. Hatef, L. Xu, M. Abduljabbai, A. Shafi, H. Subramoni, D. K. Panda, Demystifying the communication characteristics for distributed transformer models, in: 2024 IEEE Symposium on High-Performance Interconnects (HOTI), IEEE, IEEE, Piscataway, NJ, USA, 2024, pp. 57–65. [17] Z. Hu, S. Shen, T. Bonato, S. Jeaugey, C. Alexander, E. Spada, J. Dinan, J. Hammond, T. Hoefler, Demystifying nccl: An in-depth analysis of gpu communication protocols and algorithms, in: 2025 IEEE Symposium on High-Performance Interconnects (HOTI), IEEE, 2025, pp. 48–59. 31

[18] A. Gangidi, R. Miao, S. Zheng, S. J. Bondu, G. Goes, H. Morsy, R. Puri, M. Riftadi, A. J. Shetty, J. Yang, et al., Rdma over ethernet for distributed training at meta scale, in: Proceedings of the ACM SIGCOMM 2024 Conference, 2024, pp. 57–70. [19] A. Liu, B. Feng, B. Xue, B. Wang, B. Wu, C. Lu, C. Zhao, C. Deng, C. Zhang, C. Ruan, et al., Deepseek-v3 technical report, arXiv preprint arXiv:2412.19437. [20] D. Guo, D. Yang, H. Zhang, J. Song, P. Wang, Q. Zhu, R. Xu, R. Zhang, S. Ma, X. Bi, et al., Deepseek-r1: Incentivizing reasoning capability in llms via reinforcement learning, arXiv preprint arXiv:2501.12948. [21] B. Hui, J. Yang, Z. Cui, J. Yang, D. Liu, L. Zhang, T. Liu, J. Zhang, B. Yu, K. Lu, et al., Qwen2.5-coder technical report, arXiv preprint arXiv:2409.12186. [22] H. Touvron, L. Martin, K. Stone, P. Albert, A. Almahairi, Y. Babaei, N. Bashlykov, S. Batra, P. Bhargava, S. Bhosale, et al., Llama 2: Open foundation and fine-tuned chat models, arXiv preprint arXiv:2307.09288. [23] R. Anil, A. M. Dai, O. Firat, M. Johnson, D. Lepikhin, A. Passos, S. Shakeri, E. Taropa, P. Bailey, Z. Chen, et al., Palm 2 technical report (2023). arXiv:2305.10403. URL https://arxiv.org/abs/2305.10403 [24] J. Achiam, S. Adler, S. Agarwal, L. Ahmad, I. Akkaya, F. L. Aleman, D. Almeida, J. Altenschmidt, S. Altman, S. Anadkat, et al., Gpt-4 technical report (2024). arXiv:2303.08774. URL https://arxiv.org/abs/2303.08774 [25] J. Duan, S. Zhang, Z. Wang, L. Jiang, W. Qu, Q. Hu, G. Wang, Q. Weng, H. Yan, X. Zhang, et al., Efficient training of large language models on distributed infrastructures: a survey, Vicinagearth 3 (1) (2026) 9. [26] F. Liang, Z. Zhang, H. Lu, V. Leung, Y. Guo, X. Hu, Communication-efficient large-scale distributed deep learning: A comprehensive survey, arXiv preprint arXiv:2404.06114. [27] N. Tazi, F. Mom, H. Zhao, P. Nguyen, M. Mekkouri, L. von Werra, T. Wolf, The ultra-scale playbook: Training LLMs on GPU clusters, accessed: 2026-09-04 (2025). URL https://huggingface.co/spaces/nanotron/ultrascale-playbook [28] X. Song, M. Zhang, Y. Liu, J. Xia, S. Yang, X. Hu, C. Hu, M. Xu, Collective communication for distributed llm systems: Planning, runtime adaptation, and computation coordination, IEEE Network. [29] W. X. Zhao, K. Zhou, J. Li, T. Tang, X. Wang, Y. Hou, Y. Min, B. Zhang, J. Zhang, Z. Dong, et al., A survey of large language models, arXiv preprint arXiv:2303.18223 1 (2) (2023) 1–124. [30] Z. Wang, Z. Chu, T. V. Doan, S. Ni, M. Yang, W. Zhang, History, development, and principles of large language models: an introductory survey, AI and Ethics 5 (3) (2025) 1955–1971. [31] Y. Liu, J. Cao, C. Liu, K. Ding, L. Jin, Datasets for large language models: A comprehensive survey, Artificial Intelligence Review 58 (12) (2025) 403. [32] Common Crawl Foundation, Common crawl corpus (2024). URL https://commoncrawl.org [33] L. Gao, S. Biderman, S. Black, L. Golding, T. Hoppe, C. Foster, J. Phang, H. He, A. Thite, N. Nabeshima, et al., The pile: An 800gb dataset of diverse text for language modeling, arXiv preprint arXiv:2101.00027. [34] P. Gage, A new algorithm for data compression, The C Users Journal 12 (2) (1994) 23–38. [35] T. Kudo, J. Richardson, Sentencepiece: A simple and language independent subword tokenizer and detokenizer for neural text processing, in: Proceedings of the 2018 conference on empirical methods in natural language processing, 2018, pp. 66–71. [36] X. Song, A. Salcianu, Y. Song, D. Dopson, D. Zhou, Fast wordpiece tokenization, in: Proceedings of the 2021 conference on empirical methods in natural language processing, 2021, pp. 2089–2103. [37] N. Rajaraman, J. Jiao, K. Ramchandran, An analysis of tokenization: Transformers under markov data, Advances in Neural Information Processing Systems 37 (2024) 62503–62556.

32

[38] OpenAI, What are tokens and how to count them?, accessed: 2026-03-05 (2026). URL https://help.openai.com/en/articles/4936856 [39] Google, Token counting, accessed: 2026-06-12 (2026). URL https://ai.google.dev/gemini-api/docs/tokens [40] M. Shoeybi, M. Patwary, R. Puri, P. LeGresley, J. Casper, B. Catanzaro, Megatron-lm: Training multibillion parameter language models using model parallelism, arXiv preprint arXiv:1909.08053. [41] J. Rasley, S. Rajbhandari, O. Ruwase, Y. He, Deepspeed: System optimizations enable training deep learning models with over 100 billion parameters, in: Proceedings of the 26th ACM SIGKDD international conference on knowledge discovery & data mining, 2020, pp. 3505–3506. [42] S. Rajbhandari, J. Rasley, O. Ruwase, Y. He, Zero: Memory optimizations toward training trillion parameter models, in: SC20: international conference for high performance computing, networking, storage and analysis, IEEE, 2020, pp. 1–16. [43] Y. Zhao, A. Gu, R. Varma, L. Luo, C.-C. Huang, M. Xu, L. Wright, H. Shojanazeri, M. Ott, S. Shleifer, et al., Pytorch fsdp: experiences on scaling fully sharded data parallel, arXiv preprint arXiv:2304.11277. [44] Y. Huang, Y. Cheng, A. Bapna, O. Firat, D. Chen, M. Chen, H. Lee, J. Ngiam, Q. V. Le, Y. Wu, et al., Gpipe: Efficient training of giant neural networks using pipeline parallelism, Advances in neural information processing systems 32. [45] A. Harlap, D. Narayanan, A. Phanishayee, V. Seshadri, N. Devanur, G. Ganger, P. Gibbons, Pipedream: Fast and efficient pipeline parallel dnn training, arXiv preprint arXiv:1806.03377. [46] Z. Li, S. Zhuang, S. Guo, D. Zhuo, H. Zhang, D. Song, I. Stoica, Terapipe: Token-level pipeline parallelism for training large-scale language models, in: International Conference on Machine Learning, PMLR, 2021, pp. 6543–6552. [47] P. Qi, X. Wan, G. Huang, M. Lin, Zero bubble (almost) pipeline parallelism, in: International Conference on Learning Representations, Vol. 2024, 2024, pp. 48869–48884. [48] N. Shazeer, Y. Cheng, N. Parmar, D. Tran, A. Vaswani, P. Koanantakool, P. Hawkins, H. Lee, M. Hong, C. Young, et al., Mesh-tensorflow: Deep learning for supercomputers, Advances in neural information processing systems 31. [49] L. Zheng, Z. Li, H. Zhang, Y. Zhuang, Z. Chen, Y. Huang, Y. Wang, Y. Xu, D. Zhuo, E. P. Xing, et al., Alpa: Automating inter-and {Intra-Operator} parallelism for distributed deep learning, in: 16th USENIX Symposium on Operating Systems Design and Implementation (OSDI 22), 2022, pp. 559–578. [50] Y. Xu, H. Lee, D. Chen, B. Hechtman, Y. Huang, R. Joshi, M. Krikun, D. Lepikhin, A. Ly, M. Maggioni, et al., Gspmd: general and scalable parallelization for ml computation graphs, arXiv preprint arXiv:2105.04663. [51] J. Ren, S. Rajbhandari, R. Y. Aminabadi, O. Ruwase, S. Yang, M. Zhang, D. Li, Y. He, {Zero-offload}: Democratizing {billion-scale} model training, in: 2021 USENIX Annual Technical Conference (USENIX ATC 21), 2021, pp. 551–564. [52] S. Rajbhandari, O. Ruwase, J. Rasley, S. Smith, Y. He, Zero-infinity: Breaking the gpu memory wall for extreme scale deep learning, in: Proceedings of the international conference for high performance computing, networking, storage and analysis, 2021, pp. 1–14. [53] G. Wang, H. Qin, S. Jacobs, X. S. Wu, C. Holmes, Z. Yao, S. Rajbhandari, O. Ruwase, F. Yan, L. Yang, et al., Zero++: Extremely efficient collective communication for large model training, in: International conference on learning representations, Vol. 2024, 2024, pp. 50035–50053. [54] Nvidia collective communication library (2025). URL https://github.com/NVIDIA/nccl [55] NVIDIA Corporation, Collective Operations, NVIDIA Corporation, accessed: 2026-05-19 (2026). URL https://docs.nvidia.com/deeplearning/nccl/user-guide/docs/usage/collectives.html

33

[56] A. Weingram, Y. Li, H. Qi, D. Ng, L. Dai, X. Lu, xccl: A survey of industry-led collective communication libraries for deep learning, Journal of Computer Science and Technology 38 (1) (2023) 166–195. [57] Z. Cai, Z. Liu, S. Maleki, M. Musuvathi, T. Mytkowicz, J. Nelson, O. Saarikivi, Synthesizing optimal collective algorithms, in: Proceedings of the 26th ACM SIGPLAN Symposium on Principles and Practice of Parallel Programming, Association for Computing Machinery, New York, NY, USA, 2021, pp. 62–75. [58] A. Shah, V. Chidambaram, M. Cowan, S. Maleki, M. Musuvathi, T. Mytkowicz, J. Nelson, O. Saarikivi, R. Singh, {TACCL}: Guiding collective algorithm synthesis using communication sketches, in: 20th USENIX Symposium on Networked Systems Design and Implementation (NSDI 23), 2023, pp. 593–612. [59] M. Cowan, S. Maleki, M. Musuvathi, O. Saarikivi, Y. Xiong, Mscclang: Microsoft collective communication language, in: Proceedings of the 28th ACM International Conference on Architectural Support for Programming Languages and Operating Systems, Volume 2, 2023, pp. 502–514. [60] X. Liu, B. Arzani, S. K. R. Kakarla, L. Zhao, V. Liu, M. Castro, S. Kandula, L. Marshall, Rethinking machine learning collective communication as a multi-commodity flow problem, in: Proceedings of the ACM SIGCOMM 2024 Conference, 2024, pp. 16–37. [61] A. Li, S. L. Song, J. Chen, J. Li, X. Liu, N. R. Tallent, K. J. Barker, Evaluating modern gpu interconnect: Pcie, nvlink, nv-sli, nvswitch and gpudirect, IEEE Transactions on Parallel and Distributed Systems 31 (1) (2019) 94–110. [62] S. Potluri, K. Hamidouche, A. Venkatesh, D. Bureddy, D. K. Panda, Efficient inter-node mpi communication using gpudirect rdma for infiniband clusters with nvidia gpus, in: 2013 42nd International Conference on Parallel Processing, IEEE, 2013, pp. 80–89. [63] G. F. Pfister, An introduction to the infiniband architecture, High performance mass storage and parallel I/O 42 (617-632) (2001) 102. [64] A. Kalia, M. Kaminsky, D. G. Andersen, Design guidelines for high performance {RDMA} systems, in: 2016 USENIX annual technical conference (USENIX ATC 16), 2016, pp. 437–450. [65] R. Mittal, A. Shpiner, A. Panda, E. Zahavi, A. Krishnamurthy, S. Ratnasamy, S. Shenker, Revisiting network support for rdma, in: Proceedings of the 2018 Conference of the ACM Special Interest Group on Data Communication, 2018, pp. 313–326. [66] Y. Zhu, H. Eran, D. Firestone, C. Guo, M. Lipshteyn, Y. Liron, J. Padhye, S. Raindel, M. H. Yahia, M. Zhang, Congestion control for large-scale rdma deployments, ACM SIGCOMM Computer Communication Review 45 (4) (2015) 523–536. [67] K. Qian, Y. Xi, J. Cao, J. Gao, Y. Xu, Y. Guan, B. Fu, X. Shi, F. Zhu, R. Miao, et al., Alibaba hpn: A data center network for large language model training, in: Proceedings of the ACM SIGCOMM 2024 Conference, 2024, pp. 691–706. [68] N. Jouppi, G. Kurian, S. Li, P. Ma, R. Nagarajan, L. Nai, N. Patil, S. Subramanian, A. Swing, B. Towles, et al., Tpu v4: An optically reconfigurable supercomputer for machine learning with hardware support for embeddings, in: Proceedings of the 50th annual international symposium on computer architecture, 2023, pp. 1–14. [69] NVIDIA, NVIDIA Spectrum-X Network Platform Architecture: The First Ethernet Network Designed to Accelerate AI Workloads, Tech. rep., NVIDIA, accessed: 2026-05-15 (2024). URL https://resources.nvidia.com/en-us-accelerated-networking-resource-library/nvidi a-spectrum-x [70] W. Won, T. Heo, S. Rashidi, S. Sridharan, S. Srinivasan, T. Krishna, Astra-sim2. 0: Modeling hierarchical networks and disaggregated systems for large-model training at scale, in: 2023 IEEE International Symposium on Performance Analysis of Systems and Software (ISPASS), IEEE, 2023, pp. 283–294. [71] X. Wang, Q. Li, Y. Xu, G. Lu, D. Li, L. Chen, H. Zhou, L. Zheng, S. Zhang, Y. Zhu, et al., {SimAI}: unifying architecture design and performance tuning for {Large-Scale} large language model training with scalability and precision, in: 22nd USENIX Symposium on Networked Systems Design and Implementation (NSDI 25), 2025, pp. 541–558.

34

[72] G. He, Y. Jiang, W. Xiao, K. Jiang, S. Wang, J. Wang, Z. Du, Z. Jiang, X. Zhang, B. Yuan, et al., Efficient pre-training of llms via topology-aware communication alignment on more than 9600 gpus, Advances in Neural Information Processing Systems 38 (2025) 147100–147126. [73] M. Liang, H. T. Kassa, W. Fu, B. Coutinho, L. Feng, C. Delimitrou, Lumos: Efficient performance modeling and estimation for large-scale llm training, Proceedings of Machine Learning and Systems 7. [74] B. Workshop, T. L. Scao, A. Fan, C. Akiki, E. Pavlick, S. Ilić, D. Hesslow, R. Castagné, A. S. Luccioni, F. Yvon, et al., Bloom: A 176b-parameter open-access multilingual language model, arXiv preprint arXiv:2211.05100. [75] A. Chowdhery, S. Narang, J. Devlin, M. Bosma, G. Mishra, A. Roberts, P. Barham, H. W. Chung, C. Sutton, S. Gehrmann, et al., Palm: Scaling language modeling with pathways, Journal of machine learning research 24 (240) (2023) 1–113. [76] G. Team, P. Georgiev, V. I. Lei, R. Burnell, L. Bai, A. Gulati, G. Tanzer, D. Vincent, Z. Pan, S. Wang, et al., Gemini 1.5: Unlocking multimodal understanding across millions of tokens of context, arXiv preprint arXiv:2403.05530. [77] X. Bi, D. Chen, G. Chen, S. Chen, D. Dai, C. Deng, H. Ding, K. Dong, Q. Du, Z. Fu, et al., Deepseek llm: Scaling open-source language models with longtermism (2024). arXiv:2401.02954. URL https://arxiv.org/abs/2401.02954 [78] C. Raffel, N. Shazeer, A. Roberts, K. Lee, S. Narang, M. Matena, Y. Zhou, W. Li, P. J. Liu, Exploring the limits of transfer learning with a unified text-to-text transformer, Journal of machine learning research 21 (140) (2020) 1–67. [79] J. Li, A. Fang, et al., DataComp-LM: In search of the next generation of training sets for language models, in: Advances in Neural Information Processing Systems, Vol. 37, Curran Associates, Inc., 2024, pp. 14200–14282. [80] W. Xiong, J. Liu, I. Molybog, H. Zhang, P. Bhargava, R. Hou, L. Martin, R. Rungta, K. A. Sankararaman, B. Oguz, et al., Effective long-context scaling of foundation models, in: Proceedings of the 2024 Conference of the North American Chapter of the Association for Computational Linguistics: Human Language Technologies (Volume 1: Long Papers), 2024, pp. 4643–4663. [81] A. Grattafiori, A. Dubey, A. Jauhri, A. Pandey, A. Kadian, A. Al-Dahle, A. Letman, A. Mathur, A. Schelten, A. Vaughan, et al., The llama 3 herd of models, arXiv preprint arXiv:2407.21783. [82] L. Soldaini, R. Kinney, A. Bhagia, D. Schwenk, D. Atkinson, R. Authur, B. Bogin, K. Chandu, J. Dumas, Y. Elazar, et al., Dolma: An open corpus of three trillion tokens for language model pretraining research, in: Proceedings of the 62nd Annual Meeting of the Association for Computational Linguistics (Volume 1: Long Papers), 2024, pp. 15725–15788. [83] G. Penedo, H. Kydlíček, A. Lozhkov, M. Mitchell, C. A. Raffel, L. Von Werra, T. Wolf, et al., The fineweb datasets: Decanting the web for the finest text data at scale, Advances in Neural Information Processing Systems 37 (2024) 30811–30849. [84] M. Ali, M. Fromm, K. Thellmann, R. Rutmann, M. Lübbering, J. Leveling, K. Klug, J. Ebert, N. Doll, J. Buschhoff, C. Jain, A. Weber, L. Jurkschat, H. Abdelwahab, C. John, P. Ortiz Suarez, M. Ostendorff, S. Weinbach, R. Sifa, S. Kesselheim, N. Flores-Herr, Tokenizer choice for LLM training: Negligible or crucial?, in: K. Duh, H. Gomez, S. Bethard (Eds.), Findings of the Association for Computational Linguistics: NAACL 2024, Association for Computational Linguistics, Mexico City, Mexico, 2024, pp. 3907–3924. doi:10.18653/v1/2024.findings-naacl.247. [85] D. Liang, H. Gonen, Y. Mao, R. Hou, N. Goyal, M. Ghazvininejad, L. Zettlemoyer, M. Khabsa, Xlm-v: Overcoming the vocabulary bottleneck in multilingual masked language models, in: Proceedings of the 2023 Conference on Empirical Methods in Natural Language Processing, 2023, pp. 13142–13152. [86] E. Dodici, Analisi statistica della divina commedia, accessed: 2026-03-05 (2008). URL https://wp.me/p27dha-3M

35

[87] M. Lewis, Y. Liu, N. Goyal, M. Ghazvininejad, A. Mohamed, O. Levy, V. Stoyanov, L. Zettlemoyer, Bart: Denoising sequence-to-sequence pre-training for natural language generation, translation, and comprehension, in: Proceedings of the 58th annual meeting of the association for computational linguistics, 2020, pp. 7871–7880. [88] J. Devlin, M.-W. Chang, K. Lee, K. Toutanova, Bert: Pre-training of deep bidirectional transformers for language understanding, in: Proceedings of the 2019 conference of the North American chapter of the association for computational linguistics: human language technologies, volume 1 (long and short papers), 2019, pp. 4171–4186. [89] Y. Liu, M. Ott, N. Goyal, J. Du, M. Joshi, D. Chen, O. Levy, M. Lewis, L. Zettlemoyer, V. Stoyanov, Roberta: A robustly optimized bert pretraining approach, arXiv preprint arXiv:1907.11692. [90] A. Q. Jiang, A. Sablayrolles, A. Mensch, C. Bamford, D. S. Chaplot, D. d. l. Casas, F. Bressand, G. Lengyel, G. Lample, L. Saulnier, et al., Mistral 7b, arXiv preprint arXiv:2310.06825. [91] S. Raschka, Llm architecture gallery, accessed: 2026-03-17 (2024). URL https://sebastianraschka.com/llm-architecture-gallery/ [92] J. Su, M. Ahmed, Y. Lu, S. Pan, W. Bo, Y. Liu, Roformer: Enhanced transformer with rotary position embedding, Neurocomputing 568 (2024) 127063. [93] S. Zhang, S. Roller, N. Goyal, M. Artetxe, M. Chen, S. Chen, C. Dewan, M. Diab, X. Li, X. V. Lin, et al., Opt: Open pre-trained transformer language models, arXiv preprint arXiv:2205.01068. [94] S. Smith, M. Patwary, B. Norick, P. LeGresley, S. Rajbhandari, J. Casper, Z. Liu, S. Prabhumoye, G. Zerveas, V. Korthikanti, et al., Using deepspeed and megatron to train megatron-turing nlg 530b, a large-scale generative language model, arXiv preprint arXiv:2201.11990. [95] T. Dao, D. Fu, S. Ermon, A. Rudra, C. Ré, Flashattention: Fast and memory-efficient exact attention with io-awareness, Advances in neural information processing systems 35 (2022) 16344–16359. [96] T. Dao, Flashattention-2: Faster attention with better parallelism and work partitioning, in: International Conference on Learning Representations, Vol. 2024, 2024, pp. 35549–35562. [97] A. Katharopoulos, A. Vyas, N. Pappas, F. Fleuret, Transformers are rnns: Fast autoregressive transformers with linear attention, in: International conference on machine learning, PMLR, 2020, pp. 5156–5165. [98] N. Shazeer, Fast transformer decoding: One write-head is all you need, arXiv preprint arXiv:1911.02150. [99] J. Ainslie, J. Lee-Thorp, M. de Jong, Y. Zemlyanskiy, F. Lebrón, S. Sanghai, Gqa: Training generalized multi-query transformer models from multi-head checkpoints (2023). arXiv:2305.13245. URL https://arxiv.org/abs/2305.13245 [100] Y. Sun, Z. Li, Y. Zhang, T. Pan, B. Dong, Y. Guo, J. Wang, Efficient attention mechanisms for large language models: A survey, arXiv preprint arXiv:2507.19595. [101] N. Shazeer, A. Mirhoseini, K. Maziarz, A. Davis, Q. Le, G. Hinton, J. Dean, Outrageously large neural networks: The sparsely-gated mixture-of-experts layer, OpenReview.net. [102] D. Lepikhin, H. Lee, Y. Xu, D. Chen, O. Firat, Y. Huang, M. Krikun, N. Shazeer, Z. Chen, Gshard: Scaling giant models with conditional computation and automatic sharding, arXiv preprint arXiv:2006.16668. [103] W. Fedus, B. Zoph, N. Shazeer, Switch transformers: Scaling to trillion parameter models with simple and efficient sparsity, Journal of Machine Learning Research 23 (120) (2022) 1–39. [104] N. Du, Y. Huang, A. M. Dai, S. Tong, D. Lepikhin, Y. Xu, M. Krikun, Y. Zhou, A. W. Yu, O. Firat, et al., Glam: Efficient scaling of language models with mixture-of-experts, in: International conference on machine learning, PMLR, 2022, pp. 5547–5569. [105] H. Zhang, D. Morwani, N. Vyas, J. Wu, D. Zou, U. Ghai, D. Foster, S. Kakade, How does critical batch size scale in pre-training?, in: International Conference on Learning Representations, Vol. 2025, 2025, pp. 66756–66782.

36

[106] A. Karpathy, train_gpt2.py, accessed: 2026-05-06 (2024). URL https://github.com/karpathy/build-nanogpt/blob/master/train_gpt2.py [107] NVIDIA, H100 gpu, accessed: 2026-04-27 (2026). URL https://www.nvidia.com/en-eu/data-center/h100/ [108] V. A. Korthikanti, J. Casper, S. Lym, L. McAfee, M. Andersch, M. Shoeybi, B. Catanzaro, Reducing activation recomputation in large transformer models, Proceedings of Machine Learning and Systems 5 (2023) 341–353. [109] Z. R. K. Rostam, S. Szénási, G. Kertész, Achieving peak performance for large language models: A systematic review, IEEE access 12 (2024) 96017–96050. [110] Y. M. Tsung, P. Qi, M. Lin, X. Wan, Balancing pipeline parallelism with vocabulary parallelism, Proceedings of Machine Learning and Systems 7. [111] L. Won, Tensor parallelism by hand, accessed: 2026-05-11 (2024). URL https://dev.to/lewis_won/tensor-parallelism-by-hand-3eh [112] S. Singh, P. Singhania, A. K. Ranjan, Z. Sating, A. Bhatele, A 4d hybrid algorithm to scale parallel training to thousands of gpus, arXiv preprint arXiv:2305.13525. [113] F. Zeng, W. Gan, Y. Wang, P. S. Yu, Distributed training of large language models: A survey, Natural Language Processing Journal (2025) 100174. [114] C. Chen, X. Li, Q. Zhu, J. Duan, P. Sun, X. Zhang, C. Yang, Centauri: Enabling efficient scheduling for communication-computation overlap in large model training via communication partitioning, in: Proceedings of the 29th ACM International Conference on Architectural Support for Programming Languages and Operating Systems, Volume 3, 2024, pp. 178–191. [115] W. Wang, M. Ghobadi, K. Shakeri, Y. Zhang, N. Hasani, Rail-only: A low-cost high-performance network for training llms with trillion parameters, in: 2024 IEEE Symposium on High-Performance Interconnects (HOTI), IEEE, 2024, pp. 1–10. [116] Z. Wang, A. Cai, X. Xie, Z. Pan, Y. Guan, W. Chu, J. Wang, S. Li, J. Huang, C. Cai, Y. Hao, Y. Ding, WLB-LLM: Workload-Balanced 4d parallelism for large language model training, in: 19th USENIX Symposium on Operating Systems Design and Implementation (OSDI 25), USENIX Association, 2025, pp. 785–801. [117] Z. Jia, M. Zaharia, A. Aiken, Beyond data and model parallelism for deep neural networks., Proceedings of Machine Learning and Systems 1 (2019) 1–13. [118] C. Unger, Z. Jia, W. Wu, S. Lin, M. Baines, C. E. Q. Narvaez, V. Ramakrishnaiah, N. Prajapati, P. McCormick, J. Mohd-Yusof, et al., Unity: Accelerating {DNN} training through joint optimization of algebraic transformations and parallelization, in: 16th USENIX Symposium on Operating Systems Design and Implementation (OSDI 22), 2022, pp. 267–284. [119] X. Miao, Y. Wang, Y. Jiang, C. Shi, X. Nie, H. Zhang, B. Cui, Galvatron: Efficient transformer training over multiple gpus using automatic parallelism, Proceedings of the VLDB Endowment 16 (2023) 470–479. [120] G. Liu, Y. Miao, Z. Lin, X. Shi, S. Maleki, F. Yang, Y. Bao, S. Wang, Aceso: Efficient parallel dnn training through iterative bottleneck alleviation, in: Proceedings of the Nineteenth European Conference on Computer Systems, 2024, pp. 163–181. [121] W. Chu, X. Xie, J. Yu, J. Wang, A. Phanishayee, C. Tang, Y. Hao, J. Huang, M. Ozdal, J. Wang, et al., Scaling llama 3 training with efficient parallelism strategies, in: Proceedings of the 52nd Annual International Symposium on Computer Architecture, 2025, pp. 1703–1716. [122] Amd rocm collective communication library, accessed: 2026-06-24 (2025). URL https://github.com/ROCm/rccl [123] PyTorch Contributors, Gloo: Collective Communications Library, accessed: 2026-06-24 (2026). URL https://github.com/pytorch/gloo

37

[124] Intel Corporation, Intel oneAPI Collective Communications Library, accessed: 2026-06-24 (2026). URL https://www.intel.com/content/www/us/en/developer/tools/oneapi/oneccl.html [125] A. Shah, A. Jangda, et al., MSCCL++: Rethinking GPU communication abstractions for AI inference, in: Proceedings of the 31st ACM International Conference on Architectural Support for Programming Languages and Operating Systems, Association for Computing Machinery, 2026, pp. 1201–1215. [126] G. Xu, Z. Le, Y. Chen, Z. Lin, Z. Jin, Y. Miao, C. Li, {AutoCCL}: Automated collective communication tuning for accelerating distributed and parallel {DNN} training, in: 22nd USENIX Symposium on Networked Systems Design and Implementation (NSDI 25), 2025, pp. 667–683. [127] R. Thakur, R. Rabenseifner, W. Gropp, Optimization of collective communication operations in mpich, The International Journal of High Performance Computing Applications 19 (1) (2005) 49–66. [128] E. Warraich, O. Shabtai, K. Manaa, S. Vargaftik, Y. Piasetzky, M. Kadosh, L. Suresh, M. Shahbaz, Optireduce: Resilient and tail-optimal allreduce for distributed deep learning in the cloud, in: 22nd USENIX Symposium on Networked Systems Design and Implementation (NSDI 25), 2025, pp. 685–703. [129] Q. Meng, H. Zheng, Z. Zhang, C. Lao, C. Huang, B. Li, Z. Zhu, H. Lu, W. Dang, Z. Lin, et al., Astral: A datacenter infrastructure for large language model training at scale, in: Proceedings of the ACM SIGCOMM 2025 Conference, 2025, pp. 609–625. [130] T. Zhong, J. Zhao, Q. Su, G. Fox, Youmu: Efficient columnar data pipeline for llm training, Proceedings of Machine Learning and Systems 7. [131] B. Wan, M. Han, Y. Sheng, Y. Peng, H. Lin, M. Zhang, Z. Lai, M. Yu, J. Zhang, Z. Song, et al., {ByteCheckpoint}: A unified checkpointing system for large foundation model development, in: 22nd USENIX Symposium on Networked Systems Design and Implementation (NSDI 25), 2025, pp. 559–578. [132] J. Lin, Z. Jiang, Z. Song, S. Zhao, M. Yu, Z. Wang, C. Wang, Z. Shi, X. Shi, W. Jia, et al., Understanding stragglers in large model training using what-if analysis, in: 19th USENIX Symposium on Operating Systems Design and Implementation (OSDI 25), 2025, pp. 483–498. [133] D. D. Sharma, Pci express 6.0 specification: A low-latency, high-bandwidth, high-reliability, and costeffective interconnect with 64.0 gt/s pam-4 signaling, IEEE Micro 41 (1) (2020) 23–29. [134] R. Neugebauer, G. Antichi, J. F. Zazo, Y. Audzevich, S. López-Buedo, A. W. Moore, Understanding pcie performance for end host networking, in: Proceedings of the 2018 Conference of the ACM Special Interest Group on Data Communication, 2018, pp. 327–341. [135] Altera, AN 829: PCI Express Avalon-MM DMA Reference Design, accessed: 2026-05-22 (2018). URL https://docs.altera.com/r/docs/683554/current [136] J. Lawley, Understanding Performance of PCI Express Systems, Tech. rep., Xilinx, accessed: 2026-05-22 (2014). URL https://docs.amd.com/v/u/en-US/wp350 [137] E. F. Kfoury, S. Choueiri, A. Mazloum, A. AlSabeh, J. Gomez, J. Crichigno, A comprehensive survey on smartnics: Architectures, development models, applications, and research directions, IEEE Access 12 (2024) 107297–107336. [138] T. Gross, Communication in iwarp systems, in: Proceedings of the 1989 ACM/IEEE conference on Supercomputing, 1989, pp. 436–445. [139] Z. Yan, D. Li, L. Chen, D. Xiong, K. Gao, Y. Zhang, R. Yan, M. Zhang, B. Zhang, Z. Jiang, et al., From atop to zcube: Automated topology optimization pipeline and a highly cost-effective network topology for large model training, in: Proceedings of the ACM SIGCOMM 2025 Conference, 2025, pp. 861–881. [140] Y. Wei, T. Hu, C. Liang, Y. Cui, Communication optimization for distributed training: architecture, advances, and opportunities, IEEE Network 39 (3) (2024) 241–248. [141] W. Wang, M. Khazraee, Z. Zhong, M. Ghobadi, Z. Jia, D. Mudigere, Y. Zhang, A. Kewitsch, {TopoOpt}: Co-optimizing network topology and parallelization strategy for distributed training jobs, in: 20th USENIX Symposium on Networked Systems Design and Implementation (NSDI 23), 2023, pp. 739–767. 38

[142] J. Simon, D. Kunin, A. Atanasov, E. Boix-Adserà, B. Bordelon, J. Cohen, N. Ghosh, F. Guth, A. Jacot, M. Kamb, et al., There will be a scientific theory of deep learning, arXiv preprint arXiv:2604.21691. [143] Z. Gan, R. Ren, W. Yao, X. Hu, G. Xu, C. Qian, H. Tang, Z. Gong, X. Yao, P. Tang, et al., Beyond the black box: Theory and mechanism of large language models, arXiv preprint arXiv:2601.02907. [144] R. Zhao, T. Qin, D. Alvarez-Melis, S. Kakade, N. Saphra, Random scaling for emergent capabilities, arXiv preprint arXiv:2502.17356. [145] G. Wang, G. Baker, A. Gordon, D. Murfet, Embryology of a language model, arXiv preprint arXiv:2508.00331. [146] P. Qi, X. Wan, N. Amar, M. Lin, Pipeline parallelism with controllable memory, Advances in Neural Information Processing Systems 37 (2024) 46539–46566. [147] S. A. Jacobs, M. Tanaka, C. Zhang, M. Zhang, S. L. Song, S. Rajbhandari, Y. He, Deepspeed ulysses: System optimizations for enabling training of extreme long sequence transformer models, arXiv preprint arXiv:2309.14509. [148] S. Li, F. Xue, C. Baranwal, Y. Li, Y. You, Sequence parallelism: Long sequence training from system perspective, in: Proceedings of the 61st Annual Meeting of the Association for Computational Linguistics (Volume 1: Long Papers), 2023, pp. 2391–2404. [149] J. Fang, S. Zhao, Usp: A unified sequence parallelism approach for long context generative ai, arXiv preprint arXiv:2405.07719. [150] H. Liu, M. Zaharia, P. Abbeel, Ringattention with blockwise transformers for near-infinite context, in: International Conference on Learning Representations, Vol. 2024, 2024, pp. 3992–4008. [151] D. Li, R. Shao, A. Xie, E. P. Xing, X. Ma, I. Stoica, J. E. Gonzalez, H. Zhang, Distflashattn: Distributed memory-efficient attention for long-context llms training, arXiv preprint arXiv:2310.03294. [152] S. Go, J. Park, S. More, H. Wu, I. Wang, A. Jezghani, T. Krishna, D. Mahajan, Characterizing the efficiency of distributed training: A power, performance, and thermal perspective, in: Proceedings of the 58th IEEE/ACM International Symposium on Microarchitecture, 2025, pp. 626–642. [153] X. Pan, W. Lin, L. Zhang, S. Shi, Z. Tang, R. Wang, B. Li, X. Chu, Fsmoe: A flexible and scalable training system for sparse mixture-of-experts models, in: Proceedings of the 30th ACM International Conference on Architectural Support for Programming Languages and Operating Systems, Volume 1, 2025, pp. 524–539. [154] Y. Lei, D. Lee, L. Zhao, D. Kurniawan, C. Kim, H. Jeong, C. Kim, H. Choi, L. Yu, A. Krishnamurthy, et al., Flash: Fast all-to-all communication in gpu clusters, arXiv e-prints (2025) arXiv–2505.

39

Record · ID 978383 · SHA-256 0b44296d356ce3ce
Retrieved via Conceptio — every document is proof-bundled with source, license, and retrieval metadata.