首页
/ Colossal-AI 流水线并行(Pipeline Parallel)实战指南:1F1B 调度、Interleaved 虚拟流水线与 Bert 微调

Colossal-AI 流水线并行(Pipeline Parallel)实战指南:1F1B 调度、Interleaved 虚拟流水线与 Bert 微调

2026-09-08 11:15:06作者:尤峻淳Whitney

本教程面向希望在单卡放不下的巨型模型上开展训练、但又不满足于朴素数据并行的开发者。文章以 Colossal-AI 官方英文文档 pipeline_parallel.md 为主体骨架,系统讲解流水线并行的两大核心调度——Non-interleaved(1F1B)与 Interleaved(虚拟流水线)的原理与差异,并深入 Colossal-AI 的 HybridParallelPlugin 与源码实现,最终给出一个可直接运行的 Bert + GLUE 流水线微调全流程示例。阅读完成后,你将掌握流水线并行的调度原理、microbatch 相关约束,以及通过 Booster API 配置与驱动流水线训练的方法。

预备知识

在阅读本文之前,建议先掌握以下 Colossal-AI 基础概念,本文后续内容将直接复用这些 API:

流水线并行相关的公开论文工作(Megatron-LM、GPipe 及 Colossal-AI 系统论文)给出了 1F1B 调度的原始思想来源,本文中的调度图例即源自 Megatron-LM 论文的经典示意。仓库中与本主题对应的完整可运行示例位于 examples/language/bert/finetune.py

快速总览

在 Colossal-AI 中,流水线并行采用的是由 Nvidia 在 Megatron-LM 中提出的 1F1B(one forward pass followed by one backward pass) 调度,而非 GPipe 那种"整批前向完成后再统一反向"的朴素调度。由于 ViT 与 ImageNet 组合过于庞大,官方教程以 Bert + GLUE 数据集作为演示对象。本文覆盖以下三部分内容:

  1. 1F1B 流水线的原理介绍;
  2. Non-interleaved 与 Interleaved 两种调度的使用方式;
  3. 使用流水线并行微调 Bert 的完整流程。

1F1B 流水线原理

从 GPipe 说起

GPipe 是最早被广泛采用的流水线并行方案之一。在 GPipe 中,一个 batch 会被切分为若干 microbatch,各设备依次完成这些 microbatch 的前向传播,并且——关键之处在于——只有当整个 batch 内所有 microbatch 的前向都完成后,才会开始执行反向传播。这种"前向全部完成后才反向"的做法导致设备在等待期间大量闲置,形成显著的 pipeline bubble(流水线气泡),设备利用率与显存效率都不理想。

为什么是 1F1B

1F1B(一次前向后紧跟一次反向)的设计初衷正是为了缓解上述气泡问题。总体而言,1F1B 比 GPipe 更高效(在显存上更优,甚至时间与显存同时更优)。1F1B 又有两种具体调度形式:Non-interleaved(非交错,又称常规 1F1B)与 Interleaved(交错,又称虚拟流水线/Virtual Pipeline)。二者在吞吐与显存行为上存在明显差异。

Non-interleaved 调度

Non-interleaved 调度的时间线可以划分为三个阶段:

  1. Warm-up(预热)阶段:位于流水线不同位置的 worker 依次执行数量不同的前向传播,使整个流水线"灌满";
  2. 稳态阶段:每个 worker 执行"一次前向 + 一次反向"的交替操作,前后端设备在此阶段满负荷协同工作;
  3. Cooldown(冷却)阶段:worker 陆续完成剩余的反向传播,流水线逐渐"排空"。

该模式比 GPipe 更省显存——因为每个 microbatch 的反向紧随其前向之后进行,激活值不需要像 GPipe 那样全部保留到整批前向结束。不过,它完成一轮(一个 batch 的)passes 所需的总时间与 GPipe 相当,即并未减少理论流水线气泡总量。

Interleaved 调度

Interleaved(交错)调度有以下两个关键特征:

  • 要求 microbatch 数量是流水线 stage 数量的整数倍(源码层面的强制校验见下文);
  • 每个设备负责多个不连续的层子集(称为 model chunk),而非一段连续层。

例如传统切分下:device 1 拥有 layer 1–4,device 2 拥有 layer 5–8,以此类推;而交错切分后:device 1 拥有 layer 1、2、9、10,device 2 拥有 layer 3、4、11、12,以此类推。通过这种方案,流水线中的每个设备被分配了多个 stage,而每个 stage 的计算量相应变小。

由于每个设备可以更早开始下一 chunk 的前向,设备间的气泡被进一步压缩,因此该模式在显存和时间上都更高效,代价是实现更复杂、模型需要支持多 chunk 切分(一般由 Shardformer 提供切层支持)。

Colossal-AI 的流水线实现架构

在 Colossal-AI 中,流水线并行由 scheduler(调度器)Shardformer 两个组件协作完成。Shardformer 负责对模型进行层切分,并改写模型的 forward 函数,使其与 scheduler 的收发约定兼容;scheduler 则驱动各 stage 之间按 1F1B 的节奏进行 microbatch 的收发与前反向执行。

仓库中流水线相关源码位于 colossalai/pipeline 目录:

需要说明的是:官方文档正文写作时描述的是"HybridParallelPlugin 默认使用 OneForwardOneBackwardScheduleInterleavedSchedule 将后续集成"。但从当前仓库源码看,该集成早已完成——hybrid_parallel_plugin.py 中已根据 pp_style 分别实例化 InterleavedScheduleZeroBubbleVPipeScheduler,三种风格均可直接选用。

HybridParallelPlugin 如何组织流水线

在 Colossal-AI 中,HybridParallelPlugin 封装了流水线执行策略:它负责组建流水线并行的通信组并创建对应 scheduler。当使用该插件 boost 模型时,模型的层会被调用 shardformer.optimize 完成切分,随后训练循环中通过 booster.execute_pipeline(...) 让模型分片逐段执行。数据并行、张量并行与流水线并行可以在这个插件中被自由组合,且自动满足关系:

dp_size = world_size / (tp_size * pp_size)

插件内部通过 ProcessGroupMesh 建立多维设备网格,流水线方向由 pp_axis 指定(dp_outside=True 时 pp 位于第 1 轴);PipelineStageManager 依据该网格轴查询当前 stage 号、前后 rank,并创建沿轴的点对点通信组。

Stage 划分与虚拟 stage

PipelineStageManager 提供两个与层切分直接相关的关键方法:

  • distribute_layers(num_layers, num_stages, num_model_chunks):把总层数均匀切给 num_stages * num_model_chunks 个逻辑位置,余数层从中间位置开始均摊(见 stage_manager.py);
  • get_stage_index(...):返回某 stage(及其每个 model chunk)负责的层下标区间 [start, end)(见 stage_manager.py)。

对于 interleaved/虚拟流水线,num_model_chunks 大于 1,is_first_stage() / is_last_stage() 会同时考虑 stagemodel_chunk_id;在训练日志或 loss 记录处判断是否"流水线末端"时,应调用 is_last_stage()(对应示例代码中的 is_pp_last_stage 判断)。

调度器内部流程(以 1F1B 为例)

阅读 one_f_one_b.pyrun_forward_backward(约第 359–441 行),可以看到与 Megatron-LM 论文一一对应的三段式实现:

  1. 依据当前 stage 计算 num_warmup_microbatches = num_stages - stage - 1(再对 num_microbatches 取 min),先执行这批纯前向的 warmup;
  2. 进入稳态:forward_stepbackward_step 交替进行,并通过 send_forward_recv_backward / send_backward_recv_forward 将"发送张量"与"接收梯度/张量"重叠成一次通信;
  3. 最后执行 cooldown 阶段的纯反向。

中间 stage 的 forward_step 会把输出对象沿流水线向下一 stage 传递,只有 is_last_stage() 才会调用 criterion 计算 loss,并将其按 1 / num_microbatches 归一后累加;backward_step 则利用上游传入的 output_obj_grad 对当前 stage 做 optimizer.backward_by_grad。值得注意的是,每个 batch 加载时(load_batch)有一个重要断言:

num_microbatches >= num_stages   # 1F1B 训练时要求

1F1B 训练要求 microbatch 数不少于流水线 stage 数,否则流水线无法灌满。而在 InterleavedSchedule.load_batch(见 interleaved_pp.py)中,断言则变为:

num_microbatch % num_stages == 0   # interleaved 要求为 stage 数的整数倍

这也正是文档"Interleaved 要求 microbatch 数是 stage 数的整数倍"这一约束在源码中的强制落点。此外,如果同时提供 num_microbatchesmicrobatch_size,源码约定 num_microbatches 优先,另一个参数被忽略;microbatch_size 缺省时由 batch_size / num_microbatches 推导,反之亦然,且必须满足 batch_size == microbatch_size * num_microbatches

execute_pipeline 的执行包装

execute_pipeline(见 hybrid_parallel_plugin.py)在调用 scheduler 的 forward_backward_step 前,会为 HybridParallelZeroOptimizer 建立 optimizer.no_sync() 上下文(其余优化器则使用 model.no_sync()),以避免每个 microbatch 都触发冗余的梯度规约(多个 microbatch 的梯度应当累积后统一规约一次)。整个调度执行结束后,插件再统一完成 model.sync_shared_params()(同步如 embedding 等跨 stage 共享参数梯度)、model.sync_sp_grads()(序列并行梯度)以及 model.sync_dp_grads() / optimizer.sync_dp_grads()(数据并行维度梯度同步),从而保证与普通数据并行训练语义一致。

HybridParallelPlugin 关键参数速查

结合 hybrid_parallel_plugin.py 的参数文档,与流水线最相关的参数如下:

参数 默认值 说明
tp_size 必填 张量并行度,设为 1 表示不使用张量并行
pp_size 必填 流水线 stage 数,设为 1 表示不使用流水线并行
sp_size None 序列并行度
precision 'fp16' 'fp16'/'bf16' 时启用对应混合精度,否则 fp32 训练
zero_stage 0 ZeRO 阶段(0/1/2);使用流水线时必须为 0 或 1,以避免高昂的梯度同步开销
enable_all_optimization False 一键开启 Shardformer 支持的全部优化(fused normalization、flash attention、JIT)
num_microbatches None 流水线的 microbatch 数;与 microbatch_size 二者必须提供一个
microbatch_size None microbatch 大小;若提供了 num_microbatches 则此项被忽略
initial_scale 2**16 AMP 初始 loss scale
pp_style '1f1b' 流水线风格,可选 '1f1b''interleaved''zbv'
num_model_chunks 1 interleaved/zbv 流水线的模型 chunk 数;1f1b 下必须为 1,interleaved 下必须大于 1
gradient_checkpoint_config None 梯度检查点配置
overlap_p2p True 是否让流水线 p2p 通信与计算重叠
fp8_communication False 是否启用流水线间 fp8 通信

源码在插件初始化时还会做若干一致性校验(见 hybrid_parallel_plugin.py):

  • 世界大小必须能被 tp_size * pp_size 整除;
  • pp_size > 1 时,num_microbatchesmicrobatch_size 必须至少提供其一;
  • pp_style 仅支持 "1f1b" / "interleaved" / "zbv"
  • pp_size > 1zero_stage 只能取 0 或 1。

使用流水线并行微调 Bert(完整示例)

下面这段流程改编自仓库 examples/language/bert/finetune.py(数据侧实现见同目录 data.pyGLUEDataBuilder),对应官方教程的完整代码。演示采用 2 个流水线 stage,配置 microbatch 粒度为 1(这些参数可按硬件规模调整)。

第一步:导入与全局配置

import argparse
from typing import Callable, List, Union

import torch
import torch.nn as nn
from data import GLUEDataBuilder
from torch.optim import Adam, Optimizer
from torch.optim.lr_scheduler import _LRScheduler as LRScheduler
from torch.utils.data import DataLoader
from tqdm import tqdm
from transformers import (
    AlbertForSequenceClassification,
    AutoConfig,
    BertForSequenceClassification,
    get_linear_schedule_with_warmup,
)

import colossalai
from colossalai.booster import Booster
from colossalai.booster.plugin import HybridParallelPlugin
from colossalai.cluster import DistCoordinator
from colossalai.nn.optimizer import HybridAdam

# Define some config
NUM_EPOCHS = 3
BATCH_SIZE = 32
LEARNING_RATE = 2.4e-5
WEIGHT_DECAY = 0.01
WARMUP_FRACTION = 0.1

coordinator = DistCoordinator()

def move_to_cuda(batch):
    return {k: v.cuda() for k, v in batch.items()}

第二步:定义 criterion、优化器与学习率调度

流水线训练要求 criterion 接收两个参数(模型输出与输入),它会作为 execute_pipeline 的入参被调度器内部调用:

# Define 'criterion' function with two inputs, which will be passed to 'execute_pipeline'.
def _criterion(outputs, inputs):
    return outputs.loss

# Define optimizer
lr = LEARNING_RATE
no_decay = ["bias", "LayerNorm.weight"]
optimizer_grouped_parameters = [
    {
        "params": [p for n, p in model.named_parameters() if not any(nd in n for nd in no_decay)],
        "weight_decay": WEIGHT_DECAY,
    },
    {
        "params": [p for n, p in model.named_parameters() if any(nd in n for nd in no_decay)],
        "weight_decay": 0.0,
    },
]

optimizer = HybridAdam(optimizer_grouped_parameters, lr=lr, eps=1e-8)


# Define lr_scheduler
total_steps = len(train_dataloader) * NUM_EPOCHS
num_warmup_steps = int(WARMUP_FRACTION * total_steps)
lr_scheduler = get_linear_schedule_with_warmup(
    optimizer,
    num_warmup_steps=num_warmup_steps,
    num_training_steps=total_steps,
)

参数分组刻意将 biasLayerNorm.weight 排除在权重衰减之外——这是 Bert 类模型微调的标准实践;优化器选用 HybridAdam(Colossal-AI 提供的 Adam 封装,可配合混合精度与梯度规约使用)。

第三步:创建 HybridParallelPlugin

plugin = HybridParallelPlugin(tp_size=1,
                                pp_size=2,
                                num_microbatches=None,
                                microbatch_size=1,
                                enable_all_optimization=True,
                                zero_stage=1,
                                precision='fp16',
                                initial_scale=1)
booster = Booster(plugin=plugin)

此处 tp_size=1 表示不启用张量并行,仅用 2 个 stage 的流水线;zero_stage=1 为流水线各 stage 内的数据并行维度启用 ZeRO-1;precision='fp16' 开启 AMP,initial_scale=1 指定初始 loss scale。前文源码已说明:当 num_microbatches=Nonemicrobatch_size=1 时,调度器会在加载 batch 时自动推导 num_microbatches = BATCH_SIZE / 1 = 32,因此该配置满足 1F1B 对 num_microbatches >= num_stages 的要求。

第四步:定义模型与 dataloader

# Define Bert model
cfg = AutoConfig.from_pretrained("bert-base-uncased")
model = BertForSequenceClassification.from_pretrained("bert-base-uncased", config=cfg).cuda()

# Define a dataloader
data_builder = GLUEDataBuilder(model_name,
                                plugin,
                                args.task,
                                train_batch_size=BATCH_SIZE,
                                eval_batch_size=BATCH_SIZE)
train_dataloader = data_builder.train_dataloader()

注意 GLUEDataBuilder 接收 plugin 作为参数——数据并行采样器的副本数(num_replicas)与 rank 划分正是依据插件内部的 dp 通信组计算得到的(见 prepare_dataloader 实现)。

第五步:使用 booster 统一 boost 训练组件

model, optimizer, _criterion, _, lr_scheduler = booster.boost(model,
                                                                optimizer,
                                                                criterion=_criterion,
                                                                lr_scheduler=lr_scheduler)

booster.boost 内部会调用 shardformer.optimize 对模型完成按 stage 的层切分(对 pp_size>1 的模型),并把 forward 改写为可被调度器驱动的流水线形式。

第六步:流水线训练循环

# Define a train function
def train_epoch(epoch: int, model: nn.Module, optimizer: Optimizer, _criterion: Callable, lr_scheduler: LRScheduler,
                train_dataloader: DataLoader, booster: Booster, coordinator: DistCoordinator):

    is_pp_last_stage = booster.plugin.stage_manager.is_last_stage()
    total_step = len(train_dataloader)

    model.train()
    optimizer.zero_grad()
    # convert train_dataloader to a iterator
    train_dataloader_iter = iter(train_dataloader)
    with tqdm(range(total_step),
              desc=f'Epoch [{epoch + 1}/{NUM_EPOCHS}]',
              disable=not (is_pp_last_stage)) as pbar:
        # Forward pass
        for _ in pbar:
            outputs = booster.execute_pipeline(train_dataloader_iter,
                                                model,
                                                _criterion,
                                                optimizer,
                                                return_loss=True)
            # Backward and optimize
            if is_pp_last_stage:
                loss = outputs['loss']
                pbar.set_postfix({'loss': loss.item()})

            optimizer.step()
            optimizer.zero_grad()
            lr_scheduler.step()

# Train model
for epoch in range(NUM_EPOCHS):
    train_epoch(epoch, model, optimizer, _criterion, lr_scheduler, train_dataloader, booster, coordinator)

这段代码有三个值得注意的要点:

  1. 只有流水线末段才产生 lossbooster.execute_pipeline 内部由调度器按 1F1B 节奏驱动各 stage;criterion 只在 is_last_stage() 的进程上执行,因此 outputs['loss'] 仅在末段 rank 非空。这里通过 booster.plugin.stage_manager.is_last_stage() 判断,并用它控制 tqdm 进度条只在末段 rank 打印,避免多进程刷屏;
  2. execute_pipeline 一次性覆盖前向 + 反向。调度器完成整个 batch 各 microbatch 的前反向后返回,循环体内只需执行 optimizer.step()optimizer.zero_grad()lr_scheduler.step(),这些调用在所有 rank 上同步执行(gradient accumulation 逻辑由插件内部统一处理);
  3. 优化器/学习率必须在所有 rank 上创建(各 stage 仅持有其负责的层参数),这是流水线并行的通用约束。

运行该示例至少需要 2 个 GPU(对应 pp_size=2),并通过多进程启动器拉起,例如 torchrun --nproc_per_node=2 examples/language/bert/finetune.py(也可使用 Colossal-AI 提供的分布式 launcher)。官方教程中"使用 2 个流水线 stage、micro batch 为 1"即对应上述 pp_size=2microbatch_size=1 的配置,读者可以依据实际显存与设备数修改这两个参数。

更进一步:切换为 Interleaved 与更多并行风格

若希望体验文档所述"显存与时间双优"的 Interleaved 调度,只需修改插件配置并满足两个前置条件:

plugin = HybridParallelPlugin(
    tp_size=1,
    pp_size=2,
    pp_style="interleaved",   # 切换为交错/虚拟流水线
    num_model_chunks=2,       # 每个设备负责多个层 chunk,必须大于 1
    num_microbatches=8,       # 须为 pp_size * num_model_chunks 的整数倍
    enable_all_optimization=True,
    zero_stage=1,
    precision='fp16',
)

从源码可知,interleaved 模式要求 num_model_chunks > 1num_microbatch % num_stages == 0;为充分填满虚拟 stage,建议让 num_microbatches 同时覆盖多个 chunk 轮次。仓库同样支持 pp_style="zbv"(零气泡调度,需配合 scheduler_nodes 使用,见 zero_bubble_pp.py),对零气泡变体感兴趣的读者可进一步阅读 zerobubble_pipeline_parallelism.md。上述配置下模型是否被自动切为多 chunk、各 chunk 层如何分配,均由 Shardformer 配合 PipelineStageManager 完成,无需手工改动模型结构。

流水线并行相关的自动化测试位于 tests/test_pipeline,其中调度器行为、stage 间通信与 stage 管理均有对应用例覆盖,可作为二次开发与验证的参考。

小结

本文以 Colossal-AI 官方流水线并行教程为骨架,梳理了 GPipe 到 1F1B 的演进动机、Non-interleaved 与 Interleaved 两种调度的时间线差异与显存行为,并对照源码说明了 HybridParallelPlugin 中调度器与 Shardformer 的分工、microbatch 约束的强制校验位置,以及 execute_pipeline 背后的梯度规约语义。最后给出的 Bert + GLUE 微调全流程直接对接仓库 examples/language/bert/finetune.py。掌握以上内容后,你便可以在此基础上自由组合数据并行(ZeRO)、张量并行与流水线并行,将单卡放不下的模型迁移到多卡/多机环境上训练。

登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.14 K
2.76 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
860
1.35 K
docsdocs
暂无描述
Markdown
899
5.83 K
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
925
1.85 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.84 K
1.02 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
533
601
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
1.03 K
525
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.37 K
1.46 K
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
548
395