ColossalAI Pipeline Inference 管道并行推理实践:原理、MicroBatch 调度与性能评测
【免费下载链接】ColossalAIMaking large AI models cheaper, faster and more accessible项目地址: https://gitcode.com/GitHub_Trending/co/ColossalAI
本文聚焦 ColossalAI 中专门为大模型生成式推理设计的
Pipeline Inference模块(文档主体位于 colossalai/legacy/inference/pipeline/README.md)。当超大模型权重无法装入单张 GPU 时,该模块借助流水线并行把模型切分到多卡上完成推理;读者将系统掌握其"推理阶段几乎无 bubble"的设计动机、PPInferEngine/MicroBatchManager/GenerateSchedule三大组件的协作机制、PP(Pipeline Parallel)与 MicroBatch 状态机实现细节,以及可复现的快速开始与基准测试方法。
一、背景:为什么大模型推理也需要 Pipeline Parallel
在推理(Inference)阶段,虽然不再需要像训练那样保存前向传播的中间激活值用于反向传播,但一些较大模型的权重依然无法在单张 GPU 上放下。此时必须借助张量并行(Tensor Parallel)或流水线并行(Pipeline Parallel)等模型并行手段来降低单卡显存占用。
流水线并行是传统模型并行方案之一,具备两个突出优点:
- 通信开销小:无需像张量并行那样在每个 Transformer Layer 内部频繁做 All-Reduce,流水线阶段之间仅需传递 hidden states 与 KV Cache,通信模式简单规整;
- 切分布局简单:模型按层切分为若干连续 Stage,布局直观、易于理解。
传统上流水线并行最受人诟病的问题是bubble(流水线气泡),即不同 Stage 之间因为前向/反向交替而产生的空转时间。而该模块的设计核心洞察在于:推理阶段不存在引发 bubble 的反向传播,因此只要各 Stage 上序列长度一致,流水线并行在理想情况下几乎是零 bubble 的。这正是 ColossalAI 把 pipeline 范式迁移到推理场景的根本动机。
二、整体设计:三大组件协同
从 colossalai/legacy/inference/pipeline/README.md 的定义看,Pipeline Inference 由三部分构成:PPInferEngine、MicroBatchManager与generateschedule(仓库中的对应实现见 colossalai/pipeline/schedule/generate.py)。
用户输入 Batch │ ▼ ┌─────────────────────────┐ │ PPInferEngine │ 高层 API:环境初始化 + 驱动推理 │ - 构建 PipelineStageManager │ - 用 Shardformer 切分模型 / 分配各 Stage │ - 维护 MicroBatchManager 与 generate schedule └─────────────────────────┘ │ 按 MicroBatch 切分下发 ▼ ┌─────────────────────────┐ │ MicroBatchManager │ 管理每个 MicroBatch 的推理状态与信息 │ - new_tokens / kvcache │ PREFILL → GENERATE → DONE 状态流转 └─────────────────────────┘ │ ▼ ┌─────────────────────────┐ │ generate schedule │ 定义流水线推理布局与阶段间通信 │ - pp_size == 2 : P2P │ │ - pp_size > 2 : broadcast └─────────────────────────┘1. PPInferEngine:面向用户的高层 API
PPInferEngine承担两类职责:
- 初始化 pipeline 推理环境:用
PipelineStageManager组织多卡流水线拓扑,并用ShardFormer把模型按层切分、下发到各 Stage; - 运行 pipeline 推理:将 batch 切成 micro-batch 送入流水线,直至所有 micro-batch 完成生成并回收结果。
说明:本文档对应的早期实现中该引擎从
colossalai.inference导入,而在当前仓库结构中,Pipeline Inference 相关代码已归档至 legacy 目录(见 colossalai/legacy/inference/pipeline/),顶层仅导出MicroBatchManager(见 colossalai/legacy/inference/pipeline/init.py)。阅读本文案例时请注意按你所处代码分支的实际导入路径调整。
2. MicroBatchManager:MicroBatch 信息管理器
它负责跟踪流水线内每个 micro-batch 的运行轨迹:
- 记录每个 micro-batch 的信息,例如新生成的 token(new tokens)与 KV Cache;
- 记录每个 micro-batch 的推理状态,例如当前处于 prefill(预填充)、generate(逐 token 生成)还是 done(完成);
- 更新 micro-batch 信息,推进生成与序列长度。
3. generate schedule:流水线推理布局
generateschedule 实现流水线推理的具体排布。阶段间通信策略按流水线深度自适应:
- pp_size = 2(2 个流水线阶段):使用
torch.distributed.P2Pop实现阶段间通信,主要用于规避通信竞态(race communication); - pp_size > 2:改用
torch.distributed.broadcast,其速度比 P2P 方式更快。
在仓库中,GenerateSchedule派生自PipelineSchedule,通过PipelineP2PCommunication与MicroBatchManager协同完成多 Stage 推理(见 colossalai/pipeline/schedule/generate.py)。它持有 action interval buffer 用于保存 stage 之间传递的中间 hidden states 与新 token,并把 embedding/lm_head 层放在同一设备上以节省显存。
三、源码级深挖:MicroBatch 状态机与描述符
MicroBatch 的底层实现在 colossalai/legacy/inference/pipeline/microbatch_manager.py,它定义了推理流程中的核心状态语义。
3.1 状态枚举:PREFILL / GENERATE / DONE / COOLDOWN
class Status(Enum): PREFILL = 1 GENERATE = 2 DONE = 3 COOLDOWN = 4状态判断基于cur_length与目标长度的关系(MicroBatchDescription.state,microbatch_manager.py):
- 当前长度等于
target_length(= 输入长度 +max_output_len)→DONE; - 当前长度等于
target_length - 1→COOLDOWN(最后一个 token 生成前的收尾阶段); - 否则 →
GENERATE。
其中max_output_len即用户指定的new_length,即期望继续生成的 token 数。
3.2 两类描述符:头阶段与主体阶段
由于流水线首段(第 0 个 stage)负责接收原始文本并产出第一个 token,其行为与后续 stage 不同,仓库分别实现了两个描述符类:
| 描述符 | 适用阶段 | 输入 | 职责 |
|---|---|---|---|
HeadMicroBatchDescription | stage 0(头) | input_ids+attention_mask | 保存原始输入与new_tokens,负责逐 token 拼接生成结果并扩展 attention mask |
BodyMicroBatchDescription | stage 1..N-1(主体) | hidden_states+past_key_values | 只接收上一阶段传来的 hidden states,依据 KV Cache 的seq_len推断当前长度 |
关键差异在cur_length的判定逻辑(microbatch_manager.py):
- 头阶段:尚未生成 token 时长度为
mb_length;生成后为mb_length + len(new_tokens[0]); - 主体阶段:直接读取
infer_state.seq_len.max(),即由 KV Cache 记录的序列长度决定。
class HeadMicroBatchDescription(MicroBatchDescription): def _update_newtokens(self, new_token: torch.Tensor): if self.new_tokens is None: self.new_tokens = new_token else: self.new_tokens = torch.cat([self.new_tokens, new_token], dim=-1) def _update_attnmask(self): # 每生成一个新 token,向 attention_mask 追加一个有效位 self.attn_mask = torch.cat( (self.attn_mask, torch.ones((self.attn_mask.shape[0], 1), dtype=torch.int64, device="cuda")), dim=-1 )MicroBatchManager本身以环形 buffer 管理多个 micro-batch(microbatch_manager.py):
micro_batch_size:单个 micro-batch 的样本数;micro_batch_buffer_size:micro-batch buffer 深度,文档建议与流水线 stage 数量一致;- 关键行为:
step()推进当前描述符状态并返回cur_state;is_micro_batch_done()检查是否所有 micro-batch 均为DONE;export_new_tokens()把 buffer 内所有已生成 token 汇总为列表返回给上层;clear()清空描述符并释放 KV Cache。
四、快速开始:完整可运行示例
文档给出的最小示例以 Llama 为例,假设将模型切分为2 个 pipeline stage推理:
from colossalai.inference import PPInferEngine from colossalai.inference.pipeline.policies import LlamaModelInferPolicy import colossalai from transformers import LlamaForCausalLM, LlamaTokenizer colossalai.launch_from_torch() model = LlamaForCausalLM.from_pretrained("/path/to/model") tokenizer = LlamaTokenizer.from_pretrained("/path/to/model") # assume the model is inferred with 2 pipeline stages inferengine = PPInferEngine(pp_size=2, model=model, model_policy=LlamaModelInferPolicy(), new_length=32) input = ["Introduce a landmark in London", "Introduce a landmark in Singapore"] data = tokenizer(input, return_tensors='pt') output = inferengine.inference(data.to('cuda')) print(tokenizer.batch_decode(output))执行要点拆解:
colossalai.launch_from_torch()从torch.distributed启动器获取分布式环境(多卡场景需配合colossalai run,见第六节);- 加载 Hugging Face Llama 权重与 tokenizer,路径替换为本地权重目录;
- 构造
PPInferEngine:pp_size=2表示切成 2 个流水线阶段(需 2 个 GPU / 进程),model_policy传入 Llama 的推理切分策略LlamaModelInferPolicy(),new_length=32表示最多续写 32 个新 token; tokenizer对两个句子批量编码后调用engine.inference(data.to('cuda')),返回每个请求生成的 token,用tokenizer.batch_decode还原为文本。
五、PPInferEngine 关键参数与更多配置项
在文档示例的基础上,仓库自带的基准脚本 colossalai/legacy/inference/pipeline/benchmark/benchmark.py 展示了更完整的参数组合:
engine = PPInferEngine( pp_size=args.pp_size, dtype=args.dtype, micro_batch_size=args.mb_size, new_length=args.new_length, model=model, model_policy=LlamaModelInferPolicy(), verbose=True, max_batch_size=args.mb_size, max_input_len=args.seq_len, max_output_len=args.seq_len + args.new_length + 256, )核心参数含义归纳如下:
| 参数 | 含义 | 典型取值 / 说明 |
|---|---|---|
pp_size | 流水线 stage 数量 | 通常等于使用的 GPU 数;为 2 时走 P2P 通信路径,大于 2 走 broadcast 路径 |
model | 待推理模型 | 支持 Hugging Face 风格的LlamaForCausalLM等 |
model_policy | 模型切分/forward 策略 | Llama 场景传LlamaModelInferPolicy() |
new_length | 每个请求期望生成的 token 数 | 示例中为 32 |
micro_batch_size | 每个 micro-batch 的样本数 | 与吞吐/显存折中,见性能表 |
dtype | 推理精度 | fp16/bf16等 |
max_batch_size | 显存预分配的最大 batch | 需覆盖实际 batch |
max_input_len | 预分配的最大输入长度 | 应大于等于最长输入序列 |
max_output_len | 预分配的最大输出长度 | 需要覆盖seq_len + new_length,基准脚本额外 +256 作为冗余 |
verbose | 是否输出流水线运行明细 | 基准测试用于采集时间戳 |
其中micro_batch_size与 buffer 机制直接决定显存中 KV Cache 的预分配规模与吞吐上限,是调优的核心旋钮。
六、多卡启动与 Benchmark:复现文档性能数字
6.1 基准脚本的输入与输出
benchmark.py支持三种模型规模:toy(8 层随机 Llama 配置,用于快速功能自检)、7b、13b(从decapoda-research预训练权重构造)。它会在 rank 0 上统计并落盘以下指标(benchmark.py):
- Average prefill time / Average encode time
- Average micro-batch end2end time / Whole-batch end2end time
- Micro-batch / Whole-batchPer Token Latency(ms)
- Throughput(tokens/s)
- FLOPS(按参数量估算)
同时记录 GPU 的 free / allocated / reserved 显存数据,日志文件名格式为llama-{model}{dtype}_pp{pp_size}_{seq_len}_{new_length}_bsz{batch}_mbsz{mb}.log。
6.2 启动命令与多组测试矩阵
仓库提供了配套的批量启动脚本 colossalai/legacy/inference/pipeline/benchmark/run.sh,其核心启动方式为:
colossalai run --nproc_per_node 2 --master_port 29800 ./benchmark.py \ --model="7b" \ --dtype="fp16" \ --batch_size=${BATCH_SIZE} \ --seq_len=1024 \ --new_length=128 \ --mb_size=$((${BATCH_SIZE}/2)) \ --pp_size=2要点:
- 通过
colossalai run启动分布式任务,--nproc_per_node 2与pp_size=2对应(2 张 GPU 组成 2 级流水线); - 脚本按
BATCH_SIZE ∈ {2,4,8,16}循环压测,micro-batch 大小取batch_size/2,与流水线 stage 数一致,正好填满 buffer; - 除
seq_len=1024 / new_length=128外,还覆盖seq_len=512 / new_length=512长生成场景,以及 7b/13b 两种规模。
七、性能数据:Pipeline Inference vs Hugging Face Pipeline
文档在2 × A10 20G与2 × A800 80G两种环境下对比了Pipeline Inference与 Hugging Face pipeline 的吞吐(tokens/s),测试条件为 input length=1024、output length=128。以下数字均取自原文档,作为结果复述。
7.1 A10 环境(7b / 13b,fp16)
Llama-7b,fp16(表中batch_size(micro_batch size)):
| batch_size(micro_batch size) | 2(1) | 4(2) | 8(4) | 16(8) | 32(8) | 32(16) |
|---|---|---|---|---|---|---|
| Pipeline Inference | 40.35 | 77.1 | 139.03 | 232.7 | 257.81 | OOM |
| Hugging Face | 41.43 | 65.30 | 91.93 | 114.62 | OOM | OOM |
Llama-13b,fp16:
| batch_size(micro_batch size) | 2(1) | 4(2) | 8(4) | 16(4) |
|---|---|---|---|---|
| Pipeline Inference | 25.39 | 47.09 | 83.7 | 89.46 |
| Hugging Face | 23.48 | 37.59 | 53.44 | OOM |
7.2 A800 环境(7b / 13b,fp16)
Llama-7b,fp16:
| batch_size(micro_batch size) | 2(1) | 4(2) | 8(4) | 16(8) | 32(16) |
|---|---|---|---|---|---|
| Pipeline Inference | 57.97 | 110.13 | 213.33 | 389.86 | 670.12 |
| Hugging Face | 42.44 | 76.5 | 151.97 | 212.88 | 256.13 |
Llama-13b,fp16:
| batch_size(micro_batch size) | 2(1) | 4(2) | 8(4) | 16(8) | 32(16) |
|---|---|---|---|---|---|
| Pipeline Inference | 41.78 | 94.18 | 172.67 | 310.75 | 470.15 |
| Hugging Face | 36.57 | 68.4 | 105.81 | 139.51 | 166.34 |
从结果可以观察到的规律(原文档数据直接反映的事实):
- batch 越大优势越明显:在 A10 7b 上 batch 4 时两者持平,batch 8/16 后 Pipeline Inference 吞吐提升至 1.5~2 倍;这是因为流水线多卡并行分摊了显存并持续吞吐生成 token;
- 显存上限被显著推高:Hugging Face 在 A10 上 32 batch(7b)、16 batch(13b)即 OOM,而 Pipeline Inference 可推进到更大 batch 才 OOM;
- A800 大显存下差距更大:A800 7b 的 32 batch 场景,Pipeline Inference 约 670 tokens/s,约为 Hugging Face 的 2.6 倍。
需要注意这些数字有明确适用前提:2 卡流水线、fp16、给定 batch 与序列长度组合、指定硬件。实际复现时吞吐受驱动、框架版本、显存分配策略影响,建议以本仓库 benchmark 脚本在本地复测为准。
八、使用限制与阅读指引
- 代码归档位置:本仓库当前结构中,Pipeline Inference 相关代码(
MicroBatchManager、状态机描述符、benchmark)位于 colossalai/legacy/inference/pipeline/,其中描述文件即本文依据的 README,microbatch_manager.py为状态机核心实现; - schedule 实现:流水线推理调度
GenerateSchedule位于 colossalai/pipeline/schedule/generate.py,与MicroBatchManager、PipelineP2PCommunication协作,属于当前仍在维护的colossalai.pipeline.schedule模块; - 上层框架衔接:Pipeline Inference 隶属于 ColossalAI 的 legacy inference 体系,同一目录下还有基于张量并行的
engine.py/kvcache_manager.py/batch_infer_state.py(见 colossalai/legacy/inference/tensor_parallel/),两者在 KV Cache 管理(MemoryManager、BatchInferState)上复用同一套底层设施; - 示例导入路径:文档与 benchmark 中
from colossalai.inference import PPInferEngine等导入对应本文档产生时的分支布局;若你在当前 trunk 直接运行遇到导入错误,请以实际legacy目录结构或对应 release 分支为准。
总体而言,ColossalAI 的 Pipeline Inference 用一个清晰的 micro-batch 状态机 + 流水线 schedule,在几乎零 bubble 的前提下把多卡显存与算力高效转化为推理吞吐,特别适合模型权重超出单卡显存、且请求 batch 较大的生成式推理场景。
【免费下载链接】ColossalAIMaking large AI models cheaper, faster and more accessible项目地址: https://gitcode.com/GitHub_Trending/co/ColossalAI
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考