本篇主题
这两节Lecture都是关于并行训练的,主要内容包括:
- 理解训练超大模型时系统层面的复杂性
- 掌握不同的 parallelization paradigms,以及为什么人们通常会同时使用多种并行方式
- 了解大规模训练任务通常是如何进行的
Lecture 7 - Parallelism 1
单块GPU无法满足SCALING的需求,必须使用多块GPU进行训练,所以我们需要Multi-GPU、Multi-Machine的并行训练方法,如下图所示:

首先讲了一些Basic Concepts:
- All Reduce:通信操作,所有参与的进程都将输入数据进行规约(如求和、最大值等)后,将结果分发给所有进程。常用于分布式训练中的梯度同步。
- Broadcast:通信操作,数据从一个进程发送到所有其他进程。常用于分布式训练中的分发模型参数或者初始化数据。
- All Gather:通信操作,所有参与的进程将各自的数据发送给所有其他进程,最终每个进程都获得所有数据的集合。常用于分布式训练中的收集模型输出或者中间结果。和All Reduce的区别在于,All Reduce会对数据进行规约操作(如求和),而All Gather只是简单地收集数据,不进行任何计算。
- Reduce Scatter:通信操作,所有参与的进程将各自的数据发送给所有其他进程,并对数据进行规约操作(如求和),最终每个进程都获得规约后的结果的一部分。常用于分布式训练中的分布式梯度更新。

了解了初始知识之后,可以正式进入核心部分,不同的并行方式:
- Data Parallelism:最早的并行方式,模型复制到每个GPU上,每个GPU处理不同的数据batch,计算梯度后进行同步更新。优点是实现简单;缺点是通信开销大,尤其是模型参数较大时。 假设计算一个 SGD:我们会把 B 大小下的 batch 分给 M 个不同的机器,然后交换梯度去同步计算。这种情况对于 Compute scaling,每个 GPU 计算 B/M 个数据(不错!);对于 Communication overhead,每个 batch 需要转移两次梯度(发送和接受);对于 Memory scaling,完全没有,每个 GPU 都需要复制一次模型参数。 早期的data parallelism问题是只缓解了计算压力,内存压力没有缓解,需要每个GPU都复制一份模型参数,通信压力也很大。 Zero[HTTPS://arxiv.org/pdf/1910.02054]可以解决这个问题,核心的思想是将模型的state(参数和优化器状态)分布在不同的GPU上,每个GPU只存储模型的一部分,这样可以显著减少每个GPU的内存占用,同时通过通信操作来同步更新模型参数。 ZeRO分三个阶段:
- ZeRO-1:主要聚焦于 optimizer state sharding。把 optimizer state(first + second moments)分给每个 GPU,每个 GPU 都有 parameters + gradients,负责更新一部分 params。 每个 GPU 根据分配到的 batch 子集计算完整的梯度,此时只拥有局部梯度 利用 Reduce-Scatter 将所有 GPU 的局部梯度汇总并分散到每个设备上,现在每个 GPU 只持有全局梯度的一部分 每个 GPU 利用局部梯度和 optimizer state 更新该部分参数 利用 AllGather 将所有 GPU 的部分参数收集并分发给所有设备,确保每个 GPU 拥有完整的参数。 这种方法的通信开销并没有增加,而且 memory 减少了接近四倍(优化器状态是fp32,参数是fp16)。

- ZeRO-2:ZeRO 的第二阶段利用一阶段的思想(增加通信和计算),聚焦于 gradient sharding.这样操作的复杂性在于我们无法实例化一个完整的梯度向量,但每个 GPU 必须计算完整的梯度(因为 data parallel)。如何操作呢: 每张卡计算完某一层的梯度后,立即进行 reduce 操作,将对应梯度加总并发给负责该参数的那一张卡,不需要每张卡都存全部梯度,一旦梯度不再用于反向传播,立即释放内存 各 GPU 用自己负责的梯度(对应参数的梯度,与第一阶段的局部梯度不同) + optimizer state 更新参数 参数 All-Gather 同步

- ZeRO-3:ZeRO第三阶段自然就到了share everything了,参数、梯度、优化器状态都分布在不同的GPU上。每个GPU只存储模型的一部分参数和对应的梯度以及优化器状态,这样可以极大地减少每个GPU的内存占用,同时通过通信操作来同步更新模型参数。(这块太复杂暂时没看懂)


知乎文章的实践经验:第三阶段的 ZeRO 明显减少了很多显存的需求,但由于增加了通信, 等待的时间明显变久了。
- Model Parallelism:模型并行将模型的不同部分分布在不同的GPU上,每个GPU负责计算模型的一部分。优点是可以训练更大的模型;缺点是实现复杂,通信开销大。 model parallelism 可以在不改变 batch size 的情况下 Scaling up in memory,它会像 zero3 一样把参数分布给 GPU,但是 communicate activations(zero3 sends params)。因为 cut up model 的形式不同,所以总共有两类:Pipeline parallel 和 Tensor parallel.
- Pipeline parallel:很自然的我们会想到模型有很多 layer 构成,我们把 layers cut up 然后分布到 GPU,交换 activations。但这样的话利用率会非常低,每个 GPU 大部分时间都在空闲,等待反向传播的完成。我们可以尝试在一个bubble发送完之后马上计算下一个 bubble 的 forward,这样就可以提升利用率了。但是我们需要很大的 batch size。现在也有一些工作比如‘Zero bubble’pipelining 在尝试推进这方面工作,但 Pipeline parallel 在实际操作中其实是非常复杂的。


| 特性 | Pipeline Parallel (气泡问题) | Tensor Parallel (这张图) |
|---|---|---|
| 并行方式 | 时间上的流水线 | 空间上的数据并行 |
| Batch Size 要求 | 需要很大(填流水线) | 灵活,小 batch 也能跑 |
| 通信量 | 小(只传激活值) | 大(频繁 All-Reduce) |
| 扩展性 | 可扩展到 1000+ GPUs | 通常只跨 2-8 个 GPU(通信瓶颈) |
- Sequence parallel:在训练过程中,模型参数只占用一部分显存,激活值也会占用大量显存。一种与 Tensor Parallel 互补 的技术,专门用来解决 Activation Memory(激活内存)爆炸 的问题。 核心思想:把 LayerNorm、Dropout 这些与参数无关、只在序列维度上独立操作的层,沿 sequence axis(序列维度)切分到不同 GPU 上。
| Tensor Parallel | Sequence Parallel | |
|---|---|---|
| 切分对象 | 矩阵乘法(Linear)、Attention | LayerNorm、Dropout、GeLU(逐元素操作) |
| 切分维度 | 隐藏层维度(hidden dim) | 序列长度维度(seq len) |
| 计算特点 | 需要 All-Reduce 聚合结果 | 各 GPU 独立计算,仅需 All-Gather 收集 |

最后主要讲了下其他模型的使用情况,比如 DeepSeek V3 – ZeRO stage 1 with Tensor, Sequence, and Pipeline parallel(16)等,Llama3 405B 和 Gemma 2。
Lecture 8 - Parallelism 2
多Node多GPU结构回顾:在前面的Leture中讲过,我跳过了先学的Lecture7-8,后边补一个链接。

这节课主要把Lecutre提到的概念用代码实现,主要分为两个 part:
- Part 1: building blocks of distributed communication/computation,利用pytorch封装好的函数用一些例子来展示lecture7中提到的概念,比如all-reduce、broadcast、all-gather、reduce-scatter等
- Part 2: distributed training,主要展示了如何在pytorch中实现分布式训练,主要是data parallel和tenor parallel的实现。
Part 1
因为我没有在本地模拟多GPU环境,所以只贴一下代码。
# All-Reduce Example
torch.tensor([1.0, 2.0, 3.0]).cuda() # 每个 GPU 上的张量
dist.all_reduce(tensor, op=dist.ReduceOp.SUM) # 将所有 GPU 上的张量求和并分发回每个 GPU
print(tensor) # 每个 GPU 上的张量现在都是 [num_gpus, num_gpus*2, num_gpus*3]
# 输入
tensor([1., 2., 3.], device=『cuda:0』)
tensor([1., 2., 3.], device=『cuda:1』)
tensor([1., 2., 3.], device=『cuda:2』)
# 输出
tensor([4., 8., 12.], device=『cuda:0』)
tensor([4., 8., 12.], device=『cuda:1』)
tensor([4., 8., 12.], device=『cuda:2』)
# reduce-scatter Example
tensor = torch.tensor([1.0, 2.0, 3.0]).cuda() # 每个 GPU 上的张量
dist.reduce_scatter(tensor, tensor_list=[torch.tensor([1.0, 2.0, 3.0]).cuda() for _ in range(dist.get_world_size())], op=dist.ReduceOp.SUM) # 将所有 GPU 上的张量求和并分散到每个 GPU 上
print(tensor) # 每个 GPU 上的张量现在是 [num_gpus, num_gpus*2, num_gpus*3] 中的一部分
# 输入
tensor([1., 2., 3.], device=『cuda:0』)
tensor([1., 2., 3.], device=『cuda:1』)
tensor([1., 2., 3.], device=『cuda:2』)
# 输出
tensor([4.], device=『cuda:0』) # GPU 0 上的张量是所有 GPU 上的张量求和后分散到 GPU 0 上的一部分
tensor([8.], device=『cuda:1』) # GPU 1 上的张量是所有 GPU 上的张量求和后分散到 GPU 1 上的一部分
tensor([12.], device=『cuda:2』) # GPU 2 上的张量是所有 GPU 上的张量求和后分散到 GPU 2 上的一部分
# all-gather Example
tensor = torch.tensor([1.0, 2.0, 3.0]).cuda() # 每个 GPU 上的张量
dist.all_gather(tensor_list=[torch.tensor([0.0, 0.0, 0.0]).cuda() for _ in range(dist.get_world_size())], tensor=tensor) # 将每个 GPU 上的张量收集到所有 GPU 上
print(tensor_list) # 每个 GPU 上的张量列表现在包含所有 GPU 上的张量
# 输入是all-gather的输出
# 输出
tensor([1., 2., 3.], device=『cuda:0』) # GPU 0 上的张量是所有 GPU 上的张量列表中的第一个张量
tensor([1., 2., 3.], device=『cuda:1』) # GPU 1 上的张量是所有 GPU 上的张量列表中的第二个张量
tensor([1., 2., 3.], device=『cuda:2』) # GPU 2 上的张量是所有 GPU 上的张量列表中的第三个张量
通过这个例子可以观察到,All-Reduce 和 Reduce-Scatter 的区别在于,All-Reduce 会将所有 GPU 上的张量求和并分发回每个 GPU,而 Reduce-Scatter 会将所有 GPU 上的张量求和并分散到每个 GPU 上的一部分。All-Gather 则是将每个 GPU 上的张量收集到所有 GPU 上。并且,all-reduce = reduce-scatter + all-gather。
Part 2
在实际分布式训练中,我们通常会使用 PyTorch 的 torch.nn.parallel.DistributedDataParallel(DDP)来实现分布式训练。但是如果要我们自己实现分布式训练,对于前向过程和反向过程的通信操作,我们需要使用前面提到的通信原语来同步参数和梯度。
# 只实现模型前向和反向传播的通信循环过程-- 以 Data Parallel 为例
# dataparallel的核心思想是每个GPU上都有一份完整的模型参数,每个GPU处理不同的数据batch,计算梯度后进行同步更新(进行All-Reduce)。下面是一个简化的示例代码:
def data_parallel_main(rank:int,world_size:int,data:torch.Tensor,num_layers: int, num_steps: int):
for step in range(num_steps):
# Forward pass
x = data
for param in params:
x = x @ param
x = F.gelu(x)
loss = x.square().mean() # Loss function is average squared magnitude
# Backward pass
loss.backward()
# 同步各个工作节点的梯度(与标准训练的唯一区别在于 DDP)
for param in params:
dist.all_reduce(tensor=param.grad, op=dist.ReduceOp.AVG, async_op=False)
# 更新参数
optimizer.step()
# 而对于model parallel(tensor parallel),每个GPU只负责模型的一部分参数,前向需要跨GPU通信。每个 GPU 获取部分 layer(submatrix),传输所有的 data 和 activation。代码实现如下(把一个 num_dim * num_dim 的矩阵分割为 GPU 数量个 num_dim * local_num_dim)。
def tensor_parallel_main(rank:int,world_size:int,data:torch.Tensor,num_layers: int, num_steps: int):
local_num_dim = num_dim // world_size
# 每个 GPU 负责模型的一部分参数-不是基于模型长度的切分,而是基于模型宽度的切分
local_params = [torch.randn(local_num_dim, local_num_dim).cuda() for _ in range(num_layers)]
for step in range(num_steps):
# Forward pass
x = data
for param in local_params:
x = x @ param # 每个 GPU 计算自己负责的部分
x = F.gelu(x)
loss = x.square().mean() # Loss function is average squared magnitude
# 收集激活值并拼接得到完整的激活值
activations = [torch.empty(batch_size, local_num_dim, device=get_device(rank)) for _ in range(world_size)]
dist.all_gather(tensor_list=activations, tensor=x, async_op=False)
x = torch.cat(activations, dim=-1) # 拼接得到完整的激活值
其他
在实际的分布式训练中,我们通常会使用 PyTorch 的 torch.nn.parallel.DistributedDataParallel(DDP)来实现分布式训练。DDP 会自动处理前向和反向传播中的通信操作,简化了分布式训练的实现。不过,了解这些原理也是不错的,可以帮助我们更好地理解分布式训练的底层机制,以及在某些特殊情况下进行自定义优化。
现在已经有了很多分布式并行训练的库和框架,比如 DeepSpeed、Megatron-LM、FairScale 等,在之前的goole deepreaserch总结的文章中也有提到,这些库和框架提供了更高层次的抽象,简化了分布式训练的实现,同时也提供了很多优化策略来提升训练效率。