Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -222,6 +222,9 @@ link_infini_train_exe(llama3)
add_subdirectory(tools/infini_run)
set_target_properties(infini_run PROPERTIES RUNTIME_OUTPUT_DIRECTORY ${CMAKE_BINARY_DIR})

add_executable(pipeline_layout_suggest tools/pipeline_layout_suggest.cc)
target_link_libraries(pipeline_layout_suggest PRIVATE infini_train)

# Tests
if(BUILD_TEST)
add_subdirectory(tests)
Expand Down
164 changes: 164 additions & 0 deletions docs/pipeline_custom_layout_usage.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,164 @@
# Pipeline 并行自定义布局使用说明

本功能把层和特殊模块的归属统一记录在 `PipelineLayout` 中。模型构建、
`PipelineParallel` 包装器、调度器和 GPT-2/LLaMA3 checkpoint loader 都查询同一份
布局,因此不会因为不同组件各自重新计算切分而产生不一致。

## 命令行参数

GPT-2 和 LLaMA3 都支持:

```text
--pipeline_parallel=N
--virtual_pipeline_parallel=N
--pipeline_layer_partition=4,8,6,6
--pipeline_layout='Etttttt|ttttttFH'
--pipeline_embedding_stage=0
--pipeline_final_norm_stage=-1
--pipeline_lm_head_stage=-1
```

`-1` 表示最后一个 pipeline stage。特殊模块参数默认保持旧行为:embedding 在
stage 0,final norm 和 lm head 在最后一个 stage。

不提供 `--pipeline_layer_partition` 时,使用 `BuildDefault` 的均匀布局。chunk 按
`global_chunk_id = local_chunk_id * num_stages + stage_id` 交错编号,以保持现有
GPipe/1F1B 调度编号兼容。

提供 partition 时,每个逗号分隔的整数表示一个 stage 拥有的连续层数。例如
24 层模型使用 `--pipeline_parallel=4 --pipeline_layer_partition=4,8,6,6`:

```text
stage 0: layers [0,4)
stage 1: layers [4,12)
stage 2: layers [12,18)
stage 3: layers [18,24)
```

当前连续 `pipeline_layer_partition` 必须使用 `--virtual_pipeline_parallel=1`。
需要 vPP 的非对称 chunk 映射时使用 Megatron 风格布局字符串,或在 C++ 集成中
调用显式 chunk 构造接口。

### Megatron 风格布局(增强功能)

`--pipeline_layout` 可显式描述每个 stage 的 chunk 和特殊模块:`t` 表示
Transformer 层,`E`/`F`/`H` 分别表示 embedding、final norm、LM head,`|`
分隔 stage,`,` 分隔同一 stage 内的 vPP chunk,括号后使用 `*N` 重复。
例如 `--pipeline_parallel=2 --virtual_pipeline_parallel=2
--pipeline_layout='tt,tt|tt,tt'` 将 8 层映射为 4 个显式 chunk。该参数与
`--pipeline_layer_partition` 互斥(同时传入会在启动期报错);每个 stage 必须恰好给出 `vpp` 个 chunk,
所有 `t` 的总数必须等于模型层数。布局解析和校验在训练启动前完成。

如果需要在 C++ 集成中完全控制每一个 chunk 的 stage 归属,可调用
`PipelineLayout::BuildExplicit(num_layers, pp_size, vpp_size, chunks)`。该接口要求
global chunk id 连续,但以 `ChunkLayout.stage_id` 和 `local_chunk_id` 为权威,不从
global chunk id 反推 stage;因此可以表达非对称的 vPP chunk 大小和显式 ownership。

显式布局也允许空 chunk/空 stage。当前 Pipeline transport 会让空
`TransformerChunk` 原样传递激活和梯度,因此只要 layout 中保留该 stage 的
chunk id,就可以用于真实 PP 运行。例如:

```bash
--pipeline_parallel=3
--pipeline_layout='Etttttt||ttttttFH'
```

### 自动负载均衡建议与性能对比

可使用 `pipeline_layout_suggest` 根据每层代价生成连续 partition:

```bash
./build_cpu/pipeline_layout_suggest \
--num_layers=8 --pp_size=4 --layer_cost=1,1,1,1,4,1,1,1
# pipeline_layer_partition=3,1,1,3
```

不传 `--layer_cost` 时默认每层代价为 1。工具会保证层数总和正确、每个
stage 至少有一层,并倾向于将高代价层单独切分。

可用性能脚本对默认单卡和双卡显式布局做一次可重复的单步对比:

```bash
BUILD_DIR=build_cuda scripts/compare_pipeline_layout_perf.sh gpt2 cuda
```

脚本在 `targets/pipeline_layout_perf.csv` 写出 `elapsed_ms`、`tok_per_s`
和理想 GPipe bubble 比例。`BATCH_SIZE`、`SEQUENCE_LENGTH`、
`TOTAL_BATCH_SIZE` 可通过环境变量覆盖;bubble 使用
`(pipeline_parallel-1)/(micro_batches+pipeline_parallel-1)` 估算,属于
调度上界,不代替真实 profiler。

如需在相同 PP degree 下比较均匀布局和 cost-aware 自定义布局,并记录每个 stage
的真实计算时间,可运行:

```bash
BUILD_DIR=build_cuda \
LAUNCHER='torchrun --standalone --no-python' \
TOTAL_BATCH_SIZE=128 \
scripts/compare_pipeline_stage_perf.sh gpt2 cuda
```

`INFINI_PIPELINE_STAGE_TIMING=1` 会在每个 stage/chunk 的 forward、backward 完成
CUDA stream 同步后输出 `pipeline_stage_timing` 记录;脚本将这些记录汇总到
`targets/pipeline_stage_perf.csv`,并同时写出端到端 elapsed、吞吐、理想 bubble
和 stage imbalance ratio。

均衡建议工具支持直接读取 profiler/用户代价 CSV。文件可以是一列 cost,也可以
是带 `layer,cost` 表头的两列 CSV:

```bash
./build_cuda/pipeline_layout_suggest \
--num_layers=8 --pp_size=4 \
--layer_cost_file=layer_costs.csv
```

输出除了 `pipeline_layer_partition`,还会给出均匀 partition 与建议 partition
的最大 stage 代价及预计降低比例。

## 特殊模块和 checkpoint

布局查询决定 `transformer.wte/wpe`、`transformer.ln_f` 和
`transformer.lm_head` 在哪个 stage 注册和加载。层参数仍使用 canonical key:
`transformer.h.<local_index>...`。checkpoint loader 即使当前 rank 不拥有某个参数,
也会继续 seek 对应字节,避免后续权重错位。

当前 pipeline transport 的既有语义仍按 rank 0 为首 stage、最后 rank 为末 stage。
因此把 final norm/lm head 放到非末 stage 虽可由布局描述,但真实 loss/target 的跨
stage transport 尚未完全重构;部署前应保持它们在最后 stage,或先扩展 transport。

## 错误排查

布局错误使用可捕获的 `PipelineLayoutError`,不会用新的 glog `CHECK` 替代:

```text
stage count mismatch: pipeline_parallel=4, partition_entries=3
layer sum mismatch: num_layers=23, partition_sum=24
pipeline_layer_partition contains invalid character '-'
custom layer partition is not supported with vpp_size=2
```

普通空模块/非法指针参数使用 `std::invalid_argument`;训练代码中原有的
`CHECK`/`LOG(FATAL)` 仍保留其原语义。

## 校验命令

CPU 构建和单测:

```bash
cmake -S . -B build_cpu -DBUILD_TEST=ON -DUSE_CUDA=OFF -DUSE_NCCL=OFF
cmake --build build_cpu -j2
ctest --test-dir build_cpu -R 'test_pipeline_(layout|scheduler_layout|parallel_chunking)_cpu' --output-on-failure
./build_cpu/tests/transformer/test_transformer_cpu --gtest_filter='TransformerPipelineLayoutTest.*' --gtest_color=no
```

有完整数据和 launcher 时,可以运行一迭代 smoke test(这会实际执行训练步骤):

```bash
./scripts/test_pipeline_custom_layout.sh gpt2 cuda
./scripts/test_pipeline_custom_layout.sh llama3 cuda
```

脚本默认查找 `build/`、`data/<model>/` 下的可执行文件和输入;可用
`BUILD_DIR`、`INPUT_BIN`、`INPUT_VAL_BIN`、`TOKENIZER_BIN`、`LLMC_FILE`、
`OUT_DIR`、`LAUNCHER` 覆盖。`LAUNCHER=direct` 可在单进程环境执行,但不能验证
两 stage 的真实通信。没有数据、CUDA 或 launcher 时脚本会明确退出,不会静默通过。
49 changes: 37 additions & 12 deletions example/gpt2/checkpoint_loader.cc
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
#include <tuple>
#include <vector>

#include "gflags/gflags.h"
#include "glog/logging.h"

#include "infini_train/include/nn/modules/normalization.h"
Expand All @@ -28,6 +29,12 @@
using namespace infini_train;
namespace nn = infini_train::nn;

DECLARE_string(pipeline_layer_partition);
DECLARE_string(pipeline_layout);
DECLARE_int32(pipeline_embedding_stage);
DECLARE_int32(pipeline_final_norm_stage);
DECLARE_int32(pipeline_lm_head_stage);

namespace {
constexpr int kRandomSeed = 42;

Expand Down Expand Up @@ -89,6 +96,23 @@ std::shared_ptr<nn::TransformerModel> LoadFromLLMC(const std::string &filepath)
gpt2_config.n_head = n_head;
gpt2_config.n_embd = n_embd;
gpt2::SanitizeGPT2Config(gpt2_config);
nn::parallel::SpecialModulePlacement placement{
.embedding_stage = FLAGS_pipeline_embedding_stage,
.final_norm_stage = FLAGS_pipeline_final_norm_stage,
.lm_head_stage = FLAGS_pipeline_lm_head_stage,
};
auto layout = FLAGS_pipeline_layout.empty()
? nn::parallel::PipelineLayout::BuildPipelineLayout(
static_cast<int>(n_layer),
nn::parallel::global::GetPipelineParallelSize(),
nn::parallel::global::GetVirtualPipelineParallelSize(),
FLAGS_pipeline_layer_partition, placement)
: nn::parallel::PipelineLayout::ParseMegatronStyleLayout(
FLAGS_pipeline_layout, static_cast<int>(n_layer),
nn::parallel::global::GetPipelineParallelSize(),
nn::parallel::global::GetVirtualPipelineParallelSize(), placement);
layout.ValidateForCurrentPipelineTransport();
nn::parallel::global::InstallPipelineLayout(layout);
auto local_gpt2 = std::make_shared<nn::TransformerModel>(gpt2_config);

LOG(INFO) << "magic: " << magic << " version: " << version << " block_size: " << block_size
Expand All @@ -100,15 +124,15 @@ std::shared_ptr<nn::TransformerModel> LoadFromLLMC(const std::string &filepath)
CHECK_EQ(n_head % tp_size, 0) << "n_head must be divisible by TP world size.";

// ========== pp_size:num_stages; vpp_size: num_chunks_per_stage ==========
int pp_size = nn::parallel::global::GetPipelineParallelSize();
int vpp_size = nn::parallel::global::GetVirtualPipelineParallelSize();
auto pp_rank = nn::parallel::pp_rank;
auto [is_first_stage, is_last_stage, layer_ranges_per_chunk]
= nn::parallel::PipelineParallel::GetStageInfo(n_layer, pp_size, pp_rank, vpp_size);
// ========== layer to chunk ==========
const auto &installed_layout = nn::parallel::global::GetPipelineLayout();
const bool owns_embedding = installed_layout.owns(nn::parallel::SpecialModule::kEmbedding, pp_rank);
const bool owns_final_norm = installed_layout.owns(nn::parallel::SpecialModule::kFinalNorm, pp_rank);
const bool owns_lm_head = installed_layout.owns(nn::parallel::SpecialModule::kLMHead, pp_rank);

std::vector<bool> owned_layers(n_layer, false);
for (const auto &[start, end] : layer_ranges_per_chunk) {
for (int i = start; i < end; ++i) { owned_layers[i] = true; }
for (int layer = 0; layer < static_cast<int>(n_layer); ++layer) {
owned_layers[layer] = installed_layout.stage_of_layer(layer) == pp_rank;
}

auto tp_rank = nn::parallel::tp_rank;
Expand All @@ -130,14 +154,15 @@ std::shared_ptr<nn::TransformerModel> LoadFromLLMC(const std::string &filepath)
// transformer.wte.weight (also transformer.lm_head.weight)
// full: (model_vocab_size, n_embd)
// local: (vocab_size_per_partition, n_embd)
if (is_first_stage) {
if (owns_embedding) {
auto &transformer_wte_weight = state_dict[std::format("{}.{}.{}", nn::TransformerModel::kTransformerModelName,
nn::TransformerFirstStage::kWTELayerName,
nn::parallel::VocabParallelEmbedding::kParamWeightName)];
ReadMatrixRowShardFloat(ifs, static_cast<float *>(transformer_wte_weight->DataPtr()), model_vocab_size, n_embd,
v_start, vpp);
} else if (pp_size > 1 && is_last_stage) {
auto &lm_head_weight = state_dict[std::format("{}.{}", nn::TransformerLastStage::kLMHeadLayerName,
} else if (nn::parallel::global::GetPipelineParallelSize() > 1 && owns_lm_head) {
auto &lm_head_weight = state_dict[std::format("{}.{}.{}", nn::TransformerModel::kTransformerModelName,
nn::TransformerLastStage::kLMHeadLayerName,
nn::parallel::ColumnParallelLinear::kParamWeightName)];
ReadMatrixRowShardFloat(ifs, static_cast<float *>(lm_head_weight->DataPtr()), model_vocab_size, n_embd, v_start,
vpp);
Expand All @@ -151,7 +176,7 @@ std::shared_ptr<nn::TransformerModel> LoadFromLLMC(const std::string &filepath)
ifs.ignore((padded_vocab_size - model_vocab_size) * n_embd * sizeof(float));
}

if (is_first_stage) {
if (owns_embedding) {
// transformer.wpe.weight
auto &transformer_wpe_weight
= state_dict[std::format("{}.{}.{}", nn::TransformerModel::kTransformerModelName,
Expand Down Expand Up @@ -407,7 +432,7 @@ std::shared_ptr<nn::TransformerModel> LoadFromLLMC(const std::string &filepath)
}
}

if (is_last_stage) {
if (owns_final_norm) {
// transformer.ln_f.weight
auto &transformer_ln_f_weight
= state_dict[std::format("{}.{}.{}", nn::TransformerModel::kTransformerModelName,
Expand Down
Loading