
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三大组件的协作机制、PPPipeline 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面向用户的高层 APIPPInferEngine承担两类职责初始化 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. MicroBatchManagerMicroBatch 信息管理器它负责跟踪流水线内每个 micro-batch 的运行轨迹记录每个 micro-batch 的信息例如新生成的 tokennew tokens与 KV Cache记录每个 micro-batch 的推理状态例如当前处于 prefill预填充、generate逐 token 生成还是 done完成更新 micro-batch 信息推进生成与序列长度。3. generate schedule流水线推理布局generateschedule 实现流水线推理的具体排布。阶段间通信策略按流水线深度自适应pp_size 22 个流水线阶段使用torch.distributed.P2Pop实现阶段间通信主要用于规避通信竞态race communicationpp_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 / COOLDOWNclass Status(Enum): PREFILL 1 GENERATE 2 DONE 3 COOLDOWN 4状态判断基于cur_length与目标长度的关系MicroBatchDescription.statemicrobatch_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 不同仓库分别实现了两个描述符类描述符适用阶段输入职责HeadMicroBatchDescriptionstage 0头input_idsattention_mask保存原始输入与new_tokens负责逐 token 拼接生成结果并扩展 attention maskBodyMicroBatchDescriptionstage 1..N-1主体hidden_statespast_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), dtypetorch.int64, devicecuda)), dim-1 )MicroBatchManager本身以环形 buffer 管理多个 micro-batchmicrobatch_manager.pymicro_batch_size单个 micro-batch 的样本数micro_batch_buffer_sizemicro-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_size2, modelmodel, model_policyLlamaModelInferPolicy(), new_length32) input [Introduce a landmark in London, Introduce a landmark in Singapore] data tokenizer(input, return_tensorspt) output inferengine.inference(data.to(cuda)) print(tokenizer.batch_decode(output))执行要点拆解colossalai.launch_from_torch()从torch.distributed启动器获取分布式环境多卡场景需配合colossalai run见第六节加载 Hugging Face Llama 权重与 tokenizer路径替换为本地权重目录构造PPInferEnginepp_size2表示切成 2 个流水线阶段需 2 个 GPU / 进程model_policy传入 Llama 的推理切分策略LlamaModelInferPolicy()new_length32表示最多续写 32 个新 tokentokenizer对两个句子批量编码后调用engine.inference(data.to(cuda))返回每个请求生成的 token用tokenizer.batch_decode还原为文本。五、PPInferEngine 关键参数与更多配置项在文档示例的基础上仓库自带的基准脚本 colossalai/legacy/inference/pipeline/benchmark/benchmark.py 展示了更完整的参数组合engine PPInferEngine( pp_sizeargs.pp_size, dtypeargs.dtype, micro_batch_sizeargs.mb_size, new_lengthargs.new_length, modelmodel, model_policyLlamaModelInferPolicy(), verboseTrue, max_batch_sizeargs.mb_size, max_input_lenargs.seq_len, max_output_lenargs.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 数示例中为 32micro_batch_size每个 micro-batch 的样本数与吞吐/显存折中见性能表dtype推理精度fp16/bf16等max_batch_size显存预分配的最大 batch需覆盖实际 batchmax_input_len预分配的最大输入长度应大于等于最长输入序列max_output_len预分配的最大输出长度需要覆盖seq_len new_length基准脚本额外 256 作为冗余verbose是否输出流水线运行明细基准测试用于采集时间戳其中micro_batch_size与 buffer 机制直接决定显存中 KV Cache 的预分配规模与吞吐上限是调优的核心旋钮。六、多卡启动与 Benchmark复现文档性能数字6.1 基准脚本的输入与输出benchmark.py支持三种模型规模toy8 层随机 Llama 配置用于快速功能自检、7b、13b从decapoda-research预训练权重构造。它会在 rank 0 上统计并落盘以下指标benchmark.pyAverage prefill time / Average encode timeAverage micro-batch end2end time / Whole-batch end2end timeMicro-batch / Whole-batchPer Token LatencymsThroughputtokens/sFLOPS按参数量估算同时记录 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 \ --model7b \ --dtypefp16 \ --batch_size${BATCH_SIZE} \ --seq_len1024 \ --new_length128 \ --mb_size$((${BATCH_SIZE}/2)) \ --pp_size2要点通过colossalai run启动分布式任务--nproc_per_node 2与pp_size2对应2 张 GPU 组成 2 级流水线脚本按BATCH_SIZE ∈ {2,4,8,16}循环压测micro-batch 大小取batch_size/2与流水线 stage 数一致正好填满 buffer除seq_len1024 / new_length128外还覆盖seq_len512 / new_length512长生成场景以及 7b/13b 两种规模。七、性能数据Pipeline Inference vs Hugging Face Pipeline文档在2 × A10 20G与2 × A800 80G两种环境下对比了Pipeline Inference与 Hugging Face pipeline 的吞吐tokens/s测试条件为 input length1024、output length128。以下数字均取自原文档作为结果复述。7.1 A10 环境7b / 13bfp16Llama-7bfp16表中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.81OOMHugging Face41.4365.3091.93114.62OOMOOMLlama-13bfp16batch_size(micro_batch size)2(1)4(2)8(4)16(4)Pipeline Inference25.3947.0983.789.46Hugging Face23.4837.5953.44OOM7.2 A800 环境7b / 13bfp16Llama-7bfp16batch_size(micro_batch size)2(1)4(2)8(4)16(8)32(16)Pipeline Inference57.97110.13213.33389.86670.12Hugging Face42.4476.5151.97212.88256.13Llama-13bfp16batch_size(micro_batch size)2(1)4(2)8(4)16(8)32(16)Pipeline Inference41.7894.18172.67310.75470.15Hugging Face36.5768.4105.81139.51166.34从结果可以观察到的规律原文档数据直接反映的事实batch 越大优势越明显在 A10 7b 上 batch 4 时两者持平batch 8/16 后 Pipeline Inference 吞吐提升至 1.5~2 倍这是因为流水线多卡并行分摊了显存并持续吞吐生成 token显存上限被显著推高Hugging Face 在 A10 上 32 batch7b、16 batch13b即 OOM而 Pipeline Inference 可推进到更大 batch 才 OOMA800 大显存下差距更大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/其中描述文件即本文依据的 READMEmicrobatch_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),仅供参考