news 2026/9/10 4:02:48

ColossalAI Pipeline Inference 管道并行推理实践:原理、MicroBatch 调度与性能评测

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
ColossalAI Pipeline Inference 管道并行推理实践:原理、MicroBatch 调度与性能评测

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 由三部分构成:PPInferEngineMicroBatchManagergenerateschedule(仓库中的对应实现见 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承担两类职责:

  1. 初始化 pipeline 推理环境:用PipelineStageManager组织多卡流水线拓扑,并用ShardFormer把模型按层切分、下发到各 Stage;
  2. 运行 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,通过PipelineP2PCommunicationMicroBatchManager协同完成多 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 - 1COOLDOWN(最后一个 token 生成前的收尾阶段);
  • 否则 →GENERATE

其中max_output_len即用户指定的new_length,即期望继续生成的 token 数。

3.2 两类描述符:头阶段与主体阶段

由于流水线首段(第 0 个 stage)负责接收原始文本并产出第一个 token,其行为与后续 stage 不同,仓库分别实现了两个描述符类:

描述符适用阶段输入职责
HeadMicroBatchDescriptionstage 0(头)input_ids+attention_mask保存原始输入与new_tokens,负责逐 token 拼接生成结果并扩展 attention mask
BodyMicroBatchDescriptionstage 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_stateis_micro_batch_done()检查是否所有 micro-batch 均为DONEexport_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))

执行要点拆解:

  1. colossalai.launch_from_torch()torch.distributed启动器获取分布式环境(多卡场景需配合colossalai run,见第六节);
  2. 加载 Hugging Face Llama 权重与 tokenizer,路径替换为本地权重目录;
  3. 构造PPInferEnginepp_size=2表示切成 2 个流水线阶段(需 2 个 GPU / 进程),model_policy传入 Llama 的推理切分策略LlamaModelInferPolicy()new_length=32表示最多续写 32 个新 token;
  4. 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 配置,用于快速功能自检)、7b13b(从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 2pp_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 20G2 × 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 Inference40.3577.1139.03232.7257.81OOM
Hugging Face41.4365.3091.93114.62OOMOOM

Llama-13b,fp16

batch_size(micro_batch size)2(1)4(2)8(4)16(4)
Pipeline Inference25.3947.0983.789.46
Hugging Face23.4837.5953.44OOM

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 Inference57.97110.13213.33389.86670.12
Hugging Face42.4476.5151.97212.88256.13

Llama-13b,fp16

batch_size(micro_batch size)2(1)4(2)8(4)16(8)32(16)
Pipeline Inference41.7894.18172.67310.75470.15
Hugging Face36.5768.4105.81139.51166.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,与MicroBatchManagerPipelineP2PCommunication协作,属于当前仍在维护的colossalai.pipeline.schedule模块;
  • 上层框架衔接:Pipeline Inference 隶属于 ColossalAI 的 legacy inference 体系,同一目录下还有基于张量并行的engine.py/kvcache_manager.py/batch_infer_state.py(见 colossalai/legacy/inference/tensor_parallel/),两者在 KV Cache 管理(MemoryManagerBatchInferState)上复用同一套底层设施;
  • 示例导入路径:文档与 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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/10 4:01:32

Fanuc FOCAS二次开发:C#调用DLL对接数控机床

简介:本资源是面向工业自动化领域开发者与数控系统集成工程师的Fanuc数控机床二次开发核心工具包,聚焦Focas协议通信与底层API调用,解决设备数据采集、远程监控及定制化HMI开发等实际工程问题。压缩包为ZIP格式,大小26.05MB&#…

作者头像 李华
网站建设 2026/9/10 3:56:37

基于灰狼算法改进扰动观察法的光伏MPPT多峰寻优仿真

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/10 3:53:53

SpringBoot+Vue考务报名系统设计与实现全解析

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华