megatron_model_merger.py — Megatron-LM Checkpoint 合并器¶
源码路径:
verl/model_merger/megatron_model_merger.py
文件概述¶
这个文件实现了 MegatronModelMerger 类,用于将 Megatron-LM 分布式训练产生的 checkpoint 合并为标准的 HuggingFace 格式模型。这是整个 model_merger 模块中最复杂的文件,因为 Megatron-LM 的并行策略比 FSDP 更加多样化。
什么是 Megatron-LM?¶
Megatron-LM 是 NVIDIA 开发的大规模语言模型训练框架,支持三种并行策略: - 张量并行(TP):将每层的参数矩阵切分到多个 GPU - 流水线并行(PP):将模型的不同层分配到不同 GPU - 数据并行(DP):每个 GPU 处理不同的数据
Megatron 的参数命名和 HuggingFace 完全不同,而且某些参数会被融合(如 QKV 三个矩阵融合为一个)。因此合并时不仅要收集分片,还要做参数名映射和张量拆分。
关键代码讲解¶
1. 初始化分布式环境¶
def __init__(self, config: ModelMergerConfig):
super().__init__(config)
if "WORLD_SIZE" not in os.environ:
os.environ["RANK"] = "0"
os.environ["LOCAL_RANK"] = "0"
os.environ["WORLD_SIZE"] = "1"
os.environ["MASTER_ADDR"] = "localhost"
os.environ["MASTER_PORT"] = "12355"
set_numa_affinity()
torch.distributed.init_process_group(get_nccl_backend())
mpu.initialize_model_parallel(
tensor_model_parallel_size=1,
pipeline_model_parallel_size=self.world_size,
virtual_pipeline_model_parallel_size=None,
context_parallel_size=1,
expert_model_parallel_size=1,
)
与 FSDP 不同,Megatron 合并需要初始化 PyTorch 的分布式进程组和 Megatron 的模型并行组。关键点:
- 如果没有设置分布式环境变量,默认单机运行
- TP size 设为 1(合并后不需要张量并行)
- PP size 设为 world_size(多机合并时每个节点负责不同的流水线阶段)
2. 参数名映射表¶
self.params_mapping = {
# Megatron 名称 → HuggingFace 名称
"embedding.word_embeddings": "model.embed_tokens",
# 注意力层
"self_attention.linear_qkv.layer_norm_weight": "input_layernorm.weight",
"self_attention.linear_qkv": "self_attn.qkv_proj",
"self_attention.linear_proj": "self_attn.o_proj",
# MLA (Multi-head Latent Attention, 如 DeepSeek-V2)
"self_attention.linear_q_proj": "self_attn.q_proj",
"self_attention.linear_q_down_proj": "self_attn.q_a_proj",
# MLP 层
"mlp.linear_fc1.layer_norm_weight": "post_attention_layernorm.weight",
"mlp.linear_fc1": "mlp.gate_up_proj",
"mlp.linear_fc2": "mlp.down_proj",
# MoE (Mixture of Experts)
"mlp.router": "mlp.gate",
"mlp.shared_experts.linear_fc1": "mlp.shared_experts.gate_up_proj",
# 输出层
"final_layernorm": "norm",
"output_layer": "lm_head",
}
这是一个核心的名称映射字典。Megatron 和 HuggingFace 对相同层使用完全不同的命名:
- Megatron 的 self_attention.linear_qkv → HuggingFace 的 self_attn.qkv_proj
- Megatron 的 mlp.linear_fc1 → HuggingFace 的 mlp.gate_up_proj
- Megatron 的 output_layer → HuggingFace 的 lm_head
特殊处理:对于 Qwen2MoE 架构,shared experts 的命名有所不同,代码中有针对性的覆盖。
3. 流水线均匀分片算法 get_dynamic_pipeline_shards()¶
def get_dynamic_pipeline_shards(layer_num: int, pp_size: int) -> list[int]:
if pp_size == 1:
return [layer_num]
if pp_size == 2:
return [layer_num // 2, layer_num - layer_num // 2]
middle_size = pp_size - 2
shards_strategy = []
for middle_layer_num in range(layer_num):
first_last_layer_num = layer_num - middle_layer_num * middle_size
first_layer_num = first_last_layer_num // 2
last_layer_num = first_last_layer_num - first_last_layer_num // 2
if 0 < first_layer_num <= middle_layer_num and \
0 < last_layer_num <= middle_layer_num:
shards_strategy.append((
[first_layer_num] + [middle_layer_num] * middle_size + [last_layer_num],
abs(first_layer_num - middle_layer_num),
))
res = sorted(shards_strategy, key=lambda x: x[1])[0][0]
return res
这个函数将模型的层数尽可能均匀地分配到各流水线阶段。例如:
- 32 层模型、4 个 PP rank → [8, 8, 8, 8]
- 25 层模型、4 个 PP rank → [6, 7, 7, 5](尽量均匀)
算法策略:枚举所有可能的分配方案,选择首尾和中间差异最小的那个。
4. 加载 Megatron 分布式 Checkpoint¶
def _load_state_dicts(self, model_ckpt_path):
self.pipeline_shards = get_dynamic_pipeline_shards(
self.hf_config.num_hidden_layers, self.world_size)
tf_config = hf_to_mcore_config(self.hf_config, torch.bfloat16,
num_layers_in_first_pipeline_stage=self.pipeline_shards[0],
num_layers_in_last_pipeline_stage=self.pipeline_shards[-1])
def megatron_model_provider(pre_process, post_process):
from verl.models.mcore import init_mcore_model
parallel_model = init_mcore_model(
tf_config, self.hf_config, pre_process, post_process,
share_embeddings_and_output_weights=tie_word_embeddings,
value=False)
return parallel_model
with context():
whole_model = get_model(
model_provider_func=megatron_model_provider,
model_type=ModelType.encoder_or_decoder,
wrap_with_ddp=False,
transformer_config=tf_config)
# 加载 dist checkpoint
sharded_state_dict = {}
for vpp_rank, model in enumerate(whole_model):
key = f"model{vpp_rank}" if len(whole_model) > 1 else "model"
sharded_state_dict[key] = model.sharded_state_dict()
model_state_dict = load_dist_checkpointing(sharded_state_dict, model_ckpt_path)
Megatron 的 checkpoint 加载比 FSDP 复杂得多:
1. 先根据 HuggingFace config 构建 Megatron 模型配置
2. 构建一个空的 Megatron 模型
3. 从模型获取 sharded_state_dict(描述参数的分片元信息)
4. 用 Megatron 的 load_dist_checkpointing 加载实际参数
5. 张量拆分 _split_tensors()¶
def _split_tensors(self, key, tensor, config, is_value_model=False):
if "linear_fc1.weight" in key:
# gate_up 融合张量 → 拆分为 gate 和 up
gate, up = tensor.chunk(2)
return [gate, up]
elif "self_attention.linear_qkv." in key and "layer_norm" not in key:
# QKV 融合张量 → 拆分为 q, k, v
num_q_per_kv = config.num_attention_heads // config.num_key_value_heads
kv_size = tensor.shape[0] // (num_q_per_kv + 2)
q_lst, k_lst, v_lst = [], [], []
for chunk in tensor.chunk(config.num_key_value_heads):
split_size = [
kv_size * num_q_per_kv // config.num_key_value_heads,
kv_size // config.num_key_value_heads,
kv_size // config.num_key_value_heads,
]
q, k, v = chunk.split(split_size)
q_lst.append(q)
k_lst.append(k)
v_lst.append(v)
return [torch.cat(q_lst), torch.cat(k_lst), torch.cat(v_lst)]
else:
return [tensor]
Megatron 为了计算效率,会将某些参数融合存储: - gate_up 融合:MLP 的 gate projection 和 up projection 合并为一个矩阵 - QKV 融合:注意力层的 Q、K、V 投影合并为一个矩阵
HuggingFace 模型中这些是分开的,所以合并时需要拆分。QKV 的拆分需要考虑 GQA(Grouped Query Attention)中 Q 和 KV 的数量比例。
6. 参数名转换 _replace_name()¶
def _replace_name(self, megatron_name, name_mapping):
for m_name, v_name in name_mapping.items():
if m_name not in megatron_name:
continue
megatron_name = megatron_name.replace("decoder", "model")
param_name = megatron_name.replace(m_name, v_name)
return param_name
return None
将 Megatron 的参数名转换为 HuggingFace 格式。例如:
- decoder.layers.0.self_attention.linear_qkv.weight → model.layers.0.self_attn.qkv_proj.weight
7. 合并所有 state_dict _merge_state_dicts()¶
def _merge_state_dicts(self, model_state_dict_list):
state_dict = {}
layers_cum = 0
# 处理流水线并行的层号偏移
if self.world_size > 1:
pipeline_cumsum = np.cumsum(self.pipeline_shards)
layers_cum = 0 if self.rank == 0 else pipeline_cumsum[self.rank - 1]
for model_state_dict in model_state_dict_list:
for key in keys:
if "extra_state" in key:
continue
hf_name = self._replace_name(key, self.params_mapping)
# 修正层号(加上流水线偏移)
if "model.layers." in hf_name:
local_layer_no = int(hf_name.split(".")[2])
global_layer_no = local_layer_no + layers_cum
# ... 更新 hf_name 中的层号
tensor = model_state_dict[key]
split_tensor = self._split_tensors(key, tensor, self.hf_config)
if len(split_tensor) == 1:
state_dict[hf_name] = split_tensor[0]
elif len(split_tensor) == 3:
# QKV 拆分
for n, d in zip(["q", "k", "v"], split_tensor):
state_dict[hf_name.replace("qkv", n)] = d
elif len(split_tensor) == 2:
# gate_up 拆分
state_dict[hf_name.replace("gate_up", "gate")] = split_tensor[0]
state_dict[hf_name.replace("gate_up", "up")] = split_tensor[1]
这个方法的核心逻辑是遍历 Megatron 的 state_dict,对每个参数:
1. 跳过 extra_state(Megatron 的内部状态)
2. 将 Megatron 参数名映射为 HuggingFace 参数名
3. 修正流水线并行造成的层号偏移
4. 拆分融合张量(QKV / gate_up)
5. 处理 MoE 的 expert 索引重排
8. 多进程分布式保存¶
def save_hf_model_and_tokenizer(self, merged_state_dict):
if self.world_size == 1:
return super().save_hf_model_and_tokenizer(merged_state_dict)
# 多进程保存:每个 rank 保存自己负责的层
layer_this_rank = self.pipeline_shards[self.rank]
keys_chunk = np.array_split(np.array(keys), layer_this_rank * saves_per_layer)
for i, keys in enumerate(keys_chunk):
sd_to_save = {k: merged_state_dict[k] for k in keys}
save_path = target_dir / f"model-{save_idx + 1:05d}-of-{saves_total:05d}.safetensors"
save_file(sd_to_save, save_path)
# 汇总所有 rank 的索引信息
dist.all_gather_object(all_save_indexes, saves_indexes)
if self.rank == 0:
# rank 0 写 model.safetensors.index.json
对于超大模型(如 671B),支持多节点分布式保存。每个节点只保存自己负责的层的参数,最后由 rank 0 生成索引文件。
核心流程图¶
__init__()
│ 初始化分布式环境 + Megatron 并行组
│ 构建参数名映射表
▼
merge_and_save()
│
▼
_load_state_dicts()
│ 1. 计算流水线分片策略
│ 2. 构建空 Megatron 模型
│ 3. 加载分布式 checkpoint
▼
_merge_state_dicts()
│ 1. 参数名映射 (Megatron → HuggingFace)
│ 2. 修正流水线层号偏移
│ 3. 拆分 QKV / gate_up 融合张量
│ 4. 处理 MoE expert 索引
▼
save_hf_model_and_tokenizer()
│ 单进程:调用基类方法
│ 多进程:分布式保存 + 生成索引
▼
cleanup()
│ 销毁分布式进程组
核心类/函数列表¶
| 名称 | 类型 | 说明 |
|---|---|---|
get_dynamic_pipeline_shards() |
函数 | 均匀分配流水线各阶段的层数 |
noop_context() |
函数 | 空上下文管理器 |
MegatronModelMerger |
类 | Megatron checkpoint 合并器 |
_load_state_dicts() |
方法 | 构建模型并加载分布式 checkpoint |
_check_megatron_state_key() |
方法 | 验证参数名格式 |
_split_tensors() |
方法 | 拆分 QKV / gate_up 融合张量 |
_merge_state_dicts() |
方法 | 名称映射 + 张量拆分 + 层号修正 |
_replace_name() |
方法 | Megatron → HuggingFace 参数名转换 |
_validate_state_dict() |
方法 | 对比验证合并结果 |
save_hf_model_and_tokenizer() |
方法 | 支持单/多进程保存 |
与其他模块的关系¶
- 继承自
BaseModelMerger - verl/models/mcore/ - 使用
hf_to_mcore_config和init_mcore_model构建 Megatron 模型 - verl/utils/megatron/ - 使用
load_dist_checkpointing加载分布式 checkpoint - megatron.core - 使用 Megatron Core 的模型并行初始化和模型定义
- safetensors - 使用 safetensors 格式保存大模型
小结¶
MegatronModelMerger 是整个 model_merger 模块中最复杂的部分,其难点在于:
1. Megatron 使用独特的参数命名(与 HuggingFace 完全不同),需要详细的名称映射
2. Megatron 会融合 QKV 和 gate_up 张量以提升效率,合并时需要正确拆分
3. 流水线并行需要处理层号偏移
4. 支持多节点分布式合并以处理超大模型
5. 需要支持多种模型架构(标准 Transformer、MLA、MoE)
理解这个文件有助于深入理解分布式训练中模型参数的组织方式和不同框架之间的差异。