Colossal-AI 流水线并行(Pipeline Parallel)实战指南:1F1B 调度、Interleaved 虚拟流水线与 Bert 微调
本教程面向希望在单卡放不下的巨型模型上开展训练、但又不满足于朴素数据并行的开发者。文章以 Colossal-AI 官方英文文档 pipeline_parallel.md 为主体骨架,系统讲解流水线并行的两大核心调度——Non-interleaved(1F1B)与 Interleaved(虚拟流水线)的原理与差异,并深入 Colossal-AI 的 HybridParallelPlugin 与源码实现,最终给出一个可直接运行的 Bert + GLUE 流水线微调全流程示例。阅读完成后,你将掌握流水线并行的调度原理、microbatch 相关约束,以及通过 Booster API 配置与驱动流水线训练的方法。
预备知识
在阅读本文之前,建议先掌握以下 Colossal-AI 基础概念,本文后续内容将直接复用这些 API:
- 并行范式总览:了解数据并行、张量并行、流水线并行三种范式的定位差异;
- 使用 Booster 进行训练:理解
Booster的boost与训练驱动方式; - Shardformer:理解流水线并行中模型切层、
forward重写的执行者; - Booster 插件详解:本文主角
HybridParallelPlugin的完整参数手册。
流水线并行相关的公开论文工作(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 数据集作为演示对象。本文覆盖以下三部分内容:
- 1F1B 流水线的原理介绍;
- Non-interleaved 与 Interleaved 两种调度的使用方式;
- 使用流水线并行微调 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 调度的时间线可以划分为三个阶段:
- Warm-up(预热)阶段:位于流水线不同位置的 worker 依次执行数量不同的前向传播,使整个流水线"灌满";
- 稳态阶段:每个 worker 执行"一次前向 + 一次反向"的交替操作,前后端设备在此阶段满负荷协同工作;
- 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 目录:
- colossalai/pipeline/schedule/base.py:抽象基类
PipelineSchedule,定义统一的forward_backward_step接口——模型、数据迭代器、criterion、优化器作为输入,返回包含loss与outputs的字典; - colossalai/pipeline/schedule/one_f_one_b.py:
OneForwardOneBackwardSchedule(Non-interleaved 1F1B); - colossalai/pipeline/schedule/interleaved_pp.py:
InterleavedSchedule(交错调度,按num_model_chunks多 chunk 执行); - colossalai/pipeline/schedule/zero_bubble_pp.py:
ZeroBubbleVPipeScheduler(零气泡变体,pp_style="zbv"); - colossalai/pipeline/schedule/init.py:上述调度器的统一导出;
- colossalai/pipeline/stage_manager.py:
PipelineStageManager,管理流水线 stage 归属、p2p 通信组与虚拟 stage; - colossalai/pipeline/p2p.py:stage 之间的点对点张量收发实现。
需要说明的是:官方文档正文写作时描述的是"HybridParallelPlugin 默认使用 OneForwardOneBackwardSchedule,InterleavedSchedule 将后续集成"。但从当前仓库源码看,该集成早已完成——hybrid_parallel_plugin.py 中已根据 pp_style 分别实例化 InterleavedSchedule 与 ZeroBubbleVPipeScheduler,三种风格均可直接选用。
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() 会同时考虑 stage 与 model_chunk_id;在训练日志或 loss 记录处判断是否"流水线末端"时,应调用 is_last_stage()(对应示例代码中的 is_pp_last_stage 判断)。
调度器内部流程(以 1F1B 为例)
阅读 one_f_one_b.py 的 run_forward_backward(约第 359–441 行),可以看到与 Megatron-LM 论文一一对应的三段式实现:
- 依据当前 stage 计算
num_warmup_microbatches = num_stages - stage - 1(再对num_microbatches取 min),先执行这批纯前向的 warmup; - 进入稳态:
forward_step与backward_step交替进行,并通过send_forward_recv_backward/send_backward_recv_forward将"发送张量"与"接收梯度/张量"重叠成一次通信; - 最后执行 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_microbatches 与 microbatch_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_microbatches与microbatch_size必须至少提供其一;pp_style仅支持"1f1b"/"interleaved"/"zbv";pp_size > 1时zero_stage只能取 0 或 1。
使用流水线并行微调 Bert(完整示例)
下面这段流程改编自仓库 examples/language/bert/finetune.py(数据侧实现见同目录 data.py 的 GLUEDataBuilder),对应官方教程的完整代码。演示采用 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,
)
参数分组刻意将 bias 与 LayerNorm.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=None 而 microbatch_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)
这段代码有三个值得注意的要点:
- 只有流水线末段才产生 loss。
booster.execute_pipeline内部由调度器按 1F1B 节奏驱动各 stage;criterion 只在is_last_stage()的进程上执行,因此outputs['loss']仅在末段 rank 非空。这里通过booster.plugin.stage_manager.is_last_stage()判断,并用它控制 tqdm 进度条只在末段 rank 打印,避免多进程刷屏; execute_pipeline一次性覆盖前向 + 反向。调度器完成整个 batch 各 microbatch 的前反向后返回,循环体内只需执行optimizer.step()、optimizer.zero_grad()与lr_scheduler.step(),这些调用在所有 rank 上同步执行(gradient accumulation 逻辑由插件内部统一处理);- 优化器/学习率必须在所有 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=2、microbatch_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 > 1 且 num_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)、张量并行与流水线并行,将单卡放不下的模型迁移到多卡/多机环境上训练。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0631
MiniCPM5-2BMiniCPM5-2B 是一款面向端侧、本地部署和资源受限场景的 2B 稠密 Transformer,能够达到同尺寸开源模型 SOTA 水平。Markdown00
video-shotcraftAI宣传片skill,使用 Remotion 制作电影级产品视频:提供106 张镜头配方卡和可复用的视频魔板。适用于 Claude Code 与 Codex以及所有其他智能体Markdown00
HivisionIDPhotos⚡️HivisionIDPhotos: a lightweight and efficient AI ID photos tools. 一个轻量级的AI证件照制作算法。Python09
DragonOSDragonOS is an operating system developed from scratch using Rust, with Linux compatibility. It is designed for **Serverless** scenarios. 使用Rust从0自研内核,具有Linux兼容性的操作系统,面向云计算Serverless场景而设计。Rust00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00