inference/convert.py
| 1 | import os |
| 2 | import shutil |
| 3 | from argparse import ArgumentParser |
| 4 | from glob import glob |
| 5 | from tqdm import tqdm, trange |
| 6 | |
| 7 | import torch |
| 8 | from safetensors.torch import safe_open, save_file |
| 9 | |
| 10 | |
| 11 | FP4_TABLE = torch.tensor([ |
| 12 | 0.0, 0.5, 1.0, 1.5, 2.0, 3.0, 4.0, 6.0, |
| 13 | 0.0, -0.5, -1.0, -1.5, -2.0, -3.0, -4.0, -6.0 |
| 14 | ], dtype=torch.float32) |
| 15 | |
| 16 | |
| 17 | def cast_e2m1fn_to_e4m3fn(x: torch.Tensor, scale: torch.Tensor) -> tuple[torch.Tensor, torch.Tensor]: |
| 18 | """ |
| 19 | Casts a tensor from e2m1fn to e4m3fn losslessly. |
| 20 | """ |
| 21 | assert x.dtype == torch.int8 |
| 22 | assert x.ndim == 2 |
| 23 | out_dim, in_dim = x.size() |
| 24 | in_dim *= 2 |
| 25 | fp8_block_size = 128 |
| 26 | fp4_block_size = 32 |
| 27 | assert in_dim % fp8_block_size == 0 and out_dim % fp8_block_size == 0 |
| 28 | assert scale.size(0) == out_dim and scale.size(1) == in_dim // fp4_block_size |
| 29 | |
| 30 | x = x.view(torch.uint8) |
| 31 | low = x & 0x0F |
| 32 | high = (x >> 4) & 0x0F |
| 33 | x = torch.stack([FP4_TABLE[low.long()], FP4_TABLE[high.long()]], dim=-1).flatten(2) |
| 34 | |
| 35 | # max_fp4 (6.0) * MAX_OFFSET must fit in e4m3fn (max 448) |
| 36 | # 6.0 * 2^6 = 384 < 448; 6.0 * 2^7 = 768 > 448; so MAX_OFFSET_BITS = 6 |
| 37 | MAX_OFFSET_BITS = 6 |
| 38 | |
| 39 | bOut = out_dim // fp8_block_size |
| 40 | bIn = in_dim // fp8_block_size |
| 41 | # bOut, bIn, 128, 128 |
| 42 | x = x.view(bOut, fp8_block_size, bIn, fp8_block_size).transpose(1, 2) |
| 43 | # bOut, bIn, 128*4 |
| 44 | scale = scale.float().view(bOut, fp8_block_size, bIn, -1).transpose(1, 2).flatten(2) |
| 45 | ## bOut, bIn, 1 |
| 46 | scale_max_offset_bits = scale.amax(dim=-1, keepdim=True) / (2**MAX_OFFSET_BITS) |
| 47 | # bOut, bIn, 128*4 |
| 48 | offset = scale / scale_max_offset_bits |
| 49 | # bOut, bIn, 128, 128 |
| 50 | offset = offset.unflatten(-1, (fp8_block_size, -1)).repeat_interleave(fp4_block_size, dim=-1) |
| 51 | x = (x * offset).transpose(1, 2).reshape(out_dim, in_dim) |
| 52 | return x.to(torch.float8_e4m3fn), scale_max_offset_bits.squeeze(-1).to(torch.float8_e8m0fnu) |
| 53 | |
| 54 | |
| 55 | mapping = { |
| 56 | "embed": ("embed", 0), |
| 57 | "wq_b": ("wq_b", 0), |
| 58 | "wo_a": ("wo_a", 0), |
| 59 | "wo_b": ("wo_b", 1), |
| 60 | "head": ("head", 0), |
| 61 | "attn_sink": ("attn_sink", 0), |
| 62 | "weights_proj": ("weights_proj", 0), |
| 63 | "markov_w1": ("markov_w1", 0), |
| 64 | "markov_w2": ("markov_w2", 0), |
| 65 | } |
| 66 | |
| 67 | |
| 68 | def main(hf_ckpt_path, save_path, n_experts, mp, expert_dtype): |
| 69 | """ |
| 70 | Converts and saves model checkpoint files into a specified format. |
| 71 | |
| 72 | Args: |
| 73 | hf_ckpt_path (str): Path to the directory containing the input checkpoint files. |
| 74 | save_path (str): Path to the directory where the converted checkpoint files will be saved. |
| 75 | n_experts (int): Total number of experts in the model. |
| 76 | mp (int): Model parallelism factor. |
| 77 | |
| 78 | Returns: |
| 79 | None |
| 80 | """ |
| 81 | torch.set_num_threads(8) |
| 82 | n_local_experts = n_experts // mp |
| 83 | state_dicts = [{} for _ in range(mp)] |
| 84 | |
| 85 | for file_path in tqdm(glob(os.path.join(hf_ckpt_path, "*.safetensors"))): |
| 86 | with safe_open(file_path, framework="pt", device="cpu") as f: |
| 87 | for name in f.keys(): |
| 88 | param: torch.Tensor = f.get_tensor(name) |
| 89 | if name.startswith("model."): |
| 90 | name = name[len("model."):] |
| 91 | if name.startswith("mtp.") and ("emb" in name or name.endswith("head.weight")): |
| 92 | continue |
| 93 | name = name.replace("self_attn", "attn") |
| 94 | name = name.replace("mlp", "ffn") |
| 95 | name = name.replace("weight_scale_inv", "scale") |
| 96 | name = name.replace("e_score_correction_bias", "bias") |
| 97 | if any(x in name for x in ["hc", "attn_sink", "tie2eid", "ape"]): # without .weight |
| 98 | key = name.split(".")[-1] |
| 99 | else: |
| 100 | key = name.split(".")[-2] |
| 101 | if key in mapping: |
| 102 | new_key, dim = mapping[key] |
| 103 | else: |
| 104 | new_key, dim = key, None |
| 105 | name = name.replace(key, new_key) |
| 106 | for i in range(mp): |
| 107 | new_param = param |
| 108 | if "experts" in name and "shared_experts" not in name: |
| 109 | idx = int(name.split(".")[-3]) |
| 110 | if idx < i * n_local_experts or idx >= (i + 1) * n_local_experts: |
| 111 | continue |
| 112 | elif dim is not None: |
| 113 | assert param.size(dim) % mp == 0, f"Dimension {dim} must be divisible by {mp}" |
| 114 | shard_size = param.size(dim) // mp |
| 115 | new_param = param.narrow(dim, i * shard_size, shard_size).contiguous() |
| 116 | state_dicts[i][name] = new_param |
| 117 | |
| 118 | os.makedirs(save_path, exist_ok=True) |
| 119 | |
| 120 | for i in trange(mp): |
| 121 | names = list(state_dicts[i].keys()) |
| 122 | for name in names: |
| 123 | if name.endswith("wo_a.weight"): |
| 124 | weight = state_dicts[i][name] |
| 125 | scale = state_dicts[i].pop(name.replace("weight", "scale")) |
| 126 | weight = weight.unflatten(0, (-1, 128)).unflatten(-1, (-1, 128)).float() * scale[:, None, :, None].float() |
| 127 | state_dicts[i][name] = weight.flatten(2, 3).flatten(0, 1).bfloat16() |
| 128 | elif "experts" in name and state_dicts[i][name].dtype == torch.int8: |
| 129 | if expert_dtype == "fp8": |
| 130 | scale_name = name.replace("weight", "scale") |
| 131 | weight = state_dicts[i].pop(name) |
| 132 | scale = state_dicts[i].pop(scale_name) |
| 133 | state_dicts[i][name], state_dicts[i][scale_name] = cast_e2m1fn_to_e4m3fn(weight, scale) |
| 134 | else: |
| 135 | state_dicts[i][name] = state_dicts[i][name].view(torch.float4_e2m1fn_x2) |
| 136 | save_file(state_dicts[i], os.path.join(save_path, f"model{i}-mp{mp}.safetensors")) |
| 137 | |
| 138 | for file in ["tokenizer.json", "tokenizer_config.json"]: |
| 139 | old_file_path = os.path.join(hf_ckpt_path, file) |
| 140 | new_file_path = os.path.join(save_path, file) |
| 141 | if os.path.exists(old_file_path): |
| 142 | shutil.copyfile(old_file_path, new_file_path) |
| 143 | |
| 144 | |
| 145 | if __name__ == "__main__": |
| 146 | parser = ArgumentParser() |
| 147 | parser.add_argument("--hf-ckpt-path", type=str, required=True) |
| 148 | parser.add_argument("--save-path", type=str, required=True) |
| 149 | parser.add_argument("--n-experts", type=int, required=True) |
| 150 | parser.add_argument("--model-parallel", type=int, required=True) |
| 151 | parser.add_argument("--expert-dtype", type=str, choices=["fp8", "fp4"], required=False, default=None) |
| 152 | args = parser.parse_args() |
| 153 | assert args.n_experts % args.model_parallel == 0, "Number of experts must be divisible by model parallelism" |
| 154 | main(args.hf_ckpt_path, args.save_path, args.n_experts, args.model_parallel, args.expert_dtype) |
| 155 | |