简介
Accelerate 可以通过仅添加四行代码来分布式配置运行相同的 PyTorch 代码。
代码示例:
+ from accelerate import Accelerator
+ accelerator = Accelerator()
+ model, optimizer, training_dataloader, scheduler = accelerator.prepare(
+ model, optimizer, training_dataloader, scheduler
+ )
for batch in training_dataloader:
optimizer.zero_grad()
inputs, targets = batch
inputs = inputs.to(device)
targets = targets.to(device)
outputs = model(inputs)
loss = loss_function(outputs, targets)
+ accelerator.backward(loss)
optimizer.step()
scheduler.step()
# 命令行启动
accelerate launch {my_script.py}
安装
1. pip 安装
pip install accelerate
2. conda 安装
conda install -c conda-forge accelerate
3. GitHub 远程安装
pip install git+https://github.com/huggingface/accelerate
4. GitHub 本地安装
git clone https://github.com/huggingface/accelerate
cd accelerate
pip install -e .
5. 通过命令行进行配置
# 命令行进行参数配置
accelerate config
# 加载默认配置
python -c "from accelerate.utils import write_basic_config; write_basic_config(mixed_precision='fp16')"
# 参数检查
accelerate env
配置文件样例:
- `Accelerate` version: 1.2.0.dev0
- Platform: Linux-6.8.0-47-generic-x86_64-with-glibc2.35
- `accelerate` bash location: /home/zach/miniconda3/envs/accelerate/bin/accelerate
- Python version: 3.10.13
- Numpy version: 1.26.4
- PyTorch version (GPU?): 2.5.1+cu124 (True)
- PyTorch XPU available: False
- PyTorch NPU available: False
- PyTorch MLU available: False
- PyTorch MUSA available: False
- System RAM: 187.91 GB
- GPU type: NVIDIA GeForce RTX 4090
- `Accelerate` default config:
- compute_environment: LOCAL_MACHINE
- distributed_type: MULTI_GPU
- mixed_precision: no
- use_cpu: False
- debug: False
- num_processes: 2
- machine_rank: 0
- num_machines: 1
- gpu_ids: all
- rdzv_backend: static
- same_network: True
- main_training_function: main
- enable_cpu_affinity: False
- downcast_bf16: no
- tpu_use_cluster: False
- tpu_use_sudo: False
- tpu_env: []
快速教程
1. Accelerate 主要功能
- 配置
- 训练
- 推理
2. 启动界面
a. 配置
Accelerate 通过命令生成默认配置文件,自动为分布式训练框架选择合适的配置。
# 命令行进行参数配置
accelerate config
命令会在缓存文件夹中创建并保存一个 default_config.yaml 文件。
# 配置好环境后,启动测试脚本检查分布式环境
accelerate test
如果配置文件不在缓存文件夹内,可以通过 --config_file 参数指定配置文件位置。
# 启动训练脚本
accelerate launch path_to_script.py --args_for_the_script
b. 训练
对 PyTorch 脚本进行调整,以便在多个 GPU 或 TPU 运行。
示例代码:
+ from accelerate import Accelerator
+ accelerator = Accelerator()
+ device = accelerator.device
+ model, optimizer, training_dataloader, scheduler = accelerator.prepare(
+ model, optimizer, training_dataloader, scheduler
+ )
for batch in training_dataloader:
optimizer.zero_grad()
inputs, targets = batch
- inputs = inputs.to(device)
- targets = targets.to(device)
outputs = model(inputs)
loss = loss_function(outputs, targets)
+ accelerator.backward(loss)
optimizer.step()
scheduler.step()
说明:
-
训练脚本中导入并实例化
Accelerate类,Accelerate类会初始化分布式训练参数,根据代码的启动方式自动检测训练环境,例如单 GPU、多 GPU、多 TPU 等。from accelerate import Accelerator accelerator = Accelerator() -
Accelerate类会自动将模型和数据放置在合适的设备上。device = accelerator.device -
将优化器、模型、数据加载器、学习率调度器传递给
prepare()方法,此方法将模型包装在针对分布式优化的容器中,使用 Accelerate 版本的优化器和调度器,创建分片版本的数据加载器,以便在 GPU 或 TPU 设置之间分发。model, optimizer, train_dataloader, lr_scheduler = accelerator.prepare( model, optimizer, train_dataloader, lr_scheduler ) -
替换原有的
loss.backward()方法进行训练。accelerator.backward(loss)
评估指标
# 将验证数据加载器传递给 prepare() 方法,进行模型评估
validation_dataloader = accelerator.prepare(validation_dataloader)
# 分布式训练下,每个设备只接收部分评估数据,需要使用 gather_for_metrics() 方法收集所有的预测值和目标值
# 如果张量在每个进程中大小不同,需要使用 pad_accross_processes() 方法将张量填充到最大长度
# 注意,这里的张量必须是一维
for inputs, targets in validation_dataloader:
predictions = model(inputs)
# 收集所有预测和目标
all_predictions, all_targets = accelerator.gather_for_metrics((predictions, targets))
# 使用 *Datasets.Metric*
metric.add_batch(all_predictions, all_targets)
3. 推理
针对模型推理,有两个主要功能:init_empty_weights() 和 load_checkpoint_and_dispatch(),加载无法直接放入内存的大型推理模型。
a. 空权重初始化
init_empty_weights():加载一个空的模型,占用内存比完全加载模型和权重要小的多。
# 示例代码:加载一个空的 Mixtral-8x7B 模型
from accelerate import init_empty_weights
from transformers import AutoConfig, AutoModelForCausalLM
config = AutoConfig.from_pretrained("mistralai/Mixtral-8x7B-Instruct-v0.1")
with init_empty_weights():
model = AutoModelForCausalLM.from_config(config)
b. 加载权重
load_checkpoint_and_dispatch() 函数将完整或分片的检查点加载到空模型中,并自动在所有可用设备上分配参数。
# 示例代码
from accelerate import load_checkpoint_and_dispatch
model_checkpoint = "your-local-model-folder"
model = load_checkpoint_and_dispatch(
model, checkpoint=model_checkpoint, device_map="auto", no_split_module_classes=['Block']
)
说明:
device_map="auto"表示让框架决定每个模型层的放置位置,GPU 或者 CPU 上,如果内存不足,则以内存映射张量的形式放在硬盘上。no_split_module_classes参数指示哪些模块不进行拆分,通常是具有残差连接的模块。
将 Accelerate 添加到代码
1. 基本的 PyTorch 训练代码样例
device = "cuda"
model.to(device)
for batch in training_dataloader:
optimizer.zero_grad()
inputs, targets = batch
inputs = inputs.to(device)
targets = targets.to(device)
outputs = model(inputs)
loss = loss_function(outputs, targets)
loss.backward()
optimizer.step()
scheduler.step()
2. Accelerate
from accelerate import Accelerator
accelerator = Accelerator()
# 由框架决定放置在哪种设备
- device = "cuda"
+ device = accelerator.device
model.to(device)
3. 准备 PyTorch 对象
model, optimizer, training_dataloader, scheduler = accelerator.prepare(
model, optimizer, training_dataloader, scheduler
)
4. 训练循环
- inputs = inputs.to(device)
- targets = targets.to(device)
outputs = model(inputs)
loss = loss_function(outputs, targets)
- loss.backward()
+ accelerator.backward(loss)
5. 完整代码
from accelerate import Accelerator
accelerator = Accelerator()
device = accelerator.device
model, optimizer, training_dataloader, scheduler = accelerator.prepare(
model, optimizer, training_dataloader, scheduler
)
for batch in training_dataloader:
optimizer.zero_grad()
inputs, targets = batch
outputs = model(inputs)
loss = loss_function(outputs, targets)
accelerator.backward(loss)
optimizer.step()
scheduler.step()
6. 梯度累积
梯度累积功能允许在更新权重之前累积多个批次的梯度,从而以更大的批次进行训练。
# 通过 gradient_accumulation_steps 参数,进行设置
+ accelerator = Accelerator(gradient_accumulation_steps=2)
model, optimizer, training_dataloader = accelerator.prepare(model, optimizer, training_dataloader)
for input, label in training_dataloader:
+ with accelerator.accumulate(model):
predictions = model(input)
loss = loss_function(predictions, label)
accelerator.backward(loss)
optimizer.step()
scheduler.step()
optimizer.zero_grad()
7. 梯度裁剪
梯度裁剪防止”梯度爆炸”,提供了两种方式进行梯度裁剪:
clip_grad_value():将梯度剪裁为最小值和最大值clip_grad_norm():用于将梯度归一化为某个值
8. 混合精度
混合精度通过 fp16 半精度等较低精度的数据类型进行梯度计算,从而加速训练。
为了获得 Accelerate 的最佳性能,损失应该在模型内部计算,因为模型外部的计算是以全精度进行的。
+ accelerator = Accelerator(mixed_precision="fp16")
+ with accelerator.autocast():
loss = complex_loss_function(outputs, target)
9. 保存并加载
保存模型之前使用 unwrap_model() 方法解包模型,因为 prepare() 方法会将模型包装到分布式训练的接口中。如果不解包模型,保存模型状态字典也会保存较大模型中可能存在的任何额外层,并且将无法将权重加载回原有模型。
保存模型时应使用 save_model() 方法解包并保存模型状态字典。此方法还可以将模型保存为分片 checkpoint 或 safetensors 格式。
# 保存为单一 checkpoint 文件
accelerator.wait_for_everyone()
accelerator.save_model(model, save_directory)
# 加载单一 checkpoint 模型:使用 unwrap_model() 方法解包模型,然后加载权重
unwrapped_model = accelerator.unwrap_model(model)
path_to_checkpoint = os.path.join(save_directory, "pytorch_model.bin")
unwrapped_model.load_state_dict(torch.load(path_to_checkpoint))
# 保存为分片 checkpoint 文件
# 设置 safe_serialization=True 为以 safetensor 格式保存模型
accelerator.wait_for_everyone()
accelerator.save_model(model, save_directory, max_shard_size="1GB", safe_serialization=True)
# 加载分片 checkpoint 模型:使用 load_checkpoint_in_model() 方法,允许将模型加载到特定设备上
load_checkpoint_in_model(unwrapped_model, save_directory, device_map={"": device})
10. 状态
训练期间,可能希望保存模型、优化器、学习率调度器的当前状态,以便后续在同一脚本中恢复它们。
使用 save_state() 和 load_state() 方法,保存和加载状态。
执行过程
在分布式训练中,管理多 GPU 执行进程很重要。有些进程完成速度快,而有些进程在其他进程尚未完成时不启动。
框架提供了用于协调进程执行时间的工具,以确保所有内容在所有设备上保持同步。
1. 在所有进程上执行
有些代码只需要在给定的机器上运行一次,例如打印日志、本地主进程显示进度。
# 使用 accelerator.is_local_main_process 来指示只应执行一次的代码
from tqdm.auto import tqdm
progress_bar = tqdm(range(args.max_train_steps), disable=not accelerator.is_local_main_process)
# 使用装饰器来指示在所有进程中只应执行一次的函数
@accelerator.on_local_main_process
def do_my_thing():
"Something done once per server"
do_thing_once_per_server()
# 使用 accelerator.is_main_process 来指示在所有进程中只应执行一次的代码
if accelerator.is_main_process:
repo.push_to_hub()
# 使用装饰器来指示在所有进程中只应执行一次的函数
@accelerator.on_main_process
def do_my_thing():
"Something done once per server"
do_thing_once()
2. 在特定进程上执行
执行仅应在特定进程或本地进程索引上执行的功能。
# 使用 on_process() 方法并指定要执行函数的进程索引
@accelerator.on_process(process_index=0)
def do_my_thing():
"Something done on process index 0"
do_thing_on_index_zero()
# 使用 on_local_process() 方法并指定要执行函数的本地进程索引
@accelerator.on_local_process(local_process_idx=0)
def do_my_thing():
"Something done on process index 0 on each server"
do_thing_on_index_zero_on_each_server()
3. 推迟执行
同时在多个 GPU 上运行脚本时,某些代码的执行速度可能比其他代码更快。可能需要等待所有进程达到某一步骤后再执行下一组指令。例如,在确保每个进程都完成训练之后保存模型。
# 使用 wait_for_everyone() 阻塞所有已完成的进程,直到所有进程都到达同一点(对于单 GPU 或 CPU 此方法无效)
accelerator.wait_for_everyone()
TPU 训练
TPU(张量处理单元)是一种专为高效训练模型而设计的硬件。Accelerate 支持 TPU 训练。
1. 编译
TPU 会创建一个包含训练步骤中所有操作(例如前向传递、后向传递和优化器步骤)的计算图。第一步训练需要构建和编译计算图。一旦编译完成,后续步骤的速度都会加快。
注意:避免重复编译代码,否则训练速度非常缓慢:
- 批次中的所有张量必须具有相同的长度(例如,对于 NLP 任务,没有动态填充)
- 代码必须是静态(例如,没有根据输入而具有不同长度的 for 循环层,例如 LSTM)
2. 权重绑定
常见的语言模型是将嵌入层和 Softmax 层的权重绑定在一起,但是,将模型放置 TPU 会破坏权重绑定,需要重新绑定权重。
# 在脚本中为 TPU 进行权重绑定,需要设置 distributed_type 参数,然后调用 tie_weights() 方法:
if accelerator.distributed_type == DistributedType.TPU:
model.tie_weights()
启动脚本
1. 加速脚本
# 完整代码
from accelerate import Accelerator
accelerator = Accelerator()
model, optimizer, training_dataloader, scheduler = accelerator.prepare(
model, optimizer, training_dataloader, scheduler
)
for batch in training_dataloader:
optimizer.zero_grad()
inputs, targets = batch
outputs = model(inputs)
loss = loss_function(outputs, targets)
accelerator.backward(loss)
optimizer.step()
scheduler.step()
# 将上述代码封装为函数,作为脚本调用
from accelerate import Accelerator
+ def main():
accelerator = Accelerator()
model, optimizer, training_dataloader, scheduler = accelerator.prepare(
model, optimizer, training_dataloader, scheduler
)
for batch in training_dataloader:
optimizer.zero_grad()
inputs, targets = batch
outputs = model(inputs)
loss = loss_function(outputs, targets)
accelerator.backward(loss)
optimizer.step()
scheduler.step()
+ if __name__ == "__main__":
+ main()
2. 启动脚本
# 命令行中启动
accelerate launch {script_name.py} --arg1 --arg2 ...
# 单 GPU 上训练
# for cuda device:
CUDA_VISIBLE_DEVICES="0" accelerate launch {script_name.py} --arg1 --arg2 ...
# for xpu device:
ZE_AFFINITY_MASK="0" accelerate launch {script_name.py} --arg1 --arg2 ...
# 使用全部 GPU,禁用混合精度训练
accelerate launch --multi_gpu {script_name.py} {--arg1} {--arg2} ...
# 使用指定数量的 GPU 训练
accelerate launch --num_processes=2 {script_name.py} {--arg1} {--arg2} ...
# 使用全部 GPU,启用混合精度训练
accelerate launch --multi_gpu --mixed_precision=fp16 --num_processes=2 {script_name.py} {--arg1} {--arg2} ...
# 查看帮助文档
accelerate launch -h
# 使用 torchrun 在多 GPU 上训练
MIXED_PRECISION="fp16" torchrun --nproc_per_node=2 --nnodes=1 {script_name.py} {--arg1} {--arg2} ...
# 将启动脚本作为 Python 模块,以便配置更多 Python 启动参数
python -m accelerate.commands.launch --num_processes=2 {script_name.py} {--arg1} {--arg2}
# 在 CPU 上测试代码
accelerate launch --cpu {script_name.py} {--arg1} {--arg2}
# 在 accelerate config 进行参数配置后,可以直接使用 accelerate launch 启动
accelerate launch {script_name.py} {--arg1} {--arg2} ...
3. 自定义配置
accelerate config 命令会将配置文件 default_config.yaml 保存到缓存文件夹中,以 accelerate 结尾:
- 环境变量
HF_HOME目录,以accelerate为后缀 - 如果
HF_HOME不存在,则保存到XDG_CACHE_HOME以huggingface/accelerate为后缀 - 如果
XDG_CACHE_HOME也不存在,则文件夹~/.cache/huggingface/accelerate
yaml 配置文件示例:fp16 混合精度训练
compute_environment: LOCAL_MACHINE
deepspeed_config: {}
distributed_type: MULTI_GPU
fsdp_config: {}
machine_rank: 0
main_process_ip: null
main_process_port: null
main_training_function: main
mixed_precision: fp16
num_machines: 1
num_processes: 2
use_cpu: false
# 使用 yaml 配置文件,启动脚本
accelerate launch --config_file {path/to/config/my_config_file.yaml} {script_name.py} {--arg1} {--arg2} ...
4. 多节点训练
启动多节点训练运行,执行以下操作:
- 将代码和数据复制到所有节点。(或将它们放在共享文件系统上)
- 在所有节点上配置 Python 环境。
- 首先在主单节点上运行
accelerate config。指定节点数后,系统会要求指定每个节点的等级(主控节点的等级为 0),以及主进程的 IP 地址和端口。用于工作节点与主进程通信。之后,将此配置文件复制或发送到所有节点,并将等级更改machine_rank为 1、2、3 等,或者在其他非主节点运行accelerate config进行设置。
完成此操作后,在所有节点上运行 accelerate launch 或 torchrun 开始多节点训练运行。
注意:所有节点都运行该命令才能启动分布式训练,而不仅仅是从主节点运行。可以借助 SLURM 或其他进程执行器,以便通过单个命令调起所有节点进程。
注意:为了降低延迟,建议使用主节点的内网 IP,而不是公网 IP。例如
192.168.x.x或者172.x.x.x,可以通过在主节点上执行hostname -I进行查看。
从 Notebook 启动脚本
1. 配置环境
# 命令行执行以下命令,生成配置文件,一般采用默认配置
accelerate config
# 通过 write_basic_config() 方法,将 GPU 配置写入配置文件
import os
from accelerate.utils import write_basic_config
write_basic_config() # Write a config file
os._exit(00) # Restart the notebook
注意:在多 GPU 环境下,CUDA 无法多次初始化,因此,在 Notebook 中调试后,最终训练时,需要全面清理并重启。
2. 准备数据和模型
# 导入
import os, re, torch, PIL
import numpy as np
from torch.optim.lr_scheduler import OneCycleLR
from torch.utils.data import DataLoader, Dataset
from torchvision.transforms import Compose, RandomResizedCrop, Resize, ToTensor
from accelerate import Accelerator
from accelerate.utils import set_seed
from timm import create_model
# 创建函数获取文件名
import os
data_dir = "../../images"
fnames = os.listdir(data_dir)
fname = fnames[0]
print(fname)
# 从文件名中获取标签
import re
def extract_label(fname):
stem = fname.split(os.path.sep)[-1]
return re.search(r"^(.*)_\d+\.jpg$", stem).groups()[0]
# 创建数据集类型
class PetsDataset(Dataset):
def __init__(self, file_names, image_transform=None, label_to_id=None):
self.file_names = file_names
self.image_transform = image_transform
self.label_to_id = label_to_id
def __len__(self):
return len(self.file_names)
def __getitem__(self, idx):
fname = self.file_names[idx]
raw_image = PIL.Image.open(fname)
image = raw_image.convert("RGB")
if self.image_transform is not None:
image = self.image_transform(image)
label = extract_label(fname)
if self.label_to_id is not None:
label = self.label_to_id[label]
return {"image": image, "label": label}
# 构建数据集
fnames = [os.path.join("../../images", fname) for fname in fnames if fname.endswith(".jpg")]
# 获取所有标签
all_labels = [extract_label(fname) for fname in fnames]
id_to_label = list(set(all_labels))
id_to_label.sort()
label_to_id = {lbl: i for i, lbl in enumerate(id_to_label)}
# 创建 get_dataloaders() 函数获取数据加载器
def get_dataloaders(batch_size: int = 64):
"Builds a set of dataloaders with a batch_size"
random_perm = np.random.permutation(len(fnames))
cut = int(0.8 * len(fnames))
train_split = random_perm[:cut]
eval_split = random_perm[cut:]
# For training a simple RandomResizedCrop will be used
train_tfm = Compose([RandomResizedCrop((224, 224), scale=(0.5, 1.0)), ToTensor()])
train_dataset = PetsDataset([fnames[i] for i in train_split], image_transform=train_tfm, label_to_id=label_to_id)
# For evaluation a deterministic Resize will be used
eval_tfm = Compose([Resize((224, 224)), ToTensor()])
eval_dataset = PetsDataset([fnames[i] for i in eval_split], image_transform=eval_tfm, label_to_id=label_to_id)
# Instantiate dataloaders
train_dataloader = DataLoader(train_dataset, shuffle=True, batch_size=batch_size, num_workers=4)
eval_dataloader = DataLoader(eval_dataset, shuffle=False, batch_size=batch_size * 2, num_workers=4)
return train_dataloader, eval_dataloader
# 导入学习率调度器
from torch.optim.lr_scheduler import CosineAnnealingLR
3. 编写训练函数
基本训练循环:
def training_loop(mixed_precision="fp16", seed: int = 42, batch_size: int = 64):
set_seed(seed)
accelerator = Accelerator(mixed_precision=mixed_precision)
# 创建数据加载器、模型
train_dataloader, eval_dataloader = get_dataloaders(batch_size)
model = create_model("resnet50d", pretrained=True, num_classes=len(label_to_id))
# 冻结模型的编码器,微调模型的分类头
for param in model.parameters():
param.requires_grad = False
for param in model.get_classifier().parameters():
param.requires_grad = True
# 数据标准化
mean = torch.tensor(model.default_cfg["mean"])[None, :, None, None]
std = torch.tensor(model.default_cfg["std"])[None, :, None, None]
# 放置在加速设备
mean = mean.to(accelerator.device)
std = std.to(accelerator.device)
# 实例化其余 PyTorch 类
optimizer = torch.optim.Adam(params=model.parameters(), lr=3e-2 / 25)
lr_scheduler = OneCycleLR(optimizer=optimizer, max_lr=3e-2, epochs=5, steps_per_epoch=len(train_dataloader))
# 传递给 accelerator.prepare()
model, optimizer, train_dataloader, eval_dataloader, lr_scheduler = accelerator.prepare(
model, optimizer, train_dataloader, eval_dataloader, lr_scheduler
)
# 训练模型
for epoch in range(5):
model.train()
for batch in train_dataloader:
inputs = (batch["image"] - mean) / std
outputs = model(inputs)
loss = torch.nn.functional.cross_entropy(outputs, batch["label"])
accelerator.backward(loss)
optimizer.step()
lr_scheduler.step()
optimizer.zero_grad()
# 模型评估代码
model.eval()
accurate = 0
num_elems = 0
for batch in eval_dataloader:
inputs = (batch["image"] - mean) / std
with torch.no_grad():
outputs = model(inputs)
predictions = outputs.argmax(dim=-1)
# 注意:分布式训练模型评估,预测和标签需要通过 gather() 方法传递,以便所有数据在当前设备上可用,并获得正确的计算结果
accurate_preds = accelerator.gather(predictions) == accelerator.gather(batch["label"])
num_elems += accurate_preds.shape[0]
accurate += accurate_preds.long().sum()
# 计算指标,使用 print() 在主进程上打印
eval_metric = accurate.item() / num_elems
accelerator.print(f"epoch {epoch}: {100 * eval_metric:.2f}")
训练循环的完整版本:
def training_loop(mixed_precision="fp16", seed: int = 42, batch_size: int = 64):
set_seed(seed)
# Initialize accelerator
accelerator = Accelerator(mixed_precision=mixed_precision)
# Build dataloaders
train_dataloader, eval_dataloader = get_dataloaders(batch_size)
# Instantiate the model (you build the model here so that the seed also controls new weight initializations)
model = create_model("resnet50d", pretrained=True, num_classes=len(label_to_id))
# Freeze the base model
for param in model.parameters():
param.requires_grad = False
for param in model.get_classifier().parameters():
param.requires_grad = True
# You can normalize the batches of images to be a bit faster
mean = torch.tensor(model.default_cfg["mean"])[None, :, None, None]
std = torch.tensor(model.default_cfg["std"])[None, :, None, None]
# To make these constants available on the active device, set it to the accelerator device
mean = mean.to(accelerator.device)
std = std.to(accelerator.device)
# Instantiate the optimizer
optimizer = torch.optim.Adam(params=model.parameters(), lr=3e-2 / 25)
# Instantiate the learning rate scheduler
lr_scheduler = OneCycleLR(optimizer=optimizer, max_lr=3e-2, epochs=5, steps_per_epoch=len(train_dataloader))
# Prepare everything
# There is no specific order to remember, you just need to unpack the objects in the same order you gave them to the
# prepare method.
model, optimizer, train_dataloader, eval_dataloader, lr_scheduler = accelerator.prepare(
model, optimizer, train_dataloader, eval_dataloader, lr_scheduler
)
# Now you train the model
for epoch in range(5):
model.train()
for batch in train_dataloader:
inputs = (batch["image"] - mean) / std
outputs = model(inputs)
loss = torch.nn.functional.cross_entropy(outputs, batch["label"])
accelerator.backward(loss)
optimizer.step()
lr_scheduler.step()
optimizer.zero_grad()
model.eval()
accurate = 0
num_elems = 0
for batch in eval_dataloader:
inputs = (batch["image"] - mean) / std
with torch.no_grad():
outputs = model(inputs)
predictions = outputs.argmax(dim=-1)
accurate_preds = accelerator.gather(predictions) == accelerator.gather(batch["label"])
num_elems += accurate_preds.shape[0]
accurate += accurate_preds.long().sum()
eval_metric = accurate.item() / num_elems
# Use accelerator.print to print only on the main process.
accelerator.print(f"epoch {epoch}: {100 * eval_metric:.2f}")
4. 使用 notebook_launcher
# 导入
from accelerate import notebook_launcher
args = ("fp16", 42, 64)
notebook_launcher(training_loop, args, num_processes=2)
# 在分布式训练下,需要在每个节点上启用一个 jupyter 会话,并同时运行启动单元
# 假设当前分布式训练,包含 2 个节点,每个节点有 8 个 GPU 设备,主机 IP 为 172.31.43.8,主机上的 notebook_launcher 设置如下:
notebook_launcher(training_loop, args, master_addr="172.31.43.8", node_rank=0, num_nodes=2, num_processes=8)
# 另一台机器上的第二个 jupyter 会话中,注意 node_rank 与主机不同:
notebook_launcher(training_loop, args, master_addr="172.31.43.8", node_rank=1, num_nodes=2, num_processes=8)
# 如果在 TPU 上运行,则启动函数如下所示:
model = create_model("resnet50d", pretrained=True, num_classes=len(label_to_id))
args = (model, "fp16", 42, 64)
notebook_launcher(training_loop, args, num_processes=8)
# 如果启动弹性训练过程并启用容错功能,则启动函数如下所示:
notebook_launcher(
training_loop,
args,
num_processes=2,
max_restarts=3
)
运行时输出:
Launching training on 2 GPUs.
epoch 0: 88.12
epoch 1: 91.73
epoch 2: 92.58
epoch 3: 93.90
epoch 4: 94.71
5. 调试
notebook_launcher() 启动时报错:CUDA has already been initialized
错误排查: 说明之前导入过 torch.cuda 模块,可能是由于 Notebook 中 cuda 配置无法多次初始化导致。
为方便排查,可以在启动 notebook_launcher() 时,指定参数 ACCELERATE_DEBUG_MODE=yes,可以在启动时进行额外检查。
加速 - 快速教程
最新训练代码、配置文件、启动脚本可以通过官方网址查询:https://huggingface.co/docs/accelerate/usage_guides/explore
1. 训练代码
a. 通用样例
未集成 Accelerate 之前的训练代码如下所示:
# 训练代码
for batch in dataloader:
optimizer.zero_grad()
inputs, targets = batch
inputs = inputs.to(device)
targets = targets.to(device)
outputs = model(inputs)
loss = loss_function(outputs, targets)
loss.backward()
optimizer.step()
scheduler.step()
b. 基础版本
说明:不需要再调用 model.to(device) 或 inputs.to(device),后续由 accelerator.prepare() 自动完成。
+ from accelerate import Accelerator
+ accelerator = Accelerator()
+ dataloader, model, optimizer, scheduler = accelerator.prepare(
+ dataloader, model, optimizer, scheduler
+ )
for batch in dataloader:
inputs, targets = batch
- inputs = inputs.to(device)
- targets = targets.to(device)
outputs = model(inputs)
loss = loss_function(outputs, targets)
- loss.backward()
+ accelerator.backward(loss)
optimizer.step()
scheduler.step()
optimizer.zero_grad()
c. 计算模型性能指标
说明:使用 Accelerator.gather_for_metrics 方法从所有设备收集数据,然后根据收集到的数据计算模型性能指标。
import evaluate
+ from accelerate import Accelerator
+ accelerator = Accelerator()
+ train_dataloader, eval_dataloader, model, optimizer, scheduler = (
+ accelerator.prepare(
+ train_dataloader, eval_dataloader,
+ model, optimizer, scheduler
+ )
+ )
metric = evaluate.load("accuracy")
for batch in train_dataloader:
inputs, targets = batch
- inputs = inputs.to(device)
- targets = targets.to(device)
outputs = model(inputs)
loss = loss_function(outputs, targets)
loss.backward()
optimizer.step()
scheduler.step()
optimizer.zero_grad()
model.eval()
for batch in eval_dataloader:
inputs, targets = batch
- inputs = inputs.to(device)
- targets = targets.to(device)
with torch.no_grad():
outputs = model(inputs)
predictions = outputs.argmax(dim=-1)
+ predictions, references = accelerator.gather_for_metrics(
+ (predictions, references)
+ )
metric.add_batch(
predictions = predictions,
references = references
)
print(metric.compute())
d. 保存检查点
说明:Accelerator 提供了 save_state 和 load_state 方法,用于保存或加载检查点。
from accelerate import Accelerator
accelerator = Accelerator()
dataloader, model, optimizer, scheduler = accelerator.prepare(
dataloader, model, optimizer, scheduler
)
for batch in dataloader:
inputs, targets = batch
outputs = model(inputs)
loss = loss_function(outputs, targets)
accelerator.backward(loss)
optimizer.step()
scheduler.step()
optimizer.zero_grad()
+ accelerator.save_state("checkpoint_dir")
+ accelerator.load_state("checkpoint_dir")
e. 实验跟踪器
说明:要使用实验跟踪器,只需在构建 Accelerator 对象时将所需的跟踪器传递给 log_with 即可。然后通过 Accelerator.init_trackers() 初始化跟踪器,传入相关配置。然后调用 Accelerator.log 将日志信息写入到跟踪器。训练结束时,调用 accelerator.end_training() 来停止跟踪器。
from accelerate import Accelerator
- accelerator = Accelerator()
+ accelerator = Accelerator(log_with="wandb")
train_dataloader, model, optimizer, scheduler = accelerator.prepare(
dataloader, model, optimizer, scheduler
)
+ accelerator.init_trackers()
model.train()
for batch in train_dataloader:
inputs, targets = batch
outputs = model(inputs)
loss = loss_function(outputs, targets)
+ accelerator.log({"loss":loss})
accelerator.backward(loss)
optimizer.step()
scheduler.step()
optimizer.zero_grad()
+ accelerator.end_training()
f. 梯度累积
说明:将训练循环封装在 Accelerator.accumulate 上下文管理器中,训练时,梯度将自动累积和同步。
from accelerate import Accelerator
accelerator = Accelerator(
+ gradient_accumulation_steps=2,
)
dataloader, model, optimizer, scheduler = accelerator.prepare(
dataloader, model, optimizer, scheduler
)
for batch in dataloader:
+ with accelerator.accumulate(model):
inputs, targets = batch
outputs = model(inputs)
loss = loss_function(outputs, targets)
accelerator.backward(loss)
optimizer.step()
scheduler.step()
optimizer.zero_grad()
2. 配置与启动
a. AWS SageMaker
i. YAML 配置文件
+ base_job_name: accelerate-sagemaker-1
+ compute_environment: AMAZON_SAGEMAKER
distributed_type: 'NO'
dynamo_backend: 'NO'
+ ec2_instance_type: ml.p3.2xlarge
+ gpu_ids: all
+ iam_role_name: MY_IAM_ROLE_NAME
mixed_precision: 'no'
+ num_machines: 1
+ profile: MY_PROFILE_NAME
+ py_version: py38
+ pytorch_version: 1.10.2
+ region: us-east-1
+ transformers_version: 4.17.0
use_cpu: false
ii. 训练脚本调整
def parse_args():
parser = argparse.ArgumentParse(
description="sample task"
)
parser.add_argument(
"--some_bool_arg",
- action="store_true",
+ type=bool,
+ default=False,
)
iii. 启动命令
# 使用 Accelerate config 进行配置
accelerate launch {script_name.py} {--arg1} {--arg2} ...
# 使用 ~/config.yaml 文件
accelerate launch --config_file ~/config.yaml {script_name.py} {--arg1} {--arg2} ...
b. DeepSpeed
i. YAML 配置文件
compute_environment: LOCAL_MACHINE
+ deepspeed_config:
+ gradient_accumulation_steps: 1
+ gradient_clipping: 1.0
+ offload_optimizer_device: cpu
+ offload_param_device: cpu
+ zero3_init_flag: true
+ zero3_save_16bit_model: true
+ zero_stage: 3
+ distributed_type: DEEPSPEED
downcast_bf16: 'no'
dynamo_backend: 'NO'
fsdp_config: {}
machine_rank: 0
main_training_function: main
megatron_lm_config: {}
mixed_precision: fp16
+ num_machines: 1
+ num_processes: 8
rdzv_backend: static
same_network: true
use_cpu: false
ii. 训练脚本调整
from accelerate import Accelerator
def main():
accelerator = Accelerator()
model, optimizer, training_dataloader, scheduler = accelerator.prepare(
model, optimizer, training_dataloader, scheduler
)
generated_tokens = accelerator.unwrap_model(model).generate(
batch["input_ids"],
attention_mask=batch["attention_mask"],
**gen_kwargs,
+ synced_gpus=True
)
...
accelerator.unwrap_model(model).save_pretrained(
args.output_dir,
is_main_process=accelerator.is_main_process,
save_function=accelerator.save,
+ state_dict=accelerator.get_state_dict(model)
)
...
iii. 启动命令
# 使用 accelerate config 完成配置
accelerate launch {script_name.py} {--arg1} {--arg2} ...
# 使用 ~/config.yaml 文件
accelerate launch --config_file ~/config.yaml {script_name.py} {--arg1} {--arg2} ...
# 无配置文件,改为使用参数
accelerate launch \
--use_deepspeed \
--num_processes=8 \
--mixed_precision=fp16 \
--zero_stage=3 \
--gradient_accumulation_steps=1 \
--gradient_clipping=1 \
--zero3_init_flag=True \
--zero3_save_16bit_model=True \
--offload_optimizer_device=cpu \
--offload_param_device=cpu \
{script_name.py} {--arg1} {--arg2} ...
c. Megatron-LM
i. YAML 配置文件
compute_environment: LOCAL_MACHINE
deepspeed_config: {}
+ distributed_type: MEGATRON_LM
downcast_bf16: 'no'
dynamo_backend: 'NO'
fsdp_config: {}
machine_rank: 0
main_training_function: main
+ megatron_lm_config:
+ megatron_lm_gradient_clipping: 1.0
+ megatron_lm_num_micro_batches: 2
+ megatron_lm_pp_degree: 2
+ megatron_lm_recompute_activations: true
+ megatron_lm_sequence_parallelism: true
+ megatron_lm_tp_degree: 2
+ megatron_lm_use_distributed_optimizer: true
mixed_precision: bf16
num_machines: 1
num_processes: 8
rdzv_backend: static
same_network: true
use_cpu: false
ii. 训练脚本调整
from accelerate import Accelerator
+ from accelerate.utils import MegatronLMDummyScheduler
accelerator = Accelerator()
...
- lr_scheduler = get_scheduler(
- name=args.lr_scheduler_type,
- ...
- )
+ lr_scheduler = MegatronLMDummyScheduler(
+ optimizer=optimizer,
+ num_warmup_steps=...,
+ num_training_steps=...,
+ )
model, optimizer, train_dataloader, eval_dataloader, lr_scheduler = accelerator.prepare(
model, optimizer, train_dataloader, eval_dataloader, lr_scheduler
)
total_batch_size = (
- args.per_device_train_batch_size * accelerator.num_processes * args.gradient_accumulation_steps
+ accelerator.state.megatron_lm_plugin.global_batch_size
)
# in evaluation loop
for step, batch in enumerate(eval_dataloader):
with torch.no_grad():
outputs = model(**batch)
loss = outputs.loss
- losses.append(accelerator.gather_for_metrics(loss.repeat(args.per_device_eval_batch_size)))
+ losses.append(loss) # For Megatron-LM, the losses are already averaged across the data parallel group
- losses = torch.cat(losses)
+ losses = torch.tensor(losses)
iii. 启动命令
# 使用 accelerate config 进行配置
accelerate launch {script_name.py} {--arg1} {--arg2} ...
# 使用 ~/config.yaml 配置文件
accelerate launch --config_file ~/config.yaml {script_name.py} {--arg1} {--arg2} ...
# 无配置文件,改为使用参数
accelerate launch \
--use_megatron_lm \
--num_processes=8 \
--mixed_precision=bf16 \
--megatron_lm_tp_degree=2 \
--megatron_lm_pp_degree=2 \
--megatron_lm_num_micro_batches=2 \
--megatron_lm_sequence_parallelism=true \
--megatron_lm_recompute_activations=true \
--megatron_lm_use_distributed_optimizer=true \
{script_name.py} {--arg1} {--arg2} ...
d. Multi GPU
i. YAML 配置文件
compute_environment: LOCAL_MACHINE
deepspeed_config: {}
+ distributed_type: MULTI_GPU
downcast_bf16: 'no'
dynamo_backend: 'NO'
fsdp_config: {}
+ gpu_ids: all
+ machine_rank: 0
main_training_function: main
megatron_lm_config: {}
mixed_precision: 'no'
+ num_machines: 1
+ num_processes: 4
+ rdzv_backend: static
+ same_network: true
use_cpu: false
ii. 训练脚本调整
无需调整。
iii. 启动命令
# 使用 accelerate config 进行配置
accelerate launch {script_name.py} {--arg1} {--arg2} ...
# 使用 ~/config.yaml 配置文件
accelerate launch --config_file ~/config.yaml {script_name.py} {--arg1} {--arg2} ...
# 无配置文件,改为使用参数
accelerate launch --multi_gpu --num_processes=4 {script_name.py} {--arg1} {--arg2} ...
e. Multi Node Multi GPU
i. YAML 配置文件
在主机上的配置文件:
compute_environment: LOCAL_MACHINE
deepspeed_config: {}
+ distributed_type: MULTI_GPU
downcast_bf16: 'no'
dynamo_backend: 'NO'
fsdp_config: {}
gpu_ids: all
+ machine_rank: 0
+ main_process_ip: 192.168.20.1
+ main_process_port: 8080
main_training_function: main
megatron_lm_config: {}
mixed_precision: 'no'
+ num_machines: 2
+ num_processes: 8
+ rdzv_backend: static
+ same_network: true
use_cpu: false
在另一台非主机上的配置文件:
compute_environment: LOCAL_MACHINE
deepspeed_config: {}
+ distributed_type: MULTI_GPU
downcast_bf16: 'no'
dynamo_backend: 'NO'
fsdp_config: {}
gpu_ids: all
- machine_rank: 0
+ machine_rank: 1
+ main_process_ip: 192.168.20.1
+ main_process_port: 8080
main_training_function: main
megatron_lm_config: {}
mixed_precision: 'no'
+ num_machines: 2
+ num_processes: 8
+ rdzv_backend: static
+ same_network: true
use_cpu: false
ii. 训练脚本调整
无需调整。
iii. 启动命令
注意:需要在集群中的每一台机器上执行启动脚本,如果不使用 yaml 配置文件,务必注意正确指定
machine_rank参数。
# 使用 accelerate config 进行配置
accelerate launch {script_name.py} {--arg1} {--arg2} ...
# 使用 ~/config.yaml 配置文件
accelerate launch --config_file ~/config.yaml {script_name.py} {--arg1} {--arg2} ...
# 无配置文件,改为使用参数,注意将 {node_number} 替换为适当的机器号(主机为 0,否则为 1 以上)
accelerate launch --multi_gpu --num_machines=2 --num_processes=8 --main_process_ip="192.168.20.1" --main_process_port=8080 \
--machine_rank={node_number} {script_name.py} {--arg1} {--arg2} ...
f. PyTorch FSDP
i. YAML 配置文件
compute_environment: LOCAL_MACHINE
deepspeed_config: {}
+ distributed_type: FSDP
downcast_bf16: 'no'
dynamo_backend: 'NO'
+ fsdp_config:
+ fsdp_auto_wrap_policy: TRANSFORMER_BASED_WRAP
+ fsdp_backward_prefetch_policy: BACKWARD_PRE
+ fsdp_offload_params: true
+ fsdp_sharding_strategy: 1
+ fsdp_state_dict_type: FULL_STATE_DICT
+ fsdp_transformer_layer_cls_to_wrap: T5Block
machine_rank: 0
main_training_function: main
megatron_lm_config: {}
mixed_precision: bf16
num_machines: 1
+ num_processes: 8
rdzv_backend: static
same_network: true
use_cpu: false
ii. 训练脚本调整
from accelerate import Accelerator
def main():
accelerator = Accelerator()
- model, optimizer, dataloader, scheduler = accelerator.prepare(
- model, optimizer, dataloader, scheduler
- )
+ model = accelerator.prepare(model)
+ # Optimizer can be any PyTorch optimizer class
+ optimizer = torch.optim.AdamW(params=model.parameters(), lr=lr)
+ optimizer, dataloader, scheduler = accelerator.prepare(
+ optimizer, dataloader, scheduler
+ )
...
accelerator.unwrap_model(model).save_pretrained(
args.output_dir,
is_main_process=accelerator.is_main_process,
save_function=accelerator.save,
+ state_dict=accelerator.get_state_dict(model)
)
...
iii. 启动命令
# 使用 accelerate config 进行配置
accelerate launch {script_name.py} {--arg1} {--arg2} ...
# 使用 ~/config.yaml 配置文件
accelerate launch --config_file ~/config.yaml {script_name.py} {--arg1} {--arg2} ...
# 无配置文件,改为使用参数
accelerate launch \
--use_fsdp \
--num_processes=8 \
--mixed_precision=bf16 \
--fsdp_sharding_strategy=1 \
--fsdp_auto_wrap_policy=TRANSFORMER_BASED_WRAP \
--fsdp_transformer_layer_cls_to_wrap=T5Block \
--fsdp_offload_params=true \
{script_name.py} {--arg1} {--arg2} ...
加速 - 模型内存估计器
可以在官方提供的网页上查询模型需要的内存大小:https://huggingface.co/docs/accelerate/usage_guides/model_size_estimator
在探索机器上使用的潜在模型时,很难知道当前显卡的内存中可以容纳多大的模型(例如,将模型加载到 CUDA 上)。为了解决这个问题,Accelerate 提供了一个命令行界面,用于评估当前设备内存所支持的模型。
1. 查询命令
在命令行输入:accelerate estimate-memory,只会将模型加载到 meta 设备内存中,不会将模型全部权重加载到内存中,可以测试更多参数的模型,查看模型占用内存大小。
例如计算 bert-base-cased 模型占用内存大小:
# 命令行输入如下命令:
accelerate estimate-memory bert-base-cased
这将下载 config.json,并在设备上加载模型元数据,报告模型占用空间大小:
| 数据类型 | 最大层 | 总大小 | 使用 Adam 进行训练 |
|---|---|---|---|
| float32 | 84.95 MB | 418.18 MB | 1.61 GB |
| float16 | 42.47 MB | 206.59 MB | 826.36 MB |
| int8 | 21.24 MB | 103.29 MB | 413.18 MB |
| int4 | 10.62 MB | 51.65 MB | 206.59 MB |
a. 指定库进行查询
如果无法自动确定模型的来源,例如 bert-base-cased,可以传入库名。
# 使用 transformers 库
accelerate estimate-memory HuggingFaceM4/idefics-80b-instruct --library_name transformers
模型加载时的内存使用情况:
| 数据类型 | 最大层 | 总大小 | 使用 Adam 进行训练 |
|---|---|---|---|
| float32 | 3.02 GB | 297.12 GB | 1.16 TB |
| float16 | 1.51 GB | 148.56 GB | 594.24 GB |
| int8 | 772.52 MB | 74.28 GB | 297.12 GB |
| int4 | 386.26 MB | 37.14 GB | 148.56 GB |
# 使用 timm 库
accelerate estimate-memory timm/resnet50.a1_in1k --library_name timm
模型加载时的内存使用情况:
| 数据类型 | 最大层 | 总大小 | 使用 Adam 进行训练 |
|---|---|---|---|
| float32 | 9.0 MB | 97.7 MB | 390.78 MB |
| float16 | 4.5 MB | 48.85 MB | 195.39 MB |
| int8 | 2.25 MB | 24.42 MB | 97.7 MB |
| int4 | 1.12 MB | 12.21 MB | 48.85 MB |
b. 指定数据类型
通过 --dtypes 参数指定数据类型。
accelerate estimate-memory bert-base-cased --dtypes float32 float16
模型加载时的内存使用情况:
| 数据类型 | 最大层 | 总大小 | 使用 Adam 进行训练 |
|---|---|---|---|
| float32 | 84.95 MB | 413.18 MB | 1.61 GB |
| float16 | 42.47 MB | 206.59 MB | 826.36 MB |
2. 注意事项
- 这个计算结果是说明加载模型所需的内存,而不是推理所需的内存。计算精度相差 1-0.1 之间,例如,CUDA 上以全精度加载时,
bert-base-cased实际加载需要 413.68 MB,计算器估算为 413.18 MB。 - 进行推理时,应当额外增加 20%。
加速 - 模型量化
1. 安装
# 安装 bitsandbytes 库
pip install bitsandbytes
# 从源码安装最新版本
pip install git+https://github.com/huggingface/accelerate.git
# 安装 minGPT 并 huggingface_hub 运行示例
git clone https://github.com/karpathy/minGPT.git
pip install minGPT/
pip install huggingface_hub
非 CUDA 设备安装,需要参考官方文档:https://huggingface.co/docs/bitsandbytes/main/en/installation#multi-backend
2. 工作原理
# 为了节省内存,先使用 init_empty_weights() 初始化一个空模型,以 minGPT 库的 GPT2 模型为例
from accelerate import init_empty_weights
from mingpt.model import GPT
model_config = GPT.get_default_config()
model_config.model_type = 'gpt2-xl'
model_config.vocab_size = 50257
model_config.block_size = 1024
with init_empty_weights():
empty_model = GPT(model_config)
# 获取模型权重,可以是 state_dict 文件(例如,pytorch_model.bin),可以是分片检查点的文件
from huggingface_hub import snapshot_download
weights_location = snapshot_download(repo_id="marcsun13/gpt2-xl-linear-sharded")
# 使用 BnbQuantizationConfig 设置量化配置
# 8 位量化配置示例
from accelerate.utils import BnbQuantizationConfig
bnb_quantization_config = BnbQuantizationConfig(load_in_8bit=True, llm_int8_threshold=6)
# 4 位量化配置示例
from accelerate.utils import BnbQuantizationConfig
bnb_quantization_config = BnbQuantizationConfig(
load_in_4bit=True,
bnb_4bit_compute_dtype=torch.bfloat16,
bnb_4bit_use_double_quant=True,
bnb_4bit_quant_type="nf4"
)
# 使用选定的配置量化模型,需要使用 load_and_quantize_model() 方法
from accelerate.utils import load_and_quantize_model
quantized_model = load_and_quantize_model(
empty_model,
weights_location=weights_location,
bnb_quantization_config=bnb_quantization_config
)
3. 保存和加载 8 位模型
# 使用 save_model() 保存 8 位量化模型,当前不支持 4 位模型的序列化
from accelerate import Accelerator
accelerate = Accelerator()
new_weights_location = "path/to/save_directory"
accelerate.save_model(quantized_model, new_weights_location)
quantized_model_from_saved = load_and_quantize_model(
empty_model,
weights_location=new_weights_location,
bnb_quantization_config=bnb_quantization_config,
device_map="auto"
)
4. 将模型卸载到 CPU 和磁盘
如果 GPU 上没有足够的空间存储整个模型,可以将部分模块卸载到 CPU 或磁盘上,这将在底层使用大模型推理。
- 对于 8 位量化,选定的模块将转换为 8 位精度。
- 对于 4 位量化,选定的模块将保留在
BnbQuantizationConfig中的torch_dtype,后续提供 4 位模型序列化以后,再提供 4 位模型的卸载。
传递一个自定义函数 device_map,即可将模块卸载到 CPU 或磁盘上,卸载模块将会在需要时再度加载到 GPU,示例代码如下:
device_map = {
"transformer.wte": 0,
"transformer.wpe": 0,
"transformer.drop": 0,
"transformer.h": "cpu",
"transformer.ln_f": "disk",
"lm_head": "disk",
}
5. 微调量化模型
量化模型无法进行纯 8 位或 4 位训练。但是,可以利用参数高效微调方法 (PEFT) 来训练这些模型,并在其基础上训练适配器等。
量化模型无法添加适配器,但是,借助 Transformers 支持,可以对量化模型进行微调。
6. 量化前后对比
Colab 上实现的一个量化演示:https://colab.research.google.com/drive/1T1pOgewAWVpR9gKpaEWw4orOrzPFb3yM?usp=sharing
GPT2-1.5B 模型检查点位于 FP32 中,使用 6GB 内存。量化后,8 位模块占用 1.6GB 内存,4 位模块占用 1.2GB 内存。
加速 - 实验追踪器
Accelerate 提供了一个通用追踪 API,可用于在脚本运行过程中记录有用的内容。Accelerator.log()
1. 集成追踪器
当前 Accelerate 支持 7 种追踪器:
- TensorBoard
- WandB
- CometML
- Aim
- MLFlow
- ClearML
- DVCLive
# 通过 log_with 参数指定要使用的追踪器:
from accelerate import Accelerator
from accelerate.utils import LoggerType
accelerator = Accelerator(log_with="all") # For all available trackers in the environment
accelerator = Accelerator(log_with="wandb")
accelerator = Accelerator(log_with=["wandb", LoggerType.TENSORBOARD])
# 使用 Accelerator.init_trackers() 初始化追踪器
hps = {"num_iterations": 5, "learning_rate": 1e-2}
accelerator.init_trackers("my_project", config=hps)
# 使用 Accelerator.log() 记录数据,step 可用于控制每训练几次之后记录一次日志
accelerator.log({"train_loss": 1.12, "valid_loss": 0.8}, step=1)
# 完成训练时,需要执行 Accelerator.end_training()
accelerator.end_training()
完整示例:
from accelerate import Accelerator
accelerator = Accelerator(log_with="all")
config = {
"num_iterations": 5,
"learning_rate": 1e-2,
"loss_function": str(my_loss_function),
}
accelerator.init_trackers("example_project", config=config)
my_model, my_optimizer, my_training_dataloader = accelerator.prepare(my_model, my_optimizer, my_training_dataloader)
device = accelerator.device
my_model.to(device)
for iteration in range(config["num_iterations"]):
for step, batch in enumerate(my_training_dataloader):
my_optimizer.zero_grad()
inputs, targets = batch
inputs = inputs.to(device)
targets = targets.to(device)
outputs = my_model(inputs)
loss = my_loss_function(outputs, targets)
accelerator.backward(loss)
my_optimizer.step()
accelerator.log({"training_loss": loss}, step=step)
accelerator.end_training()
# 如果追踪器需要指定目录,例如 TensorBoard,则将目录路径传递给 project_dir
accelerator = Accelerator(log_with="tensorboard", project_dir=".")
# 如果需要指定多个目录,可以使用 ProjectConfiguration
config = ProjectConfiguration(project_dir=".", logging_dir="another/directory")
accelerator = Accelerator(log_with="tensorboard", project_config=config)
2. 实现自定义追踪器
实现一个新的追踪器,可以继承 GeneralTracker 来创建一个新的跟踪器,必须实现以下三个函数及三个属性:
__init__:- 存储
run_name并初始化跟踪器 - 如果跟踪器在本地存储其数据(例如 TensorBoard),则使用
logging_dir参数指定
- 存储
store_init_configuration:- 使用字典初始化存储目录配置
log:- 接受一个字典和一个
step参数,在运行时记录日志
- 接受一个字典和一个
name(str):- 跟踪器的唯一标识
requires_logging_directory(bool):- 跟踪器是否需要指定
logging_dir
- 跟踪器是否需要指定
tracker:- 作为一个
@property函数来实现 - 返回使用的内置跟踪器对象,例如:wandb 对象
- 作为一个
# 代码示例:
from accelerate.tracking import GeneralTracker, on_main_process
from typing import Optional
import wandb
class MyCustomTracker(GeneralTracker):
name = "wandb"
requires_logging_directory = False
@on_main_process
def __init__(self, run_name: str):
self.run_name = run_name
run = wandb.init(self.run_name)
@property
def tracker(self):
return self.run.run
@on_main_process
def store_init_configuration(self, values: dict):
wandb.config(values)
@on_main_process
def log(self, values: dict, step: Optional[int] = None):
wandb.log(values, step=step)
# 实例化 Accelerator 对象后,将跟踪器实例传递给参数 log_with
tracker = MyCustomTracker("some_run_name")
accelerator = Accelerator(log_with=tracker)
# 自定义跟踪器也可以与现有跟踪器混合使用,例如指定使用跟踪器的参数为 "all"
tracker = MyCustomTracker("some_run_name")
accelerator = Accelerator(log_with=[tracker, "all"])
3. 访问内部追踪器
如果需要与跟踪器进行交互,可以使用 Accelerator.get_tracker() 方法进行访问。需要传入与跟踪器的 name 属性对应的字符串,就会在主进程中返回该跟踪器。
wandb_tracker = accelerator.get_tracker("wandb")
wandb_tracker.log_artifact(some_artifact_to_log)
# 如果需要删除 Accelerator 中使用的跟踪器,可以通过以下方式实现:
wandb_tracker = accelerator.get_tracker("wandb", unwrap=True)
if accelerator.is_main_process:
wandb_tracker.log_artifact(some_artifact_to_log)
4. 当包装器无法工作时
一些特殊跟踪器,例如 Neptune.AI,可以通过 if accelerator.is_main_process 的方式手动进行日志记录。
from accelerate import Accelerator
+ import neptune
accelerator = Accelerator()
+ run = neptune.init_run(...)
my_model, my_optimizer, my_training_dataloader = accelerate.prepare(my_model, my_optimizer, my_training_dataloader)
device = accelerator.device
my_model.to(device)
for iteration in config["num_iterations"]:
for batch in my_training_dataloader:
my_optimizer.zero_grad()
inputs, targets = batch
inputs = inputs.to(device)
targets = targets.to(device)
outputs = my_model(inputs)
loss = my_loss_function(outputs, targets)
total_loss += loss
accelerator.backward(loss)
my_optimizer.step()
+ if accelerator.is_main_process:
+ run["logs/training/batch/loss"].log(loss)
加速 - 分析器
Profiler 是在训练和推理过程中收集性能指标的工具。Profiler 的上下文管理器 API 可以用于更好地了解哪些模型运算开销最大、检查输入形状等。能够深入了解模型性能,帮助优化和改进模型。
1. 分析执行时间
使用分析器来分析执行时间:
原生 PyTorch 代码:
import torch
import torchvision.models as models
from torch.profiler import profile, record_function, ProfilerActivity
model = models.resnet18()
inputs = torch.randn(5, 3, 224, 224)
with profile(activities=[ProfilerActivity.CPU], record_shapes=True) as prof:
model(inputs)
print(prof.key_averages().table(sort_by="cpu_time_total", row_limit=10))
Accelerator 代码:
from accelerate import Accelerator, ProfileKwargs
import torch
import torchvision.models as models
model = models.resnet18()
inputs = torch.randn(5, 3, 224, 224)
profile_kwargs = ProfileKwargs(
activities=["cpu"],
record_shapes=True
)
accelerator = Accelerator(cpu=True, kwargs_handlers=[profile_kwargs])
model = accelerator.prepare(model)
with accelerator.profile() as prof:
with torch.no_grad():
model(inputs)
print(prof.key_averages().table(sort_by="cpu_time_total", row_limit=10))
输出结果:
--------------------------------- ------------ ------------ ------------ ------------
Name Self CPU CPU total CPU time avg # of Calls
--------------------------------- ------------ ------------ ------------ ------------
aten::conv2d 171.000us 52.260ms 2.613ms 20
aten::convolution 227.000us 52.089ms 2.604ms 20
aten::_convolution 270.000us 51.862ms 2.593ms 20
aten::mkldnn_convolution 51.273ms 51.592ms 2.580ms 20
aten::batch_norm 118.000us 7.059ms 352.950us 20
aten::_batch_norm_impl_index 315.000us 6.941ms 347.050us 20
aten::native_batch_norm 6.305ms 6.599ms 329.950us 20
aten::max_pool2d 40.000us 4.008ms 4.008ms 1
aten::max_pool2d_with_indices 3.968ms 3.968ms 3.968ms 1
aten::add_ 780.000us 780.000us 27.857us 28
--------------------------------- ------------ ------------ ------------ ------------
Self CPU time total: 67.016ms
# 指定 group_by_input_shape=True 用于获取更细粒度的报告
print(prof.key_averages(group_by_input_shape=True).table(sort_by="cpu_time_total", row_limit=10))
2. 分析内存消耗
显示在模型算子执行期间分配(或释放)的内存量(由模型张量使用)。
PyTorch 代码:
model = models.resnet18()
inputs = torch.randn(5, 3, 224, 224)
with profile(activities=[ProfilerActivity.CPU],
profile_memory=True, record_shapes=True) as prof:
model(inputs)
print(prof.key_averages().table(sort_by="self_cpu_memory_usage", row_limit=10))
Accelerator 代码:
model = models.resnet18()
inputs = torch.randn(5, 3, 224, 224)
profile_kwargs = ProfileKwargs(
activities=["cpu"],
profile_memory=True,
record_shapes=True
)
accelerator = Accelerator(cpu=True, kwargs_handlers=[profile_kwargs])
model = accelerator.prepare(model)
with accelerator.profile() as prof:
model(inputs)
print(prof.key_averages().table(sort_by="self_cpu_memory_usage", row_limit=10))
输出:
--------------------------------- ------------ ------------ ------------
Name CPU Mem Self CPU Mem # of Calls
--------------------------------- ------------ ------------ ------------
aten::empty 94.85 Mb 94.85 Mb 205
aten::max_pool2d_with_indices 11.48 Mb 11.48 Mb 1
aten::addmm 19.53 Kb 19.53 Kb 1
aten::mean 10.00 Kb 10.00 Kb 1
aten::empty_strided 492 b 492 b 5
aten::cat 240 b 240 b 6
aten::abs 480 b 240 b 4
aten::masked_select 120 b 112 b 1
aten::ne 61 b 53 b 3
aten::eq 30 b 30 b 1
--------------------------------- ------------ ------------ ------------
Self CPU time total: 69.332ms
3. 导出 Chrome 跟踪
在 Chrome 跟踪查看器 (chrome://tracing) 中检查分析运算符和 CUDA 内核的序列。
PyTorch 代码:
model = models.resnet18().cuda()
inputs = torch.randn(5, 3, 224, 224).cuda()
with profile(activities=[ProfilerActivity.CPU, ProfilerActivity.CUDA]) as prof:
model(inputs)
prof.export_chrome_trace("trace.json")
Accelerator 代码:
model = models.resnet18()
inputs = torch.randn(5, 3, 224, 224).cuda()
profile_kwargs = ProfileKwargs(
activities=["cpu", "cuda"],
output_trace_dir="trace"
)
accelerator = Accelerator(kwargs_handlers=[profile_kwargs])
model = accelerator.prepare(model)
with accelerator.profile() as prof:
model(inputs)
# The trace will be saved to the specified directory
4. 分析长时间运行的作业
Profiler 提供了一个额外的 API 来处理长时间运行的作业(例如训练循环)。跟踪所有执行过程可能会很慢,并且会导致跟踪文件非常大。为了避免这种情况,请使用可选参数:
schedule_option:调度选项允许您控制何时启用性能分析。这对于长时间运行的作业非常有用,可以避免收集过多数据。可用的键包括wait、warmupon_trace_ready:指定一个函数,该函数将对分析器的引用作为输入,并在每次新的跟踪准备就绪时由分析器调用
PyTorch 代码:
from torch.profiler import schedule
my_schedule = schedule(
skip_first=1,
wait=5,
warmup=1,
active=3,
repeat=2
)
def trace_handler(p):
output = p.key_averages().table(sort_by="self_cuda_time_total", row_limit=10)
print(output)
p.export_chrome_trace("/tmp/trace_" + str(p.step_num) + ".json")
with profile(
activities=[ProfilerActivity.CPU, ProfilerActivity.CUDA],
schedule=my_schedule,
on_trace_ready=trace_handler
) as p:
for idx in range(8):
model(inputs)
p.step()
Accelerator 代码:
def trace_handler(p):
output = p.key_averages().table(sort_by="self_cuda_time_total", row_limit=10)
print(output)
p.export_chrome_trace("/tmp/trace_" + str(p.step_num) + ".json")
profile_kwargs = ProfileKwargs(
activities=["cpu", "cuda"],
schedule_option={"wait": 5, "warmup": 1, "active": 3, "repeat": 2, "skip_first": 1},
on_trace_ready=trace_handler
)
accelerator = Accelerator(kwargs_handlers=[profile_kwargs])
model = accelerator.prepare(model)
with accelerator.profile() as prof:
for idx in range(8):
model(inputs)
prof.step()
5. 统计运算次数
使用公式估算特定算子(矩阵乘法和二维卷积)的 FLOP(浮点运算)次数。
PyTorch 代码:
with profile(
activities=[ProfilerActivity.CPU, ProfilerActivity.CUDA],
with_flops=True
) as prof:
model(inputs)
print(prof.key_averages().table(sort_by="flops", row_limit=10))
Accelerator 代码:
profile_kwargs = ProfileKwargs(
with_flops=True
)
accelerator = Accelerator(kwargs_handlers=[profile_kwargs])
with accelerator.profile() as prof:
model(inputs)
print(prof.key_averages().table(sort_by="flops", row_limit=10))
输出:
------------------------------------------------------- ------------ ------------ ------------
Name Self CPU Self CUDA Total FLOPs
------------------------------------------------------- ------------ ------------ ------------
aten::conv2d 197.000us 0.000us 18135613440.000
aten::addmm 103.000us 17.000us 5120000.000
aten::mul 29.000us 2.000us 30.000
aten::convolution 409.000us 0.000us --
aten::_convolution 253.000us 0.000us --
aten::cudnn_convolution 5.465ms 2.970ms --
cudaEventRecord 138.000us 0.000us --
cudaStreamIsCapturing 43.000us 0.000us --
cudaStreamGetPriority 40.000us 0.000us --
cudaDeviceGetStreamPriorityRange 10.000us 0.000us --
------------------------------------------------------- ------------ ------------ ------------
Self CPU time total: 21.938ms
Self CUDA time total: 4.165ms
加速 - 检查点
使用检查点在训练期间保存和重新加载状态的简单示例:
from accelerate import Accelerator
import torch
accelerator = Accelerator(project_dir="my/save/path")
my_scheduler = torch.optim.lr_scheduler.StepLR(my_optimizer, step_size=1, gamma=0.99)
my_model, my_optimizer, my_training_dataloader = accelerator.prepare(my_model, my_optimizer, my_training_dataloader)
# Register the LR scheduler
accelerator.register_for_checkpointing(my_scheduler)
# Save the starting state
accelerator.save_state()
device = accelerator.device
my_model.to(device)
# Perform training
for epoch in range(num_epochs):
for batch in my_training_dataloader:
my_optimizer.zero_grad()
inputs, targets = batch
inputs = inputs.to(device)
targets = targets.to(device)
outputs = my_model(inputs)
loss = my_loss_function(outputs, targets)
accelerator.backward(loss)
my_optimizer.step()
my_scheduler.step()
# Restore the previous state
accelerator.load_state("my/save/path/checkpointing/checkpoint_0")
加速 - 故障排查
1. 日志记录
请使用 logging() 而不是标准 Python 模块。使用参数 logging 设置详细级别(INFO、DEBUG、WARNING、ERROR),例如,设置 log_level="INFO":
from accelerate.logging import get_logger
logger = get_logger(__name__, log_level="DEBUG")
# 默认情况下,日志仅在主进程中调用,如果要在所有进程中调用,需要指定 main_process_only=False
# 如果需要在所有进程中按顺序调用日志,需要指定 in_order=True
from accelerate.logging import get_logger
logger = get_logger(__name__, log_level="DEBUG")
# log all processes
logger.debug("thing_to_log", main_process_only=False)
# log all processes in order
logger.debug("thing_to_log", main_process_only=False, in_order=True)
2. 代码挂起和超时错误
a. 张量形状不一致
在分布式设置中运行脚本时,需要使用诸如 Accelerator.gather()、Accelerator.reduce()、torch.distributed 等函数来跨设备获取张量,以便对其进行聚合操作。
gather 操作要求张量在所有进程中具有完全相同的形状。当张量形状不一致时,代码就会挂起,最终会引发超时异常。
【推荐】使用 Accelerate 的运行调试模式立即捕获张量形状不一致的问题,以下为多种配置方法,任选其一即可:
-
命令行配置
accelerate launch --debug {my_script.py} --arg1 --arg2 -
环境变量配置
ACCELERATE_DEBUG_MODE="1" torchrun {my_script.py} --arg1 --arg2 -
config.yaml 文件配置
compute_environment: LOCAL_MACHINE debug: true
一旦启用调试模式,可以得到张量形状不一致问题的回溯:
Traceback (most recent call last):
File "/home/zach_mueller_huggingface_co/test.py", line 18, in <module>
main()
File "/home/zach_mueller_huggingface_co/test.py", line 15, in main
broadcast_tensor = broadcast(tensor)
File "/home/zach_mueller_huggingface_co/accelerate/src/accelerate/utils/operations.py", line 303, in wrapper
accelerate.utils.operations.DistributedOperationException:
Cannot apply desired operation due to shape mismatches. All shapes across devices must be valid.
Operation: `accelerate.utils.operations.broadcast`
Input shapes:
- Process 0: [1, 5]
- Process 1: [1, 2, 5]
b. 提前停止
对于分布式训练中的提前停止,如果每个进程都有特定的停止条件(例如验证损失),则可能无法在所有进程之间同步。因此,进程 0 可能发生中断,但进程 1 不会中断,这将导致代码无限期挂起,直到超时。
如果有提前停止条件,请使用 set_trigger 和 check_trigger 方法确保所有进程都正确结束。
# Assume `should_do_breakpoint` is a custom defined function that returns a conditional,
# and that conditional might be true only on process 1
if should_do_breakpoint(loss):
accelerator.set_trigger()
# Later in the training script when we need to check for the breakpoint
if accelerator.check_trigger():
break
c. MPI
如果使用 MPI 的分布式 CPU 训练作业挂起,请确保在节点之间设置了无密码 SSH(使用密钥)。这意味着,对于主机文件中的所有节点,都能够通过 SSH 从一个节点连接到另一个节点,而无需输入密码。
接下来,尝试运行 mpirun 命令进行完整性检查。例如,以下命令应该打印出每个节点的主机名。
mpirun -f hostfile -n {number of nodes} -ppn 1 hostname
3. 内存不足
出现内存不足错误,会导致整个脚本需要重启,并且所有进度都会丢失。
为了解决这个问题,Accelerate 提供了 find_executable_batch_size() 方法,此程序会重试 OOM(内存不足)的代码,并自动降低批次大小。对于每个 OOM 情况,该程序都会将批次大小减半,然后重试代码,直到成功。
def training_function(args):
accelerator = Accelerator()
+ @find_executable_batch_size(starting_batch_size=args.batch_size)
+ def inner_training_loop(batch_size):
+ nonlocal accelerator # Ensure they can be used in our context
+ accelerator.free_memory() # Free all lingering references
model = get_model()
model.to(accelerator.device)
optimizer = get_optimizer()
train_dataloader, eval_dataloader = get_dataloaders(accelerator, batch_size)
lr_scheduler = get_scheduler(
optimizer,
num_training_steps=len(train_dataloader)*num_epochs
)
model, optimizer, train_dataloader, eval_dataloader, lr_scheduler = accelerator.prepare(
model, optimizer, train_dataloader, eval_dataloader, lr_scheduler
)
train(model, optimizer, train_dataloader, lr_scheduler)
validate(model, eval_dataloader)
+ inner_training_loop()
4. 不同设备之间结果不同
在 TPU、多 GPU 和单 GPU 上的结果会有所不同。
例如,如果之前在单 GPU 上训练,批次大小为 16,现在切换到两个 GPU 设置,则需要将批次大小更改为 8,才能获得相同的有效批次大小。这是因为在使用 Accelerate 进行训练时,传递给数据加载器的批次大小是每个 GPU 的批次大小之和。
为了确保可以在设置之间重现结果,请确保使用相同的种子,相应地调整批量大小,并考虑调整学习率。
5. 不同 GPU 上的性能问题
多 GPU 训练环境由不同类型的 GPU 组成,可能会遇到性能问题:
- GPU 之间的 GPU 内存可能不平衡。在这种情况下,内存较小的 GPU 将限制批次大小或模型大小。
- 如果使用具有不同性能的 GPU,则性能将由最慢的 GPU 决定,因为其他 GPU 必须等待它完成任务。
训练 - 梯度累积
梯度累积可以训练比机器通常能够装入内存的更大的批次大小。这是通过在多个批次上累积梯度来实现的,并且仅在执行了一定数量的批次后才更新优化器。
1. 梯度累积代码
PyTorch 代码(每两批执行一次梯度累积):
device = "cuda"
model.to(device)
gradient_accumulation_steps = 2
for index, batch in enumerate(training_dataloader):
inputs, targets = batch
inputs = inputs.to(device)
targets = targets.to(device)
outputs = model(inputs)
loss = loss_function(outputs, targets)
loss = loss / gradient_accumulation_steps
loss.backward()
if (index + 1) % gradient_accumulation_steps == 0:
optimizer.step()
scheduler.step()
optimizer.zero_grad()
调整为 Accelerator 代码(无梯度累积功能):
+ from accelerate import Accelerator
+ accelerator = Accelerator()
+ model, optimizer, training_dataloader, scheduler = accelerator.prepare(
+ model, optimizer, training_dataloader, scheduler
+ )
for index, batch in enumerate(training_dataloader):
inputs, targets = batch
- inputs = inputs.to(device)
- targets = targets.to(device)
outputs = model(inputs)
loss = loss_function(outputs, targets)
loss = loss / gradient_accumulation_steps
+ accelerator.backward(loss)
if (index+1) % gradient_accumulation_steps == 0:
optimizer.step()
scheduler.step()
optimizer.zero_grad()
将 Accelerator 代码调整为梯度累积功能:
# 1. 指定 gradient_accumulation_steps 参数
from accelerate import Accelerator
- accelerator = Accelerator()
+ accelerator = Accelerator(gradient_accumulation_steps=2)
# 2. 使用 accumulate 上下文管理器自动执行梯度累积
- for index, batch in enumerate(training_dataloader):
+ for batch in training_dataloader:
+ with accelerator.accumulate(model):
inputs, targets = batch
outputs = model(inputs)
# 3. 删除所有损失值的特殊检查
- loss = loss / gradient_accumulation_steps
accelerator.backward(loss)
- if (index+1) % gradient_accumulation_steps == 0:
optimizer.step()
scheduler.step()
optimizer.zero_grad()
完整代码:
from accelerate import Accelerator
accelerator = Accelerator(gradient_accumulation_steps=2)
model, optimizer, training_dataloader, scheduler = accelerator.prepare(
model, optimizer, training_dataloader, scheduler
)
for batch in training_dataloader:
with accelerator.accumulate(model):
inputs, targets = batch
outputs = model(inputs)
loss = loss_function(outputs, targets)
accelerator.backward(loss)
optimizer.step()
scheduler.step()
optimizer.zero_grad()
代码示例:
import torch
import copy
from accelerate import Accelerator
from accelerate.utils import set_seed
from torch.utils.data import TensorDataset, DataLoader
# seed
set_seed(0)
# define toy inputs and labels
x = torch.tensor([1., 2., 3., 4., 5., 6., 7., 8.])
y = torch.tensor([2., 4., 6., 8., 10., 12., 14., 16.])
gradient_accumulation_steps = 4
per_device_batch_size = len(x) // gradient_accumulation_steps
# define dataset and dataloader
dataset = TensorDataset(x, y)
dataloader = DataLoader(dataset, batch_size=per_device_batch_size)
# define model, optimizer and loss function
class SimpleLinearModel(torch.nn.Module):
def __init__(self):
super(SimpleLinearModel, self).__init__()
self.weight = torch.nn.Parameter(torch.zeros((1, 1)))
def forward(self, inputs):
return inputs @ self.weight
model = SimpleLinearModel()
model_clone = copy.deepcopy(model)
criterion = torch.nn.MSELoss()
model_optimizer = torch.optim.SGD(model.parameters(), lr=0.02)
accelerator = Accelerator(gradient_accumulation_steps=gradient_accumulation_steps)
model, model_optimizer, dataloader = accelerator.prepare(model, model_optimizer, dataloader)
model_clone_optimizer = torch.optim.SGD(model_clone.parameters(), lr=0.02)
print(f"initial model weight is {model.weight.mean().item():.5f}")
print(f"initial model weight is {model_clone.weight.mean().item():.5f}")
for i, (inputs, labels) in enumerate(dataloader):
with accelerator.accumulate(model):
inputs = inputs.view(-1, 1)
print(i, inputs.flatten())
labels = labels.view(-1, 1)
outputs = model(inputs)
loss = criterion(outputs, labels)
accelerator.backward(loss)
model_optimizer.step()
model_optimizer.zero_grad()
loss = criterion(x.view(-1, 1) @ model_clone.weight, y.view(-1, 1))
model_clone_optimizer.zero_grad()
loss.backward()
model_clone_optimizer.step()
print(f"w/ accumulation, the final model weight is {model.weight.mean().item():.5f}")
print(f"w/o accumulation, the final model weight is {model_clone.weight.mean().item():.5f}")
2. 对大小可变的训练样本进行梯度累积
对于像因果语言模型训练这样的跨 token 级任务的梯度累积,正确的损失应该是:梯度累积步骤中所有批次的总损失,除以这些批次中所有非填充 token 的总数。这与每批次损失值的平均值不同。
示例代码:
from accelerate import Accelerator
import math
import contextlib
gradient_accumulation_steps = 2
accelerator = Accelerator(gradient_accumulation_steps=gradient_accumulation_steps)
model, optimizer, training_dataloader, scheduler = accelerator.prepare(
model, optimizer, training_dataloader, scheduler
)
training_iterator = iter(training_dataloader)
num_samples_in_epoch = len(training_dataloader)
remainder = num_samples_in_epoch % gradient_accumulation_steps
remainder = remainder if remainder != 0 else gradient_accumulation_steps
total_updates = math.ceil(num_samples_in_epoch / gradient_accumulation_steps)
total_batched_samples = 0
for update_step in range(total_updates):
# In order to correctly the total number of non-padded tokens on which we'll compute the cross-entropy loss
# we need to pre-load the full local batch - i.e the next per_device_batch_size * accumulation_steps samples
batch_samples = []
num_batches_in_step = gradient_accumulation_steps if update_step != (total_updates - 1) else remainder
for _ in range(num_batches_in_step):
batch_samples += [next(training_iterator)]
# get local num items in batch
num_items_in_batch = sum([(batch["labels"].ne(-100)).sum() for batch in batch_samples])
# to compute it correctly in a multi-device DDP training, we need to gather the total number of items in the full batch.
num_items_in_batch = accelerator.gather(num_items_in_batch).sum().item()
for i, batch in enumerate(batch_samples):
# if we perform gradient accumulation in a multi-devices set-up, we want to avoid unecessary communications when accumulating
# cf: https://muellerzr.github.io/blog/gradient_accumulation.html
if (i < len(batch_samples) - 1 and accelerator.num_processes > 1):
ctx = model.no_sync
else:
ctx = contextlib.nullcontext
total_batched_samples += 1
with ctx():
inputs, targets = batch
outputs = model(inputs)
loss = loss_function(outputs, targets) # the loss function shoud sum over samples rather than averaging
# We multiply by num_processes because the DDP calculates the average gradient across all devices whereas dividing by num_items_in_batch already takes into account all devices
# Same reason for gradient_accumulation_steps, but this times it's Accelerate that calculate the average gradient across the accumulated steps
loss = (loss * gradient_accumulation_steps * accelerator.num_processes) / num_items_in_batch
accelerator.backward(loss)
# Sync gradients and perform optimization steps once every gradient_accumulation_steps
optimizer.step()
scheduler.step()
optimizer.zero_grad()
针对因果语言模型进行调整:
import torch
import copy
from accelerate import Accelerator
from accelerate.utils import set_seed
from accelerate.logging import get_logger
from torch.utils.data import Dataset, DataLoader
import math
import contextlib
# seed
set_seed(0)
logger = get_logger(__name__)
class MyDataset(Dataset):
def __init__(self, num_samples):
super().__init__()
self.len = num_samples
def __getitem__(self, index):
input_ids = torch.arange(1, index+2, dtype=torch.float32)
labels = torch.remainder(input_ids, 2)
return {"input_ids": input_ids, "labels": labels}
def __len__(self):
return self.len
def collate_fn(features):
input_ids = torch.nn.utils.rnn.pad_sequence([f["input_ids"] for f in features], batch_first=True, padding_value=-100)
labels = torch.nn.utils.rnn.pad_sequence([f["labels"] for f in features], batch_first=True, padding_value=-100)
return {"input_ids": input_ids[..., None], "labels": labels[..., None]}
# define toy inputs and labels
gradient_accumulation_steps = 2
per_device_batch_size = 4
# define accelerator
accelerator = Accelerator(gradient_accumulation_steps=gradient_accumulation_steps)
# define dataset and dataloader
# for this toy example, we'll compute gradient descent over one single global batch
dataset = MyDataset(per_device_batch_size*gradient_accumulation_steps*accelerator.num_processes)
dataloader = DataLoader(dataset, batch_size=per_device_batch_size, collate_fn=collate_fn)
# define model, model_optimizer and loss function
model = torch.nn.Linear(1, 2, bias=False)
model_clone = copy.deepcopy(model)
criterion = torch.nn.CrossEntropyLoss(reduction="sum") # must sum over samples rather than averaging
model_optimizer = torch.optim.SGD(model.parameters(), lr=0.08)
logger.warning(f"initial model weight is {model.weight.detach().cpu().squeeze()}")
logger.warning(f"initial model clone weight is {model_clone.weight.detach().cpu().squeeze()}")
# prepare artifacts - accelerator handles device placement and dataloader splitting
model, model_optimizer = accelerator.prepare(model, model_optimizer)
dataloader = accelerator.prepare_data_loader(dataloader, device_placement=True)
training_iterator = iter(dataloader)
num_samples_in_epoch = len(dataloader)
remainder = num_samples_in_epoch % gradient_accumulation_steps
remainder = remainder if remainder != 0 else gradient_accumulation_steps
total_gradient_updates = math.ceil(num_samples_in_epoch / gradient_accumulation_steps)
total_batched_samples = 0
for update_step in range(total_gradient_updates):
# In order to correctly the total number of non-padded tokens on which we'll compute the cross-entropy loss
# we need to pre-load the full local batch - i.e the next per_device_batch_size * accumulation_steps samples
batch_samples = []
num_batches_in_step = gradient_accumulation_steps if update_step != (total_gradient_updates - 1) else remainder
for _ in range(num_batches_in_step):
batch_samples += [next(training_iterator)]
# get local num items in batch
local_num_items_in_batch = sum([(batch["labels"].ne(-100)).sum() for batch in batch_samples])
logger.warning(f"Step {update_step} - Device {accelerator.process_index} - num items in the local batch {local_num_items_in_batch}", main_process_only=False)
# to compute it correctly in a multi-device DDP training, we need to gather the total number of items in the full batch.
num_items_in_batch = accelerator.gather(local_num_items_in_batch).sum().item()
logger.warning(f"Total num items {num_items_in_batch}")
for i, batch in enumerate(batch_samples):
inputs, labels = batch["input_ids"], batch["labels"]
total_batched_samples += 1
# if we perform gradient accumulation in a multi-devices set-up, we want to avoid unecessary communications when accumulating
# cf: https://muellerzr.github.io/blog/gradient_accumulation.html
if (i < len(batch_samples) - 1 and accelerator.num_processes > 1):
ctx = model.no_sync
else:
ctx = contextlib.nullcontext
with ctx():
outputs = model(inputs)
loss = criterion(outputs.view(-1, 2), labels.view(-1).to(torch.int64))
# We multiply by num_processes because the DDP calculates the average gradient across all devices whereas dividing by num_items_in_batch already takes into account all devices
# Same reason for gradient_accumulation_steps, but this times it's Accelerate that calculate the average gradient across the accumulated steps
loss = (loss * gradient_accumulation_steps * accelerator.num_processes) / num_items_in_batch
accelerator.backward(loss)
model_optimizer.step()
model_optimizer.zero_grad()
logger.warning(f"Device {accelerator.process_index} - w/ accumulation, the final model weight is {accelerator.unwrap_model(model).weight.detach().cpu().squeeze()}", main_process_only=False)
# We know do the same operation but on a single device and without gradient accumulation
if accelerator.is_main_process:
# prepare one single entire batch
dataloader = DataLoader(dataset, batch_size=len(dataset), collate_fn=collate_fn)
full_batch_without_accum = next(iter(dataloader))
total_inputs, total_labels = full_batch_without_accum["input_ids"], full_batch_without_accum["labels"]
model_clone_optimizer = torch.optim.SGD(model_clone.parameters(), lr=0.08)
# train the cloned model
loss = torch.nn.CrossEntropyLoss(reduction="mean")(model_clone(total_inputs).view(-1, 2), total_labels.view(-1).to(torch.int64))
model_clone_optimizer.zero_grad()
loss.backward()
model_clone_optimizer.step()
# We should have the same final weights.
logger.warning(f"w/o accumulation, the final model weight is {model_clone.weight.detach().cpu().squeeze()}")
单台设备上的结果(梯度累积步数设置为 1,batch_size 设置为 8):
initial model weight is tensor([-0.0075, 0.5364])
initial model clone weight is tensor([-0.0075, 0.5364])
Step 0 - Device 0 - num items in the local batch 36
Total num items 36
Device 0 - w/ accumulation, the final model weight is tensor([0.0953, 0.4337])
w/o accumulation, the final model weight is tensor([0.0953, 0.4337])
两台设备上的结果(梯度累积步数设置为 2,batch_size 设置为 4):
initial model weight is tensor([-0.0075, 0.5364])
initial model clone weight is tensor([-0.0075, 0.5364])
Step 0 - Device 0 - num items in the local batch 52
Step 0 - Device 1 - num items in the local batch 84
Total num items 136
Device 1 - w/ accumulation, the final model weight is tensor([0.2117, 0.3172])
Device 0 - w/ accumulation, the final model weight is tensor([0.2117, 0.3172])
w/o accumulation, the final model weight is tensor([0.2117, 0.3172])
训练 - 局部 SGD
局部随机梯度下降 (SGD) 是一种分布式训练技术,其梯度并非每一步都同步。因此,每个进程都会更新自身版本的模型权重,并在给定步数后通过在所有进程之间取平均值来同步这些权重。提高了通信效率,并可以显著加快训练速度,尤其是在缺乏 NVLink 时。与梯度累积(提高通信效率需要增加有效批次大小)不同,局部随机梯度下降 (SGD) 不需要更改批次大小或学习率/学习进度。
样例 PyTorch 代码(每两批执行一次梯度累积):
device = "cuda"
model.to(device)
gradient_accumulation_steps = 2
for index, batch in enumerate(training_dataloader):
inputs, targets = batch
inputs = inputs.to(device)
targets = targets.to(device)
outputs = model(inputs)
loss = loss_function(outputs, targets)
loss = loss / gradient_accumulation_steps
loss.backward()
if (index + 1) % gradient_accumulation_steps == 0:
optimizer.step()
scheduler.step()
optimizer.zero_grad()
转换为 Accelerator 代码(不含局部 SGD 或梯度累积功能):
+ from accelerate import Accelerator
+ accelerator = Accelerator()
+ model, optimizer, training_dataloader, scheduler = accelerator.prepare(
+ model, optimizer, training_dataloader, scheduler
+ )
for index, batch in enumerate(training_dataloader):
inputs, targets = batch
- inputs = inputs.to(device)
- targets = targets.to(device)
outputs = model(inputs)
loss = loss_function(outputs, targets)
loss = loss / gradient_accumulation_steps
+ accelerator.backward(loss)
if (index+1) % gradient_accumulation_steps == 0:
optimizer.step()
scheduler.step()
调整 Accelerator 代码,使其带有局部 SGD 功能:
+ local_sgd_steps = 8
+ with LocalSGD(accelerator=accelerator, model=model, local_sgd_steps=8, enabled=True) as local_sgd:
for batch in training_dataloader:
with accelerator.accumulate(model):
inputs, targets = batch
outputs = model(inputs)
loss = loss_function(outputs, targets)
accelerator.backward(loss)
optimizer.step()
scheduler.step()
optimizer.zero_grad()
+ local_sgd.step()
说明:局部 SGD 代码在底层禁用了自动梯度同步。相反,它会在每 local_sgd_steps 一步对模型参数进行平均。
总结:局部随机梯度下降 (SGD) 收敛速度快,通信量小。
训练 - 低精度(FP8)训练
Accelerate 通过 TransformersEngine、MS-AMP 和 torchao 软件包提供使用指定支持硬件进行低精度方法训练的集成。
1. FP8 训练
FP8 训练是指模型训练的部分工作可以使用 8 位而不是 16 位来完成,并且在不明显降低最终性能的情况下做到这一点。
该功能仅在特定的 NVIDIA 硬件上启用:3000 系列之后的消费级显卡(4090)、基于 Hopper 的 GPU 架构(H100 和 H200)。
这将使得内存使用量减少,吞吐量增加。
2. 配置 Accelerator
目前支持三种不同的 FP8 后端(TransformersEngine、torchao 和 MS-AMP),每种后端都有不同的功能和配置。
# Accelerator 中的配置
from accelerate import Accelerator
accelerator = Accelerator(mixed_precision="fp8")
# 如果 MS-AMP 可用,Accelerator 将自动将其用作后端,示例代码如下:
from accelerate import Accelerator
from accelerate.utils import MSAMPRecipeKwargs
kwargs = [MSAMPRecipeKwargs()]
# Or to specify the backend as `TransformersEngine` even if MS-AMP is installed
# kwargs = [TERecipeKwargs()]
# Or to use torchao
# kwargs = [AORecipeKwargs()]
accelerator = Accelerator(mixed_precision="fp8", kwarg_handlers=kwargs)
对应的 config.yaml 配置如下:
mixed_precision: fp8
fp8_config:
amax_compute_algo: max
amax_history_len: 1024
backend: TE
fp8_format: HYBRID
interval: 1
margin: 0
override_linear_precision: (false, false, false)
use_autocast_during_eval: false
3. 配置 MS-AMP
MS-AMP 易于配置,因为只有一个参数:优化级别。当前 Accelerator 集成支持两种级别的优化:O1 和 O2(使用字母 “o”,而不是零)。
- “O1”:将权重梯度和 all_reduce 通信转换为 8 位格式,其余部分则以 16 位格式进行。这减少了 GPU 内存的总体使用量,并加快了通信带宽。
- “O2”:还会将一阶优化器状态转换为 8 位,而二阶状态则为 FP16。(目前仅 Adam 支持优化器)。这会尽力减少最终准确率的下降,并最大程度地节省内存。
# 如果指定优化级别,示例代码如下:
from accelerate import Accelerator
from accelerate.utils import FP8RecipeKwargs
kwargs = [FP8RecipeKwargs(backend="msamp", optimization_level="O2")]
accelerator = Accelerator(mixed_precision="fp8", kwarg_handlers=kwargs)
也可以在命令行中使用 accelerate launch 启动训练脚本时,指定参数 --fp8_backend=msamp --fp8_opt_level=O2。
也可以配置 config.yaml:
mixed_precision: fp8
fp8_config:
backend: MSAMP
opt_level: O2
4. 配置 TransformersEngine
TransformersEngine 提供多种选项,用于自定义 FP8 计算的执行方式和内容。
# 指定 backend="te",让 Accelerator 尝试设置合理的默认值
from accelerate import Accelerator
from accelerate.utils import FP8RecipeKwargs
kwargs = [FP8RecipeKwargs(backend="te", ...)]
accelerator = Accelerator(mixed_precision="fp8", kwarg_handlers=kwargs)
命令行中使用 accelerate launch 启动,指定参数 --fp8_backend=te,也可以使用 accelerate launch --fp8_backend=te -h 查看帮助文档。
配置 config.yaml:
mixed_precision: fp8
fp8_config:
amax_compute_algo: max
amax_history_len: 1024
backend: TE
fp8_format: HYBRID
interval: 1
margin: 0
override_linear_precision: (false, false, false)
use_autocast_during_eval: false
5. 配置 torchao
torchao 的主要特点是:为了保持数值的稳定性,将模型的第一层和最后一层保持在常规精度(无论是 FP32 还是 BF16),然后将其他层量化到 FP8。
# 指定 mixed_precision="fp8"
from accelerate import Accelerator
from accelerate.utils import AORecipeKwargs
kwargs = [AORecipeKwargs()]
accelerator = Accelerator(mixed_precision="fp8", kwarg_handlers=kwargs)
训练 - DeepSpeed
1. DeepSpeed
DeepSpeed 实现了 ZeRO 论文中描述的所有内容,一些突出的优化包括:
- 优化器状态分区(ZeRO 第 1 阶段)
- 梯度分区(ZeRO 第 2 阶段):主要用于训练
- 参数分区(ZeRO 第 3 阶段):主要用于推理
- 自定义混合精度训练处理
- 一系列基于 CUDA 扩展的快速优化器
- ZeRO-卸载到 CPU 和磁盘/NVMe
- 模型参数的分层划分(ZeRO++)
Accelerate 通过 2 个选项集成 DeepSpeed:
deepspeed_config_file:只需提供自定义配置文件或使用默认配置即可。deepspeed_plugin:支持 DeepSpeed 的部分功能,其余配置使用默认选项,用户无需更改任何代码。
2. 集成
Accelerate 集成了 DeepSpeed ZeRO 的所有功能,包括 ZeRO 的第 1、2 和 3 阶段,以及 ZeRO-Offload、ZeRO-Infinity(可卸载到磁盘/NVMe)和 ZeRO++。以下是使用 ZeRO - 零冗余优化器实现数据并行的简要说明:
- 第 1 阶段:跨数据并行工作器/GPU 的分片优化器状态
- 第 2 阶段:跨数据并行工作器/GPU 的分片优化器状态 + 梯度
- 第 3 阶段:跨数据并行工作器/GPU 的分片优化器状态 + 梯度 + 模型参数
- 优化器卸载:将梯度 + 优化器状态卸载到 ZeRO Stage 2 之上的 CPU/磁盘
- Param Offload:将模型参数卸载到 ZeRO Stage 3 之上的 CPU/磁盘
- 分层分区:基于 ZeRO Stage 3 构建,通过跨节点的数据并行训练和节点内的 ZeRO-3 分片,实现高效的多节点训练。
注意:对于磁盘卸载,磁盘应该是 NVME,以获得不错的速度,但从技术上讲,它可以在任何磁盘上运行。
推理:DeepSpeed ZeRO Inference 通过 ZeRO-Infinity 支持 ZeRO 第 3 阶段。它使用与训练相同的 ZeRO 协议,但不使用优化器和 LRR 调度器,并且仅与第 3 阶段相关。更多详情请参阅:deepspeed-zero-inference。
3. DeepSpeed 版本
研究安装 DeepSpeed 版本 >= 0.6.5。
4. DeepSpeed 插件
# 使用命令生成配置文件
accelerate config
注意:当询问是否需要使用 DeepSpeed 配置文件时,回答”否”,这将生成一个默认的 DeepSpeed 配置,在执行以下命令时,默认的 DeepSpeed 配置文件会被自动加载:
accelerate launch my_script.py --args_to_my_script
使用 ZeRO Stage-2 策略的配置文件示例:
compute_environment: LOCAL_MACHINE
deepspeed_config:
gradient_accumulation_steps: 1
gradient_clipping: 1.0
offload_optimizer_device: none
offload_param_device: none
zero3_init_flag: true
zero_stage: 2
distributed_type: DEEPSPEED
fsdp_config: {}
machine_rank: 0
main_process_ip: null
main_process_port: null
main_training_function: main
mixed_precision: fp16
num_machines: 1
num_processes: 2
use_cpu: false
# 命令行启动
accelerate launch examples/nlp_example.py --mixed_precision fp16
使用 ZeRO Stage-3 + CPU Offload 策略的配置文件示例:
compute_environment: LOCAL_MACHINE
deepspeed_config:
gradient_accumulation_steps: 1
gradient_clipping: 1.0
offload_optimizer_device: cpu
offload_param_device: cpu
zero3_init_flag: true
zero3_save_16bit_model: true
zero_stage: 3
distributed_type: DEEPSPEED
fsdp_config: {}
machine_rank: 0
main_process_ip: null
main_process_port: null
main_training_function: main
mixed_precision: fp16
num_machines: 1
num_processes: 2
use_cpu: false
# 命令行启动
accelerate launch examples/nlp_example.py --mixed_precision fp16
Accelerate config 命令支持的配置选项:
| 选项 | 说明 |
|---|---|
zero_stage | [0] 禁用, [1] 优化器状态分区, [2] 优化器+梯度状态分区, [3] 优化器+梯度+参数分区 |
gradient_accumulation_steps | 在平均和应用梯度之前累积梯度的训练步骤数 |
gradient_clipping | 启用梯度裁剪并指定值 |
offload_optimizer_device | [none] 禁用优化器卸载, [cpu] 卸载到 CPU, [nvme] 卸载到 NVMe SSD。仅适用于 ZeRO >= Stage-2 |
offload_optimizer_nvme_path | 指定卸载优化器状态的 Nvme 路径。如果未指定,默认为 ‘none’ |
offload_param_device | [none] 禁用参数卸载, [cpu] 卸载到 CPU, [nvme] 卸载到 NVMe SSD。仅适用于 ZeRO Stage-3 |
offload_param_nvme_path | 指定卸载参数的 Nvme 路径。如果未指定,默认为 ‘none’ |
zero3_init_flag | 决定是否启用 deepspeed.zero.Init 来构建大型模型。仅适用于 ZeRO Stage-3 |
zero3_save_16bit_model | 决定使用 ZeRO Stage-3 时是否保存 16 位模型权重 |
mixed_precision | no 表示 FP32 训练, fp16 表示 FP16 混合精度训练, bf16 表示 BF16 混合精度训练 |
deepspeed_moe_layer_cls_names | 逗号分隔的 Transformer Mixture-of-Experts (MoE) 层类名列表(区分大小写),例如 MixtralSparseMoeBlock, Qwen2MoeSparseMoeBlock 等 |
deepspeed_hostfile | 用于配置多节点计算资源的 DeepSpeed hostfile |
deepspeed_exclusion_filter | 使用多节点设置时的 DeepSpeed 排除过滤器字符串 |
deepspeed_inclusion_filter | 使用多节点设置时的 DeepSpeed 包含过滤器字符串 |
deepspeed_multinode_launcher | 要使用的 DeepSpeed 多节点启动器,例如 pdsh, standard, openmpi 等。如果未指定,默认为 pdsh |
deepspeed_config_file | DeepSpeed 配置文件的路径(JSON 格式) |
5. DeepSpeed 配置
如果需要更多配置选项,需要使用 DeepSpeed 配置文件。
# 命令行配置界面
accelerate config
当询问是否使用 DeepSpeed 配置文件,回答”是”,并提供 DeepSpeed 配置文件路径,这将生成配置文件,该文件在执行以下命令时自动加载:
accelerate launch my_script.py --args_to_my_script
使用 ZeRO Stage-2 策略的 config.yaml 配置文件示例:
compute_environment: LOCAL_MACHINE
deepspeed_config:
deepspeed_config_file: /home/ubuntu/accelerate/examples/deepspeed_config_templates/zero_stage2_config.json
zero3_init_flag: true
distributed_type: DEEPSPEED
fsdp_config: {}
machine_rank: 0
main_process_ip: null
main_process_port: null
main_training_function: main
mixed_precision: fp16
num_machines: 1
num_processes: 2
use_cpu: false
DeepSpeed 配置文件 zero_stage2_config.json 的内容:
{
"fp16": {
"enabled": true,
"loss_scale": 0,
"loss_scale_window": 1000,
"initial_scale_power": 16,
"hysteresis": 2,
"min_loss_scale": 1
},
"optimizer": {
"type": "AdamW",
"params": {
"lr": "auto",
"weight_decay": "auto",
"torch_adam": true,
"adam_w_mode": true
}
},
"scheduler": {
"type": "WarmupDecayLR",
"params": {
"warmup_min_lr": "auto",
"warmup_max_lr": "auto",
"warmup_num_steps": "auto",
"total_num_steps": "auto"
}
},
"zero_optimization": {
"stage": 2,
"allgather_partitions": true,
"allgather_bucket_size": 2e8,
"overlap_comm": true,
"reduce_scatter": true,
"reduce_bucket_size": "auto",
"contiguous_gradients": true
},
"gradient_accumulation_steps": 1,
"gradient_clipping": "auto",
"steps_per_print": 2000,
"train_batch_size": "auto",
"train_micro_batch_size_per_gpu": "auto",
"wall_clock_breakdown": false
}
# 命令行启动
accelerate launch examples/by_feature/deepspeed_with_config_support.py \
--config_name "gpt2-large" \
--tokenizer_name "gpt2-large" \
--dataset_name "wikitext" \
--dataset_config_name "wikitext-2-raw-v1" \
--block_size 128 \
--output_dir "./clm/clm_deepspeed_stage2_accelerate" \
--learning_rate 5e-4 \
--per_device_train_batch_size 24 \
--per_device_eval_batch_size 24 \
--num_train_epochs 3 \
--with_tracking \
--report_to "wandb"
使用 ZeRO Stage-3 + CPU 卸载功能的配置文件示例:
compute_environment: LOCAL_MACHINE
deepspeed_config:
deepspeed_config_file: /home/ubuntu/accelerate/examples/deepspeed_config_templates/zero_stage3_offload_config.json
zero3_init_flag: true
distributed_type: DEEPSPEED
fsdp_config: {}
machine_rank: 0
main_process_ip: null
main_process_port: null
main_training_function: main
mixed_precision: fp16
num_machines: 1
num_processes: 2
use_cpu: false
zero_stage3_offload_config.json 的内容:
{
"fp16": {
"enabled": true,
"loss_scale": 0,
"loss_scale_window": 1000,
"initial_scale_power": 16,
"hysteresis": 2,
"min_loss_scale": 1
},
"optimizer": {
"type": "AdamW",
"params": {
"lr": "auto",
"weight_decay": "auto"
}
},
"scheduler": {
"type": "WarmupDecayLR",
"params": {
"warmup_min_lr": "auto",
"warmup_max_lr": "auto",
"warmup_num_steps": "auto",
"total_num_steps": "auto"
}
},
"zero_optimization": {
"stage": 3,
"offload_optimizer": {
"device": "cpu",
"pin_memory": true
},
"offload_param": {
"device": "cpu",
"pin_memory": true
},
"overlap_comm": true,
"contiguous_gradients": true,
"reduce_bucket_size": "auto",
"stage3_prefetch_bucket_size": "auto",
"stage3_param_persistence_threshold": "auto",
"sub_group_size": 1e9,
"stage3_max_live_parameters": 1e9,
"stage3_max_reuse_distance": 1e9,
"stage3_gather_16bit_weights_on_model_save": "auto"
},
"gradient_accumulation_steps": 1,
"gradient_clipping": "auto",
"steps_per_print": 2000,
"train_batch_size": "auto",
"train_micro_batch_size_per_gpu": "auto",
"wall_clock_breakdown": false
}
# 命令行启动
accelerate launch examples/by_feature/deepspeed_with_config_support.py \
--config_name "gpt2-large" \
--tokenizer_name "gpt2-large" \
--dataset_name "wikitext" \
--dataset_config_name "wikitext-2-raw-v1" \
--block_size 128 \
--output_dir "./clm/clm_deepspeed_stage3_offload_accelerate" \
--learning_rate 5e-4 \
--per_device_train_batch_size 32 \
--per_device_eval_batch_size 32 \
--num_train_epochs 3 \
--with_tracking \
--report_to "wandb"
6. 保存和加载
ZeRO Stage-1 和 Stage-2 的模型保存和加载方式不变。
对于 ZeRO Stage-3,有以下两种方式:
a. 保存整个 16 位模型权重
以便后续直接加载 model.load_state_dict(torch.load(pytorch_model.bin)):
- 将
zero_optimization.stage3_gather_16bit_weights_on_model_save在 DeepSpeed 配置文件中设置为True - 将
zero3_save_16bit_model在 DeepSpeed 插件中设置为True
unwrapped_model = accelerator.unwrap_model(model)
# New Code #
# Saves the whole/unpartitioned fp16 model when in ZeRO Stage-3 to the output directory if
# `stage3_gather_16bit_weights_on_model_save` is True in DeepSpeed Config file or
# `zero3_save_16bit_model` is True in DeepSpeed Plugin.
# For Zero Stages 1 and 2, models are saved as usual in the output directory.
# The model name saved is `pytorch_model.bin`
unwrapped_model.save_pretrained(
args.output_dir,
is_main_process=accelerator.is_main_process,
save_function=accelerator.save,
state_dict=accelerator.get_state_dict(model),
)
b. 获取 32 位权重
保存模型 model.save_checkpoint():
success = model.save_checkpoint(PATH, ckpt_id, checkpoint_state_dict)
status_msg = f"checkpointing: PATH={PATH}, ckpt_id={ckpt_id}"
if success:
logging.info(f"Success {status_msg}")
else:
logging.warning(f"Failure {status_msg}")
这将在检查点目录中创建 ZeRO 模型和优化器区分以及 zero_to_fp32.py 脚本,使用示例:
$ cd /path/to/checkpoint_dir
$ ./zero_to_fp32.py . pytorch_model.bin
Processing zero checkpoint at global_step1
Detected checkpoint of type zero stage 3, world_size: 2
Saving fp32 state dict to pytorch_model.bin (total_numel=60506624)
加载模型:
from deepspeed.utils.zero_to_fp32 import load_state_dict_from_zero_checkpoint
unwrapped_model = accelerator.unwrap_model(model)
fp32_model = load_state_dict_from_zero_checkpoint(unwrapped_model, checkpoint_dir)
加载模型状态:
from deepspeed.utils.zero_to_fp32 import get_fp32_state_dict_from_zero_checkpoint
state_dict = get_fp32_state_dict_from_zero_checkpoint(checkpoint_dir)
注意:所有这些功能都需要 checkpoint 大小约 2 倍的内存(通用 RAM)。
7. ZeRO 推理
通过 Accelerator,获取模型和数据加载器:
model, eval_dataloader = accelerator.prepare(model, eval_dataloader)
8. 当前不支持的功能
- 当前集成不支持 DeepSpeed 的管道并行性。
- 当前集成不支持 mpu,限制了 Megatron-LM 中支持的张量并行性。
- 当前集成不支持多种模型。
训练 - DDP 通信机制
分布式数据并行(DDP)通信钩子提供了一个通用接口。它提供了一些内置通信钩子,用于优化通信:
- FP16 压缩钩:通过将梯度转换为半精度浮点格式(
torch.float16)来压缩梯度,从而减少通信开销。 - BF16 压缩钩:与 FP16 类似,但使用 Brain 浮点格式(
torch.bfloat16),在某些硬件上效率更高。 - PowerSGD 钩子:一种先进的梯度压缩算法,可提供高压缩率并可加速受带宽限制的分布式训练。
1. FP16 压缩钩
PyTorch 代码:
import torch
from torch.nn.parallel import DistributedDataParallel as DDP
from torch.distributed.algorithms.ddp_comm_hooks import default_hooks
from accelerate.test_utils.testing import get_backend
device_type, _, _ = get_backend()
device_id = getattr(torch, device_type, torch.cuda).current_device()
class MyModel(torch.nn.Module):
def __init__(self):
super().__init__()
self.layer = torch.nn.Linear(10, 10)
def forward(self, x):
return self.layer(x)
model = MyModel()
model = DDP(model, device_ids=[device_id])
model.register_comm_hook(state=None, hook=default_hooks.fp16_compress_hook)
# Training loop
for data, targets in data_loader:
outputs = model(data)
loss = criterion(outputs, targets)
loss.backward()
optimizer.step()
optimizer.zero_grad()
Accelerator 代码:
from accelerate import Accelerator, DDPCommunicationHookType, DistributedDataParallelKwargs
import torch
class MyModel(torch.nn.Module):
def __init__(self):
super().__init__()
self.layer = torch.nn.Linear(10, 10)
def forward(self, x):
return self.layer(x)
# DDP Communication Hook setup
ddp_kwargs = DistributedDataParallelKwargs(comm_hook=DDPCommunicationHookType.FP16)
accelerator = Accelerator(kwargs_handlers=[ddp_kwargs])
model = MyModel()
optimizer = torch.optim.Adam(model.parameters())
data_loader = DataLoader(dataset, batch_size=16)
model, optimizer, data_loader = accelerator.prepare(model, optimizer, data_loader)
# Training loop
for data, targets in data_loader:
outputs = model(data)
loss = criterion(outputs, targets)
accelerator.backward(loss)
optimizer.step()
optimizer.zero_grad()
2. BF16 压缩钩
PyTorch 代码:
import torch
from torch.nn.parallel import DistributedDataParallel as DDP
from torch.distributed.algorithms.ddp_comm_hooks import default_hooks
from accelerate.test_utils.testing import get_backend
device_type, _, _ = get_backend()
device_id = getattr(torch, device_type, torch.cuda).current_device()
class MyModel(torch.nn.Module):
def __init__(self):
super().__init__()
self.layer = torch.nn.Linear(10, 10)
def forward(self, x):
return self.layer(x)
model = MyModel()
model = DDP(model, device_ids=[device_id])
model.register_comm_hook(state=None, hook=default_hooks.bf16_compress_hook)
# Training loop
for data, targets in data_loader:
outputs = model(data)
loss = criterion(outputs, targets)
loss.backward()
optimizer.step()
optimizer.zero_grad()
Accelerator 代码:
from accelerate import Accelerator, DDPCommunicationHookType, DistributedDataParallelKwargs
import torch
class MyModel(torch.nn.Module):
def __init__(self):
super().__init__()
self.layer = torch.nn.Linear(10, 10)
def forward(self, x):
return self.layer(x)
# DDP Communication Hook setup
ddp_kwargs = DistributedDataParallelKwargs(comm_hook=DDPCommunicationHookType.BF16)
accelerator = Accelerator(kwargs_handlers=[ddp_kwargs])
model = MyModel()
optimizer = torch.optim.Adam(model.parameters())
data_loader = DataLoader(dataset, batch_size=16)
model, optimizer, data_loader = accelerator.prepare(model, optimizer, data_loader)
# Training loop
for data, targets in data_loader:
outputs = model(data)
loss = criterion(outputs, targets)
accelerator.backward(loss)
optimizer.step()
optimizer.zero_grad()
3. PowerSGD 钩子
PyTorch 代码:
import torch
from torch.nn.parallel import DistributedDataParallel as DDP
from torch.distributed.algorithms.ddp_comm_hooks import powerSGD_hook
from accelerate.test_utils.testing import get_backend
device_type, _, _ = get_backend()
device_id = getattr(torch, device_type, torch.cuda).current_device()
class MyModel(torch.nn.Module):
def __init__(self):
super().__init__()
self.layer = torch.nn.Linear(10, 10)
def forward(self, x):
return self.layer(x)
model = MyModel()
model = DDP(model, device_ids=[device_id])
state = powerSGD_hook.PowerSGDState(process_group=None)
model.register_comm_hook(state=state, hook=powerSGD_hook.powerSGD_hook)
# Training loop
for data, targets in data_loader:
outputs = model(data)
loss = criterion(outputs, targets)
loss.backward()
optimizer.step()
optimizer.zero_grad()
Accelerator 代码:
from accelerate import Accelerator, DDPCommunicationHookType, DistributedDataParallelKwargs
import torch
class MyModel(torch.nn.Module):
def __init__(self):
super().__init__()
self.layer = torch.nn.Linear(10, 10)
def forward(self, x):
return self.layer(x)
# DDP Communication Hook setup
ddp_kwargs = DistributedDataParallelKwargs(comm_hook=DDPCommunicationHookType.POWER_SGD)
accelerator = Accelerator(kwargs_handlers=[ddp_kwargs])
model = MyModel()
optimizer = torch.optim.Adam(model.parameters())
data_loader = DataLoader(dataset, batch_size=16)
model, optimizer, data_loader = accelerator.prepare(model, optimizer, data_loader)
# Training loop
for data, targets in data_loader:
outputs = model(data)
loss = criterion(outputs, targets)
accelerator.backward(loss)
optimizer.step()
optimizer.zero_grad()
4. DDP 通信钩子高级用法
comm_wrapper 包装器:将通信钩子包装成附加功能的选项。例如,它可以用于将 FP16 压缩与其他通信策略结合使用。目前支持的包装器包括 NO、FP16 和 BF16。
from accelerate import Accelerator, DDPCommunicationHookType, DistributedDataParallelKwargs
import torch
class MyModel(torch.nn.Module):
def __init__(self):
super().__init__()
self.layer = torch.nn.Linear(10, 10)
def forward(self, x):
return self.layer(x)
# DDP Communication Hook setup
ddp_kwargs = DistributedDataParallelKwargs(
comm_hook=DDPCommunicationHookType.POWER_SGD,
comm_wrapper=DDPCommunicationHookType.FP16
)
accelerator = Accelerator(kwargs_handlers=[ddp_kwargs])
model = MyModel()
optimizer = torch.optim.Adam(model.parameters())
data_loader = DataLoader(dataset, batch_size=16)
model, optimizer, data_loader = accelerator.prepare(model, optimizer, data_loader)
# Training loop
for data, targets in data_loader:
outputs = model(data)
loss = criterion(outputs, targets)
accelerator.backward(loss)
optimizer.step()
optimizer.zero_grad()
comm_state_option 状态选项:允许传递某些通信钩子所需的额外状态信息。这对于像 PowerSGD 这样的状态钩子尤其有用,因为它们需要在训练步骤中维护超参数和内部状态。
from accelerate import Accelerator, DDPCommunicationHookType, DistributedDataParallelKwargs
import torch
class MyModel(torch.nn.Module):
def __init__(self):
super().__init__()
self.layer = torch.nn.Linear(10, 10)
def forward(self, x):
return self.layer(x)
# DDP Communication Hook setup
ddp_kwargs = DistributedDataParallelKwargs(
comm_hook=DDPCommunicationHookType.POWER_SGD,
comm_state_option={"matrix_approximation_rank": 2}
)
accelerator = Accelerator(kwargs_handlers=[ddp_kwargs])
model = MyModel()
optimizer = torch.optim.Adam(model.parameters())
data_loader = DataLoader(dataset, batch_size=16)
model, optimizer, data_loader = accelerator.prepare(model, optimizer, data_loader)
# Training loop
for data, targets in data_loader:
outputs = model(data)
loss = criterion(outputs, targets)
accelerator.backward(loss)
optimizer.step()
optimizer.zero_grad()
训练 - 完全分片数据并行(FSDP)
为了加速在更大批量上训练大型模型,可以使用完全分片的数据并行模型。这种数据并行范式通过对优化器状态、梯度和参数进行分片,可以拟合更多数据和更大的模型。
1. 使用方式
# 从命令行开始参数配置
accelerate config
回答问题并生成配置文件后,命令行启动任务:
accelerate launch my_script.py --args_to_my_script
启动 FSDP 的示例配置文件:
compute_environment: LOCAL_MACHINE
debug: false
distributed_type: FSDP
downcast_bf16: 'no'
fsdp_config:
fsdp_auto_wrap_policy: TRANSFORMER_BASED_WRAP
fsdp_backward_prefetch_policy: BACKWARD_PRE
fsdp_forward_prefetch: false
fsdp_cpu_ram_efficient_loading: true
fsdp_offload_params: false
fsdp_sharding_strategy: FULL_SHARD
fsdp_state_dict_type: SHARDED_STATE_DICT
fsdp_sync_module_states: true
fsdp_transformer_layer_cls_to_wrap: BertLayer
fsdp_use_orig_params: true
machine_rank: 0
main_training_function: main
mixed_precision: bf16
num_machines: 1
num_processes: 2
rdzv_backend: static
same_network: true
tpu_env: []
tpu_use_cluster: false
tpu_use_sudo: false
use_cpu: false
# 命令行执行
accelerate launch examples/nlp_example.py
accelerate config 当前支持的配置参数:
| 参数 | 说明 |
|---|---|
fsdp_sharding_strategy | [1] FULL_SHARD(分片优化器状态、梯度和参数),[2] SHARD_GRAD_OP(分片优化器状态和梯度),[3] NO_SHARD (DDP),[4] HYBRID_SHARD(分片优化器状态、梯度和参数,每个节点都有完整副本),[5] HYBRID_SHARD_ZERO2(分片优化器状态和梯度,每个节点都有完整副本) |
fsdp_offload_params | 决定是否将参数和梯度卸载到 CPU |
fsdp_auto_wrap_policy | [1] 基于 TRANSFORMER_BASED_WRAP,[2] 基于 SIZE_BASED_WRAP,[3] 无 WRAP |
fsdp_transformer_layer_cls_to_wrap | 仅适用于 Transformer。使用 fsdp_auto_wrap_policy=TRANSFORMER_BASED_WRAP 时,提供要包装的 Transformer 层类名(区分大小写)的逗号分隔字符串,例如 BertLayer、GPTJBlock、T5Block 等。共享权重(如嵌入层)的子模块不应该位于不同的 FSDP 包装单元中。使用此策略,每个包含多头注意力后跟 MLP 层的块都会进行包装,其余层包装在最外层 FSDP 单元中。也可以使用 model._no_split_modules 来自动包装 |
fsdp_min_num_params | 使用 fsdp_auto_wrap_policy=SIZE_BASED_WRAP 时的最小参数数量 |
fsdp_backward_prefetch_policy | [1] BACKWARD_PRE(后向预取),[2] BACKWARD_POST,[3] NO_PREFETCH |
fsdp_forward_prefetch | 如果为 True,FSDP 会在前向传递执行时预取下一个全收集操作。仅适用于静态图模型 |
fsdp_state_dict_type | [1] FULL_STATE_DICT,[2] LOCAL_STATE_DICT,[3] SHARDED_STATE_DICT |
fsdp_use_orig_params | 如果为 True,允许非均匀 requires_grad 分布(支持冻结/可训练参数),适合参数高效微调。也允许多优化器参数组。应在创建优化器之前设置为 True |
fsdp_cpu_ram_efficient_loading | 仅适用于 Transformers 模型。True 时只有第一个进程加载预训练模型检查点,其他进程权重为空。需同时设置 fsdp_sync_module_states=true,并在调用 from_pretrained 前初始化分布式进程组 |
fsdp_sync_module_states | 如果为 True,每个 FSDP 单元从 rank 0 广播模块参数 |
通过 FullyShardedDataParallelPlugin 进行更精细控制:
from accelerate import FullyShardedDataParallelPlugin
from torch.distributed.fsdp.fully_sharded_data_parallel import FullOptimStateDictConfig, FullStateDictConfig
fsdp_plugin = FullyShardedDataParallelPlugin(
state_dict_config=FullStateDictConfig(offload_to_cpu=False, rank0_only=False),
optim_state_dict_config=FullOptimStateDictConfig(offload_to_cpu=False, rank0_only=False),
)
accelerator = Accelerator(fsdp_plugin=fsdp_plugin)
2. 保存和加载
保存模型中间结果:
accelerator.save_state("ckpt")
查看保存的内容:
ls ckpt
# 输出:optimizer_0 pytorch_model_0 random_states_0.pkl random_states_1.pkl scheduler.bin
cd ckpt
ls optimizer_0
# 输出:__0_0.distcp __1_0.distcp
ls pytorch_model_0
# 输出:__0_0.distcp __1_0.distcp
加载并恢复训练:
accelerator.load_state("ckpt")
保存状态词典:
unwrapped_model.save_pretrained(
args.output_dir,
is_main_process=accelerator.is_main_process,
save_function=accelerator.save,
state_dict=accelerator.get_state_dict(model), # 新增
)
3. FSDP 分片与 ZeRO 阶段的对比
| FSDP 分片策略 | DeepSpeed ZeRO 阶段 | 说明 |
|---|---|---|
FULL_SHARD | ZeRO Stage-3 | 分片优化器状态、梯度和参数 |
SHARD_GRAD_OP | ZeRO Stage-2 | 分片优化器状态和梯度 |
NO_SHARD | ZeRO Stage-0 | 无需分片,每个 GPU 拥有完整副本 |
HYBRID_SHARD | ZeRO++ Stage-3 (zero_hpz_partition_size=<num_gpus_per_node>) | 节点内分片优化器状态、梯度和参数,节点间完整副本 |
训练 - 编译
1. 完整编译
使用 TorchDynamoPlugin 集成 torch.compile,实现模型编译:
from accelerate import Accelerator
from accelerate.utils import TorchDynamoPlugin
# Configure the compilation backend
dynamo_plugin = TorchDynamoPlugin(
backend="inductor", # Options: "inductor", "aot_eager", "aot_nvfuser", etc.
mode="default", # Options: "default", "reduce-overhead", "max-autotune"
fullgraph=True,
dynamic=False
)
# Initialize accelerator with the plugin
accelerator = Accelerator(dynamo_plugin=dynamo_plugin)
# This will apply torch.compile to your model
model = accelerator.prepare(model)
2. 局部编译
区域编译会针对同一类的重复块,并按顺序编译它们以命中编译器的缓存。
# TorchDynamoPlugin 中配置 use_regional_compilation=True,启用局部编译
# Configure the compilation backend
dynamo_plugin = TorchDynamoPlugin(
use_regional_compilation=True,
... # other parameters
)
# Initialize accelerator with the plugin
accelerator = Accelerator(dynamo_plugin=dynamo_plugin)
# This will apply compile_regions to your model
model = accelerator.prepare(model)
区域编译的好处:
- 可比性能:区域编译提供与完整编译类似的性能加速,特别是对于较大的模型。
- 更快的编译:区域编译显著减少了编译模型所需的时间,使其成为更有效的部署选择。
- 批次大小影响:随着批次大小的增大,编译策略之间的性能差异会减小,这表明在这些情况下编译开销的影响较小。
- 模型大小考虑:区域编译的好处在较大的模型中更为明显,因为可以节省大量编译时间。
- 实际应用:对于实际应用,区域编译是优化训练冷启动时间的实用选择,尤其是在处理大型模型时。
推理 - 大模型推理
Accelerate 允许使用不完全适合显卡的大模型进行推理。
1. 基础用法
# 加载 pytorch 模型,ModelClass 是超出设备 GPU 内存的模型
import torch
my_model = ModelClass(...)
state_dict = torch.load(checkpoint_file)
my_model.load_state_dict(state_dict)
使用大模型推理,第一步使用上下文管理器初始化空模型,由于是无参数的空模型,init_empty_weights 不需要任何内存:
from accelerate import init_empty_weights
with init_empty_weights():
my_model = ModelClass(...)
将权重加载到模型中进行推理:
load_checkpoint_and_dispatch() 方法在空模型中加载一个检查点,并在所有可用设备上分配每一层的权重。首先从最快的设备(GPU、MPS、XPU、NPU、MLU、SDAA、MUSA)开始,然后再移动到较慢的设备(CPU 和硬盘)。
设置 device_map="auto" 会首先填充 GPU 上的所有可用空间,然后是 CPU,最后如果内存仍然不足则填充硬盘。
from accelerate import load_checkpoint_and_dispatch
model = load_checkpoint_and_dispatch(
model, checkpoint=checkpoint_file, device_map="auto"
)
模型使用:
input = torch.randn(2, 3)
device_type = next(iter(model.parameters())).device.type
input = input.to(device_type)
output = model(input)
说明:每次输入经过某个层时,它都会从 CPU 发送到 GPU(或从磁盘发送到 CPU 再发送到 GPU),计算输出,然后该层从 GPU 中移除,如此循环往复。虽然这会增加一些推理开销,但它允许您在系统上运行任何大小的模型,只要最大的层能够装进 GPU。
完整代码:
import torch
from accelerate import init_empty_weights, load_checkpoint_and_dispatch
with init_empty_weights():
model = MyModel(...)
model = load_checkpoint_and_dispatch(
model, checkpoint=checkpoint_file, device_map="auto"
)
input = torch.randn(2, 3)
device_type = next(iter(model.parameters())).device.type
input = input.to(device_type)
output = model(input)
2. Hugging Face 生态
HuggingFace 的库(如 Transformers 或 Diffusers)在 from_pretrained 构造函数中支持大模型推理,添加 device_map="auto" 即可启用大模型推理。
加载 Big Sciences T0pp 110 亿参数模型:
from transformers import AutoModelForSeq2SeqLM
model = AutoModelForSeq2SeqLM.from_pretrained("bigscience/T0pp", device_map="auto")
加载模型后,之前的空模型初始化和数据调度步骤将被执行,此时模型已完全准备好利用机器上的所有资源。
通过 from_pretrained() 构造函数,还可以通过指定 torch_dtype 参数以较低的精度加载模型来节省更多内存:
from transformers import AutoModelForSeq2SeqLM
model = AutoModelForSeq2SeqLM.from_pretrained("bigscience/T0pp", device_map="auto", torch_dtype=torch.float16)
推理 - 分布式推理
分布式推理可以分为三类:
- 将整个模型加载到每个 GPU 上,并通过每个 GPU 的模型副本一次发送一批数据块
- 将模型的各个部分加载到每个 GPU 上并一次处理单个输入
- 将模型的各个部分加载到每个 GPU 上,并使用所谓的预定流水线并行性来结合两种现有技术
1. 每个 GPU 保留模型的完整副本
这是最耗费内存的解决方案,因为它要求每个 GPU 在给定时间在内存中保留模型的完整副本。
通常在执行此操作时,用户将模型发送到特定设备以从 CPU 加载它,然后将每个提示移动到不同的设备。
示例代码:
import torch
import torch.distributed as dist
from diffusers import DiffusionPipeline
pipe = DiffusionPipeline.from_pretrained("runwayml/stable-diffusion-v1-5", torch_dtype=torch.float16)
# 根据提示词进行推理
def run_inference(rank, world_size):
dist.init_process_group("nccl", rank=rank, world_size=world_size)
pipe.to(rank)
if torch.distributed.get_rank() == 0:
prompt = "a dog"
elif torch.distributed.get_rank() == 1:
prompt = "a cat"
result = pipe(prompt).images[0]
result.save(f"result_{rank}.png")
使用 Accelerator 优化上述代码:
import torch
from accelerate import PartialState # Can also be Accelerator or AcceleratorState
from diffusers import DiffusionPipeline
pipe = DiffusionPipeline.from_pretrained("runwayml/stable-diffusion-v1-5", torch_dtype=torch.float16)
distributed_state = PartialState()
pipe.to(distributed_state.device)
# Assume two processes
with distributed_state.split_between_processes(["a dog", "a cat"]) as prompt:
result = pipe(prompt).images[0]
result.save(f"result_{distributed_state.process_index}.png")
# 命令行执行 accelerate config 之后启动
accelerate launch distributed_inference.py
# 使用指定 config 文件启动
accelerate launch --config_file my_config.json distributed_inference.py
说明:如果有 3 个指令,但只有 2 个 GPU,在上下文管理器下,第一个 GPU 将接收前两个提示词,第二个 GPU 将接收第三个提示词:
import torch
from accelerate import PartialState # Can also be Accelerator or AcceleratorState
from diffusers import DiffusionPipeline
pipe = DiffusionPipeline.from_pretrained("runwayml/stable-diffusion-v1-5", torch_dtype=torch.float16)
distributed_state = PartialState()
pipe.to(distributed_state.device)
# Assume two processes
with distributed_state.split_between_processes(["a dog", "a cat", "a chicken"], apply_padding=True) as prompt:
result = pipe(prompt).images
2. 内存高效的流水线并行
# 在 CPU 上创建模型
from transformers import GPT2ForSequenceClassification, GPT2Config
config = GPT2Config()
model = GPT2ForSequenceClassification(config)
model.eval()
# 构建样例数据
input = torch.randint(
low=0,
high=config.vocab_size,
size=(2, 1024), # bs x seq_len
device="cpu",
dtype=torch.int64,
requires_grad=False,
)
执行跟踪并加载模型:
使用 inference.prepare_pippy() 函数,自动包装模型以实现流水线并行:
from accelerate.inference import prepare_pippy
example_inputs = {"input_ids": input}
model = prepare_pippy(model, example_args=(input,))
prepare_pippy 方法支持的参数:
split_points:决定在哪些层上分割模型。默认情况下,使用device_map="auto"。num_chunks:确定批次如何分割并输入模型(如num_chunks=1,四个分割点/四个 GPU 会有一个简单的 MP,其中单个输入在四个层分割点之间传递)。
分布式推理:
args = some_more_arguments
with torch.no_grad():
output = model(*args)
所有数据将仅保存在最后一个进程中:
from accelerate import PartialState
if PartialState().is_last_process:
print(output)