Article

加速训练 Accelerate 官方文档

更新于:2026-07-20

简介

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()

说明:

  1. 训练脚本中导入并实例化 Accelerate 类,Accelerate 类会初始化分布式训练参数,根据代码的启动方式自动检测训练环境,例如单 GPU、多 GPU、多 TPU 等。

    from accelerate import Accelerator
    
    accelerator = Accelerator()
  2. Accelerate 类会自动将模型和数据放置在合适的设备上。

    device = accelerator.device
  3. 将优化器、模型、数据加载器、学习率调度器传递给 prepare() 方法,此方法将模型包装在针对分布式优化的容器中,使用 Accelerate 版本的优化器和调度器,创建分片版本的数据加载器,以便在 GPU 或 TPU 设置之间分发。

    model, optimizer, train_dataloader, lr_scheduler = accelerator.prepare(
        model, optimizer, train_dataloader, lr_scheduler
    )
  4. 替换原有的 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_HOMEhuggingface/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 launchtorchrun 开始多节点训练运行。

注意:所有节点都运行该命令才能启动分布式训练,而不仅仅是从主节点运行。可以借助 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_stateload_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 进行训练
float3284.95 MB418.18 MB1.61 GB
float1642.47 MB206.59 MB826.36 MB
int821.24 MB103.29 MB413.18 MB
int410.62 MB51.65 MB206.59 MB

a. 指定库进行查询

如果无法自动确定模型的来源,例如 bert-base-cased,可以传入库名。

# 使用 transformers 库
accelerate estimate-memory HuggingFaceM4/idefics-80b-instruct --library_name transformers

模型加载时的内存使用情况:

数据类型最大层总大小使用 Adam 进行训练
float323.02 GB297.12 GB1.16 TB
float161.51 GB148.56 GB594.24 GB
int8772.52 MB74.28 GB297.12 GB
int4386.26 MB37.14 GB148.56 GB
# 使用 timm 库
accelerate estimate-memory timm/resnet50.a1_in1k --library_name timm

模型加载时的内存使用情况:

数据类型最大层总大小使用 Adam 进行训练
float329.0 MB97.7 MB390.78 MB
float164.5 MB48.85 MB195.39 MB
int82.25 MB24.42 MB97.7 MB
int41.12 MB12.21 MB48.85 MB

b. 指定数据类型

通过 --dtypes 参数指定数据类型。

accelerate estimate-memory bert-base-cased --dtypes float32 float16

模型加载时的内存使用情况:

数据类型最大层总大小使用 Adam 进行训练
float3284.95 MB413.18 MB1.61 GB
float1642.47 MB206.59 MB826.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 参数,在运行时记录日志
  • namestr
    • 跟踪器的唯一标识
  • requires_logging_directorybool
    • 跟踪器是否需要指定 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:调度选项允许您控制何时启用性能分析。这对于长时间运行的作业非常有用,可以避免收集过多数据。可用的键包括 waitwarmup
  • on_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 设置详细级别(INFODEBUGWARNINGERROR),例如,设置 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 的运行调试模式立即捕获张量形状不一致的问题,以下为多种配置方法,任选其一即可:

  1. 命令行配置

    accelerate launch --debug {my_script.py} --arg1 --arg2
  2. 环境变量配置

    ACCELERATE_DEBUG_MODE="1" torchrun {my_script.py} --arg1 --arg2
  3. 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_triggercheck_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_precisionno 表示 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_fileDeepSpeed 配置文件的路径(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 压缩与其他通信策略结合使用。目前支持的包装器包括 NOFP16BF16

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 层类名(区分大小写)的逗号分隔字符串,例如 BertLayerGPTJBlockT5Block 等。共享权重(如嵌入层)的子模块不应该位于不同的 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_SHARDZeRO Stage-3分片优化器状态、梯度和参数
SHARD_GRAD_OPZeRO Stage-2分片优化器状态和梯度
NO_SHARDZeRO Stage-0无需分片,每个 GPU 拥有完整副本
HYBRID_SHARDZeRO++ 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)

推理 - 分布式推理

分布式推理可以分为三类:

  1. 将整个模型加载到每个 GPU 上,并通过每个 GPU 的模型副本一次发送一批数据块
  2. 将模型的各个部分加载到每个 GPU 上并一次处理单个输入
  3. 将模型的各个部分加载到每个 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)