Preserve backend arguments and checkpoint origin in HF conversion

Materialize M48 G22 from the frozen revision-003 patch and source map 86efe0f5. Preserve the reviewed patch boundary and keep unit-test migrations in the final GU operation.
This commit is contained in:
Tom
2026-10-01 15:43:36 +08:00
parent c374736551
commit ad84212202
2 changed files with 33 additions and 22 deletions
+3 -5
View File
@@ -103,16 +103,13 @@ def get_named_params(args, state_dict):
def save_tensors(args, model_name, state_dict, output_dir, chunk_size, vocab_size=None):
# for miles update_weight compatible
args.sglang_enable_ep_moe = False
print(f"start saving to {output_dir}")
os.makedirs(output_dir, exist_ok=True)
# 2GB
current_size = 0
total_size = 0
modeltensors = [{}]
for name, param in get_named_params(args, state_dict):
for name, param in get_named_params(args.backend, state_dict):
if vocab_size:
param = remove_padding(name, param, vocab_size)
converted_named_tensors = convert_to_hf(args, model_name, name, param)
@@ -209,7 +206,8 @@ if __name__ == "__main__":
)
print(f"model loaded in {time.time()-t:.2f} sec.")
save_tensors(megatron_args, args.model_name, state_dict, args.output_dir, args.chunk_size, args.vocab_size)
conversion_args = argparse.Namespace(backend=megatron_args, hf_checkpoint=args.origin_hf_dir)
save_tensors(conversion_args, args.model_name, state_dict, args.output_dir, args.chunk_size, args.vocab_size)
if args.origin_hf_dir:
copy_assets(args.origin_hf_dir, args.output_dir)
+30 -17
View File
@@ -85,6 +85,7 @@ import socket
import sys
import time
from dataclasses import dataclass
from pathlib import Path
from typing import Any, cast
import ray
@@ -976,7 +977,7 @@ def prepare_moe_block_task_tensors(
def prepare_whole_source_task_tensors(
task: TaskSpec,
input_dir: str,
megatron_args: Any,
conversion_args: argparse.Namespace,
model_name: str,
metadata: dist_cp.metadata.Metadata,
) -> PreparedTaskTensors:
@@ -985,10 +986,10 @@ def prepare_whole_source_task_tensors(
groups: list[PreparedTensorGroup] = []
try:
for name, param in get_named_params(megatron_args, state_dict):
if getattr(megatron_args, "vocab_size", None) is not None:
param = m2hf.remove_padding(name, param, megatron_args.vocab_size)
converted_named_tensors = m2hf._convert_to_hf_core(megatron_args, model_name, name, param)
for name, param in get_named_params(conversion_args.backend, state_dict):
if getattr(conversion_args.backend, "vocab_size", None) is not None:
param = m2hf.remove_padding(name, param, conversion_args.backend.vocab_size)
converted_named_tensors = m2hf._convert_to_hf_core(conversion_args, model_name, name, param)
groups.append(PreparedTensorGroup(name, tuple(converted_named_tensors)))
return PreparedTaskTensors(
tuple(groups),
@@ -1002,7 +1003,7 @@ def write_prepared_tensor_groups(
staging_dir: str,
task_id: int,
groups: tuple[PreparedTensorGroup, ...],
megatron_args: Any,
conversion_args: argparse.Namespace,
quantization_config: dict[str, Any] | None,
max_file_bytes: int,
cuda_device_id: int | None,
@@ -1020,7 +1021,10 @@ def write_prepared_tensor_groups(
torch.cuda.set_device(cuda_device_id)
converted_named_tensors = tuple(
m2hf.quantize_params(
megatron_args, group.source_name, list(converted_named_tensors), quantization_config
conversion_args,
group.source_name,
list(converted_named_tensors),
quantization_config,
)
)
shard_idx, current_size, added_size = append_to_shards(
@@ -1062,7 +1066,7 @@ class ConversionWorker:
actor_id: int,
input_dir: str,
staging_dir: str,
megatron_args: Any,
conversion_args: argparse.Namespace,
model_name: str,
quantization_config: dict[str, Any] | None,
max_file_bytes: int,
@@ -1073,7 +1077,7 @@ class ConversionWorker:
self.cuda_device_id = initialize_worker_cuda_device(actor_id, quantization_config)
self.input_dir = input_dir
self.staging_dir = staging_dir
self.megatron_args = megatron_args
self.conversion_args = conversion_args
self.model_name = model_name
self.quantization_config = quantization_config
self.max_file_bytes = max_file_bytes
@@ -1086,13 +1090,17 @@ class ConversionWorker:
prepared = prepare_moe_block_task_tensors(task, self.input_dir, self.metadata)
else:
prepared = prepare_whole_source_task_tensors(
task, self.input_dir, self.megatron_args, self.model_name, self.metadata
task,
self.input_dir,
self.conversion_args,
self.model_name,
self.metadata,
)
shards, total_size = write_prepared_tensor_groups(
self.staging_dir,
task.task_id,
prepared.groups,
self.megatron_args,
self.conversion_args,
self.quantization_config,
self.max_file_bytes,
self.cuda_device_id,
@@ -1152,7 +1160,7 @@ def collect_ray_results(
tasks: list[TaskSpec],
input_dir: str,
staging_dir: str,
megatron_args: Any,
conversion_args: argparse.Namespace,
model_name: str,
quantization_config: dict[str, Any] | None,
max_file_bytes: int,
@@ -1176,7 +1184,7 @@ def collect_ray_results(
actor_id,
input_dir,
staging_dir,
megatron_args,
conversion_args,
model_name,
quantization_config,
max_file_bytes,
@@ -1270,14 +1278,14 @@ def prepare_output_dir(output_dir: str, force: bool) -> str:
def load_megatron_args(input_dir: str, model_name_override: str | None, vocab_size: int | None) -> tuple[Any, str]:
megatron_args = torch.load(os.path.join(input_dir, "common.pt"), weights_only=False)["args"]
megatron_args = argparse.Namespace(
**vars(torch.load(os.path.join(input_dir, "common.pt"), weights_only=False)["args"])
)
model_name = model_name_override or getattr(megatron_args, "original_hf_model_name", None)
if model_name is None:
raise ValueError("Model name is required when common.pt does not include original_hf_model_name")
if vocab_size is not None:
megatron_args.vocab_size = vocab_size
if not hasattr(megatron_args, "sglang_enable_ep_moe"):
megatron_args.sglang_enable_ep_moe = False
return megatron_args, model_name
@@ -1317,11 +1325,16 @@ def convert_torch_dist_to_hf_ray(args: Args) -> str:
staging_dir = prepare_output_dir(args.output_dir, args.force)
initialize_ray()
metadata_ref = ray.put(metadata)
conversion_args = argparse.Namespace(
backend=megatron_args,
hf_checkpoint=Path(args.origin_hf_dir) if args.origin_hf_dir is not None else None,
extra_high_precision_layers_megatron=getattr(megatron_args, "extra_high_precision_layers_megatron", None),
)
task_results = collect_ray_results(
tasks,
args.input_dir,
staging_dir,
megatron_args,
conversion_args,
model_name,
quantization_config,
args.max_file_bytes,