From 0ab7c2f3615f30e8023b00e3b368aa5d4d7a9a48 Mon Sep 17 00:00:00 2001 From: addsubmuldiv Date: Tue, 8 Sep 2026 12:20:03 +0000 Subject: [PATCH] adapt_megatron_018 --- src/mcore_bridge/config/model_config.py | 41 +---------- src/mcore_bridge/tuners/lora.py | 94 ++++++------------------- src/mcore_bridge/tuners/npu_lora.py | 18 ++--- 3 files changed, 31 insertions(+), 122 deletions(-) diff --git a/src/mcore_bridge/config/model_config.py b/src/mcore_bridge/config/model_config.py index c5231404..315a221f 100644 --- a/src/mcore_bridge/config/model_config.py +++ b/src/mcore_bridge/config/model_config.py @@ -269,47 +269,8 @@ class ModelConfig(TransformerConfig): num_labels: Optional[int] = None mlp_padding_free: bool = False - _mindspeed_defaults_cache = None - - def _augment_mindspeed_defaults(self): - if not is_torch_npu_available(): - return - - if ModelConfig._mindspeed_defaults_cache is None: - defaults = {} - try: - import mindspeed.features_manager as mfm - import sys - from argparse import ArgumentParser - from mindspeed.arguments import process_args - - original_features = list(mfm.FEATURES_LIST) - full_features = mfm.create_features_list() - mfm.FEATURES_LIST.clear() - mfm.FEATURES_LIST.extend(full_features) - try: - parser = ArgumentParser() - process_args(parser) - # Parse args from sys.argv - args, _ = parser.parse_known_args([]) - defaults = vars(args) - finally: - mfm.FEATURES_LIST.clear() - mfm.FEATURES_LIST.extend(original_features) - except Exception as e: - logger.warning(f'Failed to get MindSpeed defaults, which may cause issues on NPU: {e}') - defaults = {} - ModelConfig._mindspeed_defaults_cache = defaults - - for name, value in ModelConfig._mindspeed_defaults_cache.items(): - if not hasattr(self, name): - setattr(self, name, value) - elif hasattr(self, name) and getattr(self, name) is None and value is not None: - setattr(self, name, value) - def __post_init__(self): from mcore_bridge.model import get_mcore_model_type, get_model_meta - self._augment_mindspeed_defaults() self._format_config() if self.experimental_attention_variant is not None: require_version('megatron-core>=0.16.0.dev', @@ -402,7 +363,7 @@ def _check_npu(self): required_ep = (num_experts + MAX_NPU_EXPERTS_PER_EP - 1) // MAX_NPU_EXPERTS_PER_EP if expert_model_parallel_size < required_ep: logger.warning(f'{">" * 20} WARNING {"<" * 20}\n' - f'MindSpeed on NPU supports up to {MAX_NPU_EXPERTS_PER_EP} experts per EP group. ' + f'NPU grouped matmul supports up to {MAX_NPU_EXPERTS_PER_EP} experts per EP group. ' f'num_experts={num_experts}, ' f'expert_model_parallel_size={expert_model_parallel_size}. ' f'Please set expert_model_parallel_size (EP) to {required_ep} ' diff --git a/src/mcore_bridge/tuners/lora.py b/src/mcore_bridge/tuners/lora.py index 8b150936..374988dc 100644 --- a/src/mcore_bridge/tuners/lora.py +++ b/src/mcore_bridge/tuners/lora.py @@ -7,7 +7,6 @@ import torch.nn.functional as F import warnings from contextlib import contextmanager, nullcontext -from importlib import metadata from megatron.core import parallel_state from megatron.core.dist_checkpointing.mapping import ShardedStateDict from megatron.core.extensions.transformer_engine import (TEColumnParallelGroupedLinear, TEColumnParallelLinear, @@ -33,36 +32,9 @@ mcore_016 = version.parse(megatron.core.__version__) >= version.parse('0.16.0rc0') peft_019 = version.parse(peft.__version__) >= version.parse('0.19.0') -MINDSPEED_015 = version.parse('0.15.0') - - -def _get_mindspeed_version(): - try: - return version.parse(metadata.version('mindspeed')) - except metadata.PackageNotFoundError: - return None - except Exception: - return None - - -def _use_legacy_npu_local_linear() -> bool: - if not is_torch_npu_available(): - return False - mindspeed_version = _get_mindspeed_version() - if mindspeed_version is None: - # Fall back to the conservative path when the version is unknown so we - # do not force an older NPU stack onto the 0.15 TE semantics. - return True - return mindspeed_version < MINDSPEED_015 def _build_local_te_linear(input_size: int, output_size: int, bias: bool, **kwargs): - if _use_legacy_npu_local_linear(): - return nn.Linear( - in_features=input_size, - out_features=output_size, - bias=bias, - ) local_kwargs = dict(kwargs) local_kwargs.pop('tp_group', None) return TELinear( @@ -76,19 +48,27 @@ def _build_local_te_linear(input_size: int, output_size: int, bias: bool, **kwar def _get_tensor_parallel_group_for_lora(base_layer): - """Resolve the tensor-parallel group across TE and MindSpeed TE variants. - - Megatron's TE layers expose ``tp_group`` directly, but MindSpeed 0.15.x - replaces some TE classes (for example - ``MindSpeedTELayerNormColumnParallelLinear``) with implementations that keep - the same tensor-parallel semantics under ``parallel_group`` instead. LoRA - still needs to forward the right group into the newly created parallel - adapter layers, otherwise adapter injection fails before training starts. - """ + """Resolve the tensor-parallel group from the standard TE/MCore contract.""" tp_group = getattr(base_layer, 'tp_group', None) if tp_group is not None: return tp_group - return getattr(base_layer, 'parallel_group', None) + return getattr(base_layer, '_tp_group', None) + + +def _forward_npu_layernorm_column(base_layer, x, args, kwargs): + """Run TENPU LayerNormLinear once and retain its normalized activation.""" + original_return_layernorm_output = base_layer.return_layernorm_output + base_layer.return_layernorm_output = True + try: + out = base_layer(x, *args, **kwargs) + finally: + base_layer.return_layernorm_output = original_return_layernorm_output + + if base_layer.te_return_bias: + result, bias, layernorm_output = out + else: + (result, layernorm_output), bias = out + return result, bias, layernorm_output class LoraParallelLinear(MegatronModule, LoraLayer): @@ -224,10 +204,8 @@ def update_layer(self, adapter_name, r, *, lora_alpha, **kwargs): lora_b = _build_local_te_linear(r, self.out_features, lora_bias, **kwargs) lora_a.parallel_mode = self.base_layer.parallel_mode # fix moe_shared_expert_overlap else: - if is_torch_npu_available(): - out_features = self.out_features - else: - out_features = self.out_features * self.tp_size + # PEFT reports local features; MCore's constructor expects the global size. + out_features = self.out_features * self.tp_size if self.is_grouped: if is_torch_npu_available(): lora_a = NpuGroupedLoraLinear( @@ -380,37 +358,11 @@ def forward(self, x: torch.Tensor, *args: Any, **kwargs: Any): self.base_layer.return_layernorm_output = False result, bias = self.base_layer(x, *args, **kwargs) else: - self.base_layer.return_layernorm_output = True if is_torch_npu_available(): - # NPU: base_layer only returns (output, bias); it does not expose the LayerNorm/RMSNorm output. - inp = x # Keep the original pre-norm input. - result, bias = self.base_layer(inp, *args, **kwargs) - - # Key: For LoRA we need the same "x" as in the non-NPU branch, i.e. the post-norm activation - # (LayerNorm/RMSNorm output, which is the actual input to the fused linear). - if hasattr(self.base_layer, 'config') and (hasattr(self.base_layer, '_layernorm') - or hasattr(self.base_layer, '_rmsnorm')): - norm_type = getattr(self.base_layer.config, 'normalization', None) - - if norm_type == 'LayerNorm': - if not hasattr(self.base_layer, '_layernorm'): - raise RuntimeError( - 'NPU LoRA path expects base_layer to provide `_layernorm`, but it is missing. ' - 'Cannot reconstruct the post-LayerNorm activation for LoRA.') - x = self.base_layer._layernorm(inp) - else: - # Default to RMSNorm path when normalization is not LayerNorm. - if not hasattr(self.base_layer, '_rmsnorm'): - raise RuntimeError( - 'NPU LoRA path expects base_layer to provide `_rmsnorm`, but it is missing. ' - 'Cannot reconstruct the post-RMSNorm activation for LoRA.') - x = self.base_layer._rmsnorm(inp) - else: - raise RuntimeError('NPU LoRA path requires base_layer to expose post-norm activations ' - '(LayerNorm/RMSNorm output). Expected base_layer to have `config` ' - 'and either `_layernorm` or `_rmsnorm`. ' - f'Got base_layer type: {type(self.base_layer)}. ') + result, bias, x = _forward_npu_layernorm_column( + self.base_layer, x, args, kwargs) else: + self.base_layer.return_layernorm_output = True (result, x), bias = self.base_layer(x, *args, **kwargs) elif isinstance(self.base_layer, (TELinear, TEGroupedLinear)): result, bias = self.base_layer(x, *args, **kwargs) diff --git a/src/mcore_bridge/tuners/npu_lora.py b/src/mcore_bridge/tuners/npu_lora.py index 099b6a83..9a3431b0 100644 --- a/src/mcore_bridge/tuners/npu_lora.py +++ b/src/mcore_bridge/tuners/npu_lora.py @@ -18,12 +18,6 @@ _GMM_GROUP_LIST_IS_EXPERT_SIZES = 1 -def _is_mindspeed_grouped_linear(base_layer) -> bool: - if not (is_torch_npu_available() and isinstance(base_layer, TEGroupedLinear)): - return False - return type(base_layer).__module__.startswith('mindspeed.') - - def _has_moe_local_expert_grouping(base_layer) -> bool: config = getattr(base_layer, 'config', None) num_moe_experts = getattr(config, 'num_moe_experts', None) @@ -41,16 +35,16 @@ def is_expert_layer(base_layer) -> bool: is_expert = getattr(base_layer, 'is_expert', None) if is_expert is not None: return bool(is_expert) - if _is_mindspeed_grouped_linear(base_layer): + if is_torch_npu_available() and isinstance(base_layer, TEGroupedLinear): if getattr(base_layer, 'explicit_expert_comm', False): return True has_moe_local_expert_grouping = _has_moe_local_expert_grouping(base_layer) if getattr(base_layer, 'expert_parallel', False): return has_moe_local_expert_grouping - # MindSpeedTEGroupedLinear receives is_expert but does not keep it as - # an attribute. When EP/ETP does not trigger explicit expert comm, - # fall back to the TEGroupedMLP invariant: one grouped slot per local - # expert. + # MCore's TEGroupedLinear requires is_expert=True at construction but + # does not retain that argument. With EP=ETP=1 there is no explicit + # expert communication, so use the TEGroupedMLP invariant: one grouped + # GEMM slot per local expert. return has_moe_local_expert_grouping return False @@ -137,6 +131,8 @@ def sharded_state_dict( param, f'{key_prefix}{param_name}', prepend_offsets=new_sharded_offsets, + tp_group=parallel_state.get_expert_tensor_parallel_group(), + dp_cp_group=parallel_state.get_expert_data_parallel_group(), ) sharded_state_dict[f'{prefix}{local_name}'] = self._set_expert_replica_id(sharded_tensor)