diff --git a/examples/academic_paper_scripts/detxoify_lm/generate_samples_gpt.py b/examples/academic_paper_scripts/detxoify_lm/generate_samples_gpt.py index 2a2b1d63a21..f9fa4b2096b 100644 --- a/examples/academic_paper_scripts/detxoify_lm/generate_samples_gpt.py +++ b/examples/academic_paper_scripts/detxoify_lm/generate_samples_gpt.py @@ -13,6 +13,7 @@ from megatron.training import get_tokenizer from megatron.training import print_rank_0 from megatron.training.checkpointing import load_checkpoint +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core import mpu from megatron.training.arguments import parse_and_validate_args from megatron.training.initialize import initialize_megatron @@ -74,7 +75,9 @@ def model_provider(pre_process=True, post_process=True) -> GPTModel: share_embeddings_and_output_weights=not args.untie_embeddings_and_output_weights, position_embedding_type=args.position_embedding_type, rotary_percent=args.rotary_percent - ) + , + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) return model diff --git a/examples/bert/pretrain_bert.py b/examples/bert/pretrain_bert.py index 9bb3e653e22..7a10dd3d1f1 100644 --- a/examples/bert/pretrain_bert.py +++ b/examples/bert/pretrain_bert.py @@ -10,6 +10,7 @@ from megatron.training import get_args from megatron.training import print_rank_0 from megatron.training import get_timers +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core import tensor_parallel from megatron.core.enums import ModelType from megatron.core.models.bert.bert_model import BertModel @@ -55,7 +56,9 @@ def model_provider(pre_process=True, post_process=True, vp_stage=None, config=No parallel_output=True, pre_process=pre_process, post_process=post_process, - vp_stage=vp_stage) + vp_stage=vp_stage, + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) return model diff --git a/examples/export/trtllm_export/distributed_export/gpt_distributed_gpu_export.py b/examples/export/trtllm_export/distributed_export/gpt_distributed_gpu_export.py index 57d44f9f628..d667f576e5f 100644 --- a/examples/export/trtllm_export/distributed_export/gpt_distributed_gpu_export.py +++ b/examples/export/trtllm_export/distributed_export/gpt_distributed_gpu_export.py @@ -1,5 +1,6 @@ import os import torch +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core import parallel_state from megatron.core import dist_checkpointing from megatron.core.export.model_type import ModelType @@ -42,7 +43,9 @@ def model_provider(): transformer_layer_spec=get_gpt_layer_local_spec(), vocab_size=_VOCAB_SIZE, max_sequence_length=_SEQUENCE_LENGTH, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) return gpt_model diff --git a/examples/export/trtllm_export/single_device_export/gpt_single_device_cpu_export.py b/examples/export/trtllm_export/single_device_export/gpt_single_device_cpu_export.py index 587e7cfdd32..e647ac1f256 100644 --- a/examples/export/trtllm_export/single_device_export/gpt_single_device_cpu_export.py +++ b/examples/export/trtllm_export/single_device_export/gpt_single_device_cpu_export.py @@ -1,5 +1,6 @@ import os import torch +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core import parallel_state from megatron.core import dist_checkpointing from megatron.core.export.model_type import ModelType @@ -43,7 +44,9 @@ def model_provider(): transformer_layer_spec=get_gpt_layer_local_spec(), vocab_size=100, max_sequence_length=_SEQUENCE_LENGTH, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) return gpt_model diff --git a/examples/run_simple_mcore_train_loop.py b/examples/run_simple_mcore_train_loop.py index 1ba5e10cfc9..f88efbab59c 100644 --- a/examples/run_simple_mcore_train_loop.py +++ b/examples/run_simple_mcore_train_loop.py @@ -7,6 +7,7 @@ from functools import partial from pathlib import Path from typing import Any, Callable, Dict, Tuple, Iterator +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core import parallel_state from megatron.core import dist_checkpointing from megatron.core.pipeline_parallel.schedules import get_forward_backward_func @@ -74,7 +75,9 @@ def model_provider() -> GPTModel: transformer_layer_spec=get_gpt_layer_local_spec(), vocab_size=100, max_sequence_length=_SEQUENCE_LENGTH, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) return gpt_model diff --git a/examples/t5/pretrain_t5.py b/examples/t5/pretrain_t5.py index 4b33386e2d4..93bb0de4d74 100644 --- a/examples/t5/pretrain_t5.py +++ b/examples/t5/pretrain_t5.py @@ -9,6 +9,7 @@ import torch import megatron +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core import mpu, tensor_parallel from megatron.core.tokenizers.utils.build_tokenizer import build_tokenizer from megatron.core.datasets.blended_megatron_dataset_builder import BlendedMegatronDatasetBuilder @@ -132,7 +133,9 @@ def model_provider( relative_attention_max_distance=args.relative_attention_max_distance, add_encoder=add_encoder, add_decoder=add_decoder, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) return model diff --git a/megatron/core/models/T5/t5_model.py b/megatron/core/models/T5/t5_model.py index b2feb974643..d7677be9712 100644 --- a/megatron/core/models/T5/t5_model.py +++ b/megatron/core/models/T5/t5_model.py @@ -158,7 +158,7 @@ def __init__( pg_collection: ProcessGroupCollection = None, ): - super(T5Model, self).__init__(config=config) + super(T5Model, self).__init__(config=config, pg_collection=pg_collection) self.config: TransformerConfig = config self.encoder_config: TransformerConfig = encoder_config diff --git a/megatron/core/models/common/language_module/language_module.py b/megatron/core/models/common/language_module/language_module.py index 53522dd8b2b..dcc09bc4009 100644 --- a/megatron/core/models/common/language_module/language_module.py +++ b/megatron/core/models/common/language_module/language_module.py @@ -47,8 +47,12 @@ def __init__( ) -> None: super().__init__(config=config) self._set_attention_backend() - if pg_collection is None: - pg_collection = ProcessGroupCollection.use_mpu_process_groups() + assert pg_collection is not None, ( + "LanguageModule requires an explicit pg_collection. The global parallel grid is not a " + "safe default: a model built on independent grids (vision encoder + LLM, GTP, MIMO) " + "would silently get the wrong one. " + "See docs/developer/parallel-state-deprecation.md" + ) self.pg_collection = pg_collection self.cp_group = pg_collection.cp self.tp_group = get_tensor_model_parallel_group_if_none(pg_collection.tp) diff --git a/megatron/elastification/pretrain_hybrid_flex.py b/megatron/elastification/pretrain_hybrid_flex.py index 13eeca1f7c8..f60f33fc27e 100644 --- a/megatron/elastification/pretrain_hybrid_flex.py +++ b/megatron/elastification/pretrain_hybrid_flex.py @@ -7,6 +7,7 @@ import torch +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core import mpu, parallel_state from megatron.core.datasets.blended_megatron_dataset_builder import BlendedMegatronDatasetBuilder from megatron.core.datasets.gpt_dataset import GPTDataset, GPTDatasetConfig, MockGPTDataset @@ -140,7 +141,9 @@ def model_provider(pre_process=True, post_process=True, vp_stage: Optional[int] rotary_percent=args.rotary_percent, rotary_base=args.rotary_base, vp_stage=vp_stage - ) + , + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) from megatron.elastification.flextron_utils import ( inject_flextron_forward_logic, setup_flextron_model, diff --git a/tests/unit_tests/a2a_overlap/test_cuda_graphed_schedule_chunk_1f1b.py b/tests/unit_tests/a2a_overlap/test_cuda_graphed_schedule_chunk_1f1b.py index b4351cbbe1e..e08ce894c20 100644 --- a/tests/unit_tests/a2a_overlap/test_cuda_graphed_schedule_chunk_1f1b.py +++ b/tests/unit_tests/a2a_overlap/test_cuda_graphed_schedule_chunk_1f1b.py @@ -6,6 +6,7 @@ import pytest import torch +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core.enums import ModelType from megatron.core.models.gpt.gpt_layer_specs import ( get_gpt_decoder_block_spec, @@ -110,7 +111,9 @@ def model_provider( position_embedding_type=args.position_embedding_type, rotary_percent=args.rotary_percent, mtp_block_spec=mtp_block_spec, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) def create_test_args( self, cuda_graph_impl, cuda_graph_modules, cuda_graph_warmup_steps, ep_size, **kwargs diff --git a/tests/unit_tests/a2a_overlap/test_schedule_chunk_1f1b.py b/tests/unit_tests/a2a_overlap/test_schedule_chunk_1f1b.py index bf8cf45f6be..e8970f9ccb2 100644 --- a/tests/unit_tests/a2a_overlap/test_schedule_chunk_1f1b.py +++ b/tests/unit_tests/a2a_overlap/test_schedule_chunk_1f1b.py @@ -4,6 +4,7 @@ import pytest import torch +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core.models.common.model_chunk_schedule_plan import TransformerModelChunkSchedulePlan from megatron.core.models.gpt.gpt_layer_specs import ( get_gpt_decoder_block_spec, @@ -62,7 +63,9 @@ def build_model(config, use_padding_mask=False): pre_process=True, post_process=True, max_sequence_length=max_seq_len, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) f_schedule_plan = gpt_model.build_schedule_plan(**data) return gpt_model, f_schedule_plan, data diff --git a/tests/unit_tests/a2a_overlap/test_schedule_layer_1f1b.py b/tests/unit_tests/a2a_overlap/test_schedule_layer_1f1b.py index fe36e052e03..8af4b592149 100644 --- a/tests/unit_tests/a2a_overlap/test_schedule_layer_1f1b.py +++ b/tests/unit_tests/a2a_overlap/test_schedule_layer_1f1b.py @@ -5,6 +5,7 @@ import torch import torch.nn.functional as F +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core.fp8_utils import get_fp8_context from megatron.core.models.common.model_chunk_schedule_plan import TransformerLayerSchedulePlan from megatron.core.models.gpt.gpt_layer_specs import ( @@ -304,7 +305,9 @@ def test_transformer_layer_overlap_dense(self): pre_process=True, post_process=True, max_sequence_length=300, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) params = reset_model(gpt_model) input_tensors = [build_data() for _ in range(microbatches)] @@ -347,7 +350,9 @@ def test_transformer_layer_overlap_shared_expert(self): pre_process=True, post_process=True, max_sequence_length=300, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) params = reset_model(gpt_model) input_tensors = [build_data() for _ in range(microbatches)] @@ -366,7 +371,9 @@ def test_transformer_layer_overlap_shared_expert(self): pre_process=True, post_process=True, max_sequence_length=300, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) reset_model(gpt_model, params) capture_a2a_overlap = run_transformer_layer_a2a_overlap_with_capture( gpt_model, input_tensors, microbatches @@ -399,7 +406,9 @@ def test_transformer_layer_overlap_early_attn_memory_release(self): pre_process=True, post_process=True, max_sequence_length=300, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) params = reset_model(gpt_model) input_tensors = [build_data() for _ in range(microbatches)] @@ -418,7 +427,9 @@ def test_transformer_layer_overlap_early_attn_memory_release(self): pre_process=True, post_process=True, max_sequence_length=300, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) reset_model(gpt_model, params) capture_a2a_overlap = run_transformer_layer_a2a_overlap_with_capture( gpt_model, input_tensors, microbatches @@ -452,7 +463,9 @@ def test_transformer_layer_overlap(self, dispatcher_type, flex_backend, fp8_flag pre_process=True, post_process=True, max_sequence_length=300, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) params = reset_model(gpt_model) input_tensors = [build_data() for _ in range(microbatches)] @@ -512,7 +525,9 @@ def test_transformer_layer_overlap_zero_copy(self): pre_process=True, post_process=True, max_sequence_length=300, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) params = reset_model(gpt_model) input_tensors = [build_data() for _ in range(microbatches)] @@ -584,7 +599,9 @@ def test_mtp_layer_overlap(self, dispatcher_type, flex_backend, fp8_flag): pre_process=True, post_process=True, max_sequence_length=300, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) gpt_model.decoder.final_layernorm = None gpt_model.cuda() params = reset_model(gpt_model) diff --git a/tests/unit_tests/a2a_overlap/utils.py b/tests/unit_tests/a2a_overlap/utils.py index c4d6a2844e1..7fee06ea61e 100644 --- a/tests/unit_tests/a2a_overlap/utils.py +++ b/tests/unit_tests/a2a_overlap/utils.py @@ -5,6 +5,7 @@ import torch +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core import config from megatron.core.transformer.transformer_config import MLATransformerConfig from megatron.core.utils import is_te_min_version @@ -303,7 +304,9 @@ def build_gpt_model(config, vocab_size=512, max_seq_len=300): pre_process=True, post_process=True, max_sequence_length=max_seq_len, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) model.cuda() return model diff --git a/tests/unit_tests/determinism/correctness/test_gpt_model.py b/tests/unit_tests/determinism/correctness/test_gpt_model.py index acf81ef8508..b08beea2154 100644 --- a/tests/unit_tests/determinism/correctness/test_gpt_model.py +++ b/tests/unit_tests/determinism/correctness/test_gpt_model.py @@ -16,6 +16,7 @@ import pytest import torch +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core.models.gpt.gpt_layer_specs import get_gpt_layer_with_transformer_engine_spec from megatron.core.models.gpt.gpt_model import GPTModel from megatron.core.transformer.transformer_config import TransformerConfig @@ -41,7 +42,9 @@ def build_gpt(overrides, pre_process=True, post_process=True, vp_stage=None, **_ post_process=post_process, vp_stage=vp_stage, position_embedding_type="rope", - ).cuda() + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ).cuda() def make_gpt_inputs(): diff --git a/tests/unit_tests/determinism/correctness/test_hybrid_model.py b/tests/unit_tests/determinism/correctness/test_hybrid_model.py index d2dd8757dc8..48597e20595 100644 --- a/tests/unit_tests/determinism/correctness/test_hybrid_model.py +++ b/tests/unit_tests/determinism/correctness/test_hybrid_model.py @@ -10,6 +10,7 @@ import pytest import torch +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core.models.hybrid.hybrid_layer_specs import hybrid_stack_spec from megatron.core.models.hybrid.hybrid_model import HybridModel from megatron.core.transformer.transformer_config import TransformerConfig @@ -108,7 +109,9 @@ def build(overrides, pre_process=True, post_process=True, vp_stage=None, **_): pre_process=pre_process, post_process=post_process, vp_stage=vp_stage, - ).cuda() + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ).cuda() runner = BitExactRunner( build_model=build, diff --git a/tests/unit_tests/dist_checkpointing/models/test_bert_model.py b/tests/unit_tests/dist_checkpointing/models/test_bert_model.py index 81b01c8f886..2425d279aeb 100644 --- a/tests/unit_tests/dist_checkpointing/models/test_bert_model.py +++ b/tests/unit_tests/dist_checkpointing/models/test_bert_model.py @@ -5,6 +5,7 @@ import pytest import torch +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core import parallel_state as ps from megatron.core.models.bert.bert_layer_specs import ( bert_layer_local_spec, @@ -49,7 +50,9 @@ def initialize_bert_model( pre_process=pre_process, post_process=post_process, num_tokentypes=0, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) with torch.no_grad(): for p in model.parameters(): diff --git a/tests/unit_tests/dist_checkpointing/models/test_gpt_hybrid_interop.py b/tests/unit_tests/dist_checkpointing/models/test_gpt_hybrid_interop.py index dae973f4a7b..693227aae38 100644 --- a/tests/unit_tests/dist_checkpointing/models/test_gpt_hybrid_interop.py +++ b/tests/unit_tests/dist_checkpointing/models/test_gpt_hybrid_interop.py @@ -19,6 +19,7 @@ import pytest import torch +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core import parallel_state as ps from megatron.core.dist_checkpointing import load, load_plain_tensors, save from megatron.core.dist_checkpointing.dict_utils import diff @@ -349,7 +350,9 @@ def initialize_gpt_model(seed, num_gpt_layers, parallel, moe, glu=False): post_process=ps.is_pipeline_last_stage(), position_embedding_type='rope', share_embeddings_and_output_weights=True, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) with torch.no_grad(): for param in model.parameters(): param.random_() @@ -372,7 +375,9 @@ def initialize_hybrid_model(seed, pattern, parallel, moe, glu=False): post_process=ps.is_pipeline_last_stage(), position_embedding_type='rope', share_embeddings_and_output_weights=True, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) def _snapshot_fresh_layers(hybrid_model, layer_maps): @@ -546,7 +551,9 @@ def gpt_provider_for_opt( post_process=post_process, position_embedding_type='rope', share_embeddings_and_output_weights=True, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) def hybrid_provider_for_opt( @@ -565,7 +572,9 @@ def hybrid_provider_for_opt( post_process=post_process, position_embedding_type='rope', share_embeddings_and_output_weights=True, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) def _inner_optimizers(optimizer): diff --git a/tests/unit_tests/dist_checkpointing/models/test_gpt_model.py b/tests/unit_tests/dist_checkpointing/models/test_gpt_model.py index e18d3b4683b..fd85c116a19 100644 --- a/tests/unit_tests/dist_checkpointing/models/test_gpt_model.py +++ b/tests/unit_tests/dist_checkpointing/models/test_gpt_model.py @@ -8,6 +8,7 @@ import pytest import torch +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core import parallel_state as ps from megatron.core.models.gpt.gpt_layer_specs import get_gpt_layer_local_spec as gpt_local_spec from megatron.core.models.gpt.gpt_layer_specs import ( @@ -56,7 +57,9 @@ def initialize_gpt_model(seed, layer_spec_fn=gpt_te_spec, vocab_size=128, **conf max_sequence_length=4, pre_process=pre_process, post_process=post_process, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) with torch.no_grad(): for p in model.parameters(): diff --git a/tests/unit_tests/dist_checkpointing/models/test_t5_model.py b/tests/unit_tests/dist_checkpointing/models/test_t5_model.py index e393c806a94..e83c86c7b80 100644 --- a/tests/unit_tests/dist_checkpointing/models/test_t5_model.py +++ b/tests/unit_tests/dist_checkpointing/models/test_t5_model.py @@ -3,6 +3,7 @@ import pytest import torch +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core import parallel_state as ps from megatron.core.dist_checkpointing import load, save from megatron.core.dist_checkpointing.validation import StrictHandling @@ -64,7 +65,9 @@ def initialize_t5_model(seed, encoder_decoder_spec_fn, num_layers=8, **config_kw post_process=post_process, add_encoder=add_encoder, add_decoder=add_decoder, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) with torch.no_grad(): for p in model.parameters(): diff --git a/tests/unit_tests/dist_checkpointing/test_layer_wise_optimizer.py b/tests/unit_tests/dist_checkpointing/test_layer_wise_optimizer.py index 42ef0a401ee..52e10153cdf 100644 --- a/tests/unit_tests/dist_checkpointing/test_layer_wise_optimizer.py +++ b/tests/unit_tests/dist_checkpointing/test_layer_wise_optimizer.py @@ -110,7 +110,9 @@ def initialize_real_model( pre_process=pre_process, post_process=post_process, vp_stage=vp_stage, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) return this_model diff --git a/tests/unit_tests/dist_checkpointing/test_optimizer.py b/tests/unit_tests/dist_checkpointing/test_optimizer.py index f93e09a43b7..b6e01a1e48d 100644 --- a/tests/unit_tests/dist_checkpointing/test_optimizer.py +++ b/tests/unit_tests/dist_checkpointing/test_optimizer.py @@ -10,6 +10,7 @@ import torch from torch.optim import Adam +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core import parallel_state from megatron.core.dist_checkpointing import ShardedTensor, load, load_plain_tensors, save from megatron.core.dist_checkpointing.dict_utils import diff, nested_values @@ -329,7 +330,9 @@ def initialize_real_model( pre_process=pre_process, post_process=post_process, vp_stage=vp_stage, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) return this_model diff --git a/tests/unit_tests/dist_checkpointing/test_pipeline_parallel_layout.py b/tests/unit_tests/dist_checkpointing/test_pipeline_parallel_layout.py index 4ed91aa2cb6..eaccc04c80a 100644 --- a/tests/unit_tests/dist_checkpointing/test_pipeline_parallel_layout.py +++ b/tests/unit_tests/dist_checkpointing/test_pipeline_parallel_layout.py @@ -6,6 +6,7 @@ import pytest import torch +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core import mpu from megatron.core.dist_checkpointing.strategies.cached_metadata_filesystem_reader import ( CachedMetadataFileSystemReader, @@ -73,7 +74,9 @@ def initialize_gpt_model( pre_process=pre_process, post_process=post_process, vp_stage=i, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) this_model.model_type = ModelType.encoder_or_decoder model.append(this_model) diff --git a/tests/unit_tests/dist_checkpointing/utils.py b/tests/unit_tests/dist_checkpointing/utils.py index ba774b34fd2..09de1a6e19b 100644 --- a/tests/unit_tests/dist_checkpointing/utils.py +++ b/tests/unit_tests/dist_checkpointing/utils.py @@ -6,6 +6,7 @@ import torch +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core.dist_checkpointing.strategies.cached_metadata_filesystem_reader import ( CachedMetadataFileSystemReader, ) @@ -54,7 +55,9 @@ def initialize_gpt_model( max_sequence_length=4, pre_process=pre_process, post_process=post_process, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) with torch.no_grad(): for p in model.parameters(): @@ -106,7 +109,9 @@ def initialize_moe_model( max_sequence_length=4, pre_process=pre_process, post_process=post_process, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) model.bfloat16() with torch.no_grad(): diff --git a/tests/unit_tests/distributed/test_finalize_model_grads.py b/tests/unit_tests/distributed/test_finalize_model_grads.py index 372f8d0d293..3ac3eebfafc 100644 --- a/tests/unit_tests/distributed/test_finalize_model_grads.py +++ b/tests/unit_tests/distributed/test_finalize_model_grads.py @@ -240,7 +240,9 @@ def init_model(self, share_embeddings_and_output_weights: bool = False): vocab_size=100, max_sequence_length=4, share_embeddings_and_output_weights=share_embeddings_and_output_weights, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) def setup_method(self, method): os.environ.pop('NVTE_FUSED_ATTN', None) diff --git a/tests/unit_tests/export/trtllm/test_distributed_fp8.py b/tests/unit_tests/export/trtllm/test_distributed_fp8.py index 09732e98263..db8239214d2 100644 --- a/tests/unit_tests/export/trtllm/test_distributed_fp8.py +++ b/tests/unit_tests/export/trtllm/test_distributed_fp8.py @@ -8,6 +8,7 @@ from torch.optim import Adam from torch.utils.data import DataLoader +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core.datasets.blended_megatron_dataset_builder import BlendedMegatronDatasetBuilder from megatron.core.datasets.gpt_dataset import GPTDatasetConfig, MockGPTDataset from megatron.core.datasets.utils import compile_helpers @@ -50,7 +51,9 @@ def _model_provider(): transformer_layer_spec=get_gpt_layer_with_transformer_engine_spec(), vocab_size=VOCAB_SIZE, max_sequence_length=SEQUENCE_LENGTH, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) return gpt_model diff --git a/tests/unit_tests/export/trtllm/test_single_device_fp8.py b/tests/unit_tests/export/trtllm/test_single_device_fp8.py index 2e0537894de..e6602b89793 100644 --- a/tests/unit_tests/export/trtllm/test_single_device_fp8.py +++ b/tests/unit_tests/export/trtllm/test_single_device_fp8.py @@ -8,6 +8,7 @@ from torch.optim import Adam from torch.utils.data import DataLoader +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core.datasets.blended_megatron_dataset_builder import BlendedMegatronDatasetBuilder from megatron.core.datasets.gpt_dataset import GPTDatasetConfig, MockGPTDataset from megatron.core.datasets.utils import compile_helpers @@ -47,7 +48,9 @@ def _model_provider(): transformer_layer_spec=get_gpt_layer_with_transformer_engine_spec(), vocab_size=100, max_sequence_length=SEQUENCE_LENGTH, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) return gpt_model diff --git a/tests/unit_tests/export/trtllm/test_trtllm_distributed_gpu_converter.py b/tests/unit_tests/export/trtllm/test_trtllm_distributed_gpu_converter.py index 6a5ccb04a23..b5953bca9e0 100644 --- a/tests/unit_tests/export/trtllm/test_trtllm_distributed_gpu_converter.py +++ b/tests/unit_tests/export/trtllm/test_trtllm_distributed_gpu_converter.py @@ -1,6 +1,7 @@ import torch from pytest_mock import mocker +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core.export.data_type import DataType from megatron.core.export.trtllm.model_to_trllm_mapping.default_conversion_dict import ( DEFAULT_CONVERSION_DICT, @@ -46,7 +47,9 @@ def setup_method(self, method): transformer_layer_spec=get_gpt_layer_local_spec(), vocab_size=_VOCAB_SIZE, max_sequence_length=_SEQUENCE_LENGTH, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) def teardown_method(self, method): """ diff --git a/tests/unit_tests/extension/test_te_lmhead_column_parallel_linear.py b/tests/unit_tests/extension/test_te_lmhead_column_parallel_linear.py index dbde4c67df4..d62e057f35e 100644 --- a/tests/unit_tests/extension/test_te_lmhead_column_parallel_linear.py +++ b/tests/unit_tests/extension/test_te_lmhead_column_parallel_linear.py @@ -7,6 +7,7 @@ import pytest import torch +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core import tensor_parallel from megatron.core.extensions.transformer_engine import HAVE_TE, TELMHeadColumnParallelLinear from megatron.core.fp8_utils import is_mxfp8_output_proj_active @@ -114,7 +115,9 @@ def test_default_uses_column_parallel_linear(self): transformer_layer_spec=get_gpt_layer_with_transformer_engine_spec(), vocab_size=100, max_sequence_length=4, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) assert isinstance(model.output_layer, tensor_parallel.ColumnParallelLinear) assert not isinstance(model.output_layer, TELMHeadColumnParallelLinear) @@ -140,5 +143,7 @@ def test_mxfp8_active_uses_te_lm_head(self): transformer_layer_spec=get_gpt_layer_with_transformer_engine_spec(), vocab_size=128, max_sequence_length=8, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) assert isinstance(model.output_layer, TELMHeadColumnParallelLinear) diff --git a/tests/unit_tests/generalized_tensor_parallel/test_gtp_muon_dcp.py b/tests/unit_tests/generalized_tensor_parallel/test_gtp_muon_dcp.py index b26d8a974ce..7de5a56d403 100644 --- a/tests/unit_tests/generalized_tensor_parallel/test_gtp_muon_dcp.py +++ b/tests/unit_tests/generalized_tensor_parallel/test_gtp_muon_dcp.py @@ -10,6 +10,7 @@ import torch +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core.dist_checkpointing import load, save from tests.unit_tests.dist_checkpointing import TempNamedDir, setup_model_and_optimizer from tests.unit_tests.test_utilities import Utils @@ -90,7 +91,9 @@ def _initialize_native_fp8_moe_model( max_sequence_length=4, pre_process=pre_process, post_process=post_process, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) with torch.no_grad(): for p in model.parameters(): if not is_float8tensor(p): diff --git a/tests/unit_tests/inference/engines/test_dynamic_engine.py b/tests/unit_tests/inference/engines/test_dynamic_engine.py index af91409ef7b..b989a7f89a1 100644 --- a/tests/unit_tests/inference/engines/test_dynamic_engine.py +++ b/tests/unit_tests/inference/engines/test_dynamic_engine.py @@ -18,6 +18,7 @@ from tqdm import tqdm from transformer_engine.pytorch.fp8 import check_fp8_support +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core import parallel_state from megatron.core.inference.config import ( AsyncScheduleMode, @@ -430,7 +431,9 @@ def _build_test_env(cls, test_config): post_process=parallel_state.is_pipeline_last_stage(), mtp_block_spec=mtp_block_spec, position_embedding_type=test_config.position_embedding_type, - ).cuda() + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ).cuda() elif test_config.model_provider == "hybrid": pp_size = test_config.pipeline_model_parallel_size # Transformer config. @@ -497,7 +500,9 @@ def _build_test_env(cls, test_config): hybrid_layer_pattern=mamba_pattern, pre_process=parallel_state.is_pipeline_first_stage(), post_process=parallel_state.is_pipeline_last_stage(), - ).cuda() + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ).cuda() else: raise ValueError(f"Invalid model provider {test_config.model_provider}") @@ -5854,7 +5859,9 @@ def _create_model(self, model_provider, num_cuda_graphs): parallel_output=True, pre_process=parallel_state.is_pipeline_first_stage(), post_process=parallel_state.is_pipeline_last_stage(), - ).cuda() + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ).cuda() elif model_provider == "hybrid": config = TransformerConfig( params_dtype=torch.bfloat16, @@ -5880,7 +5887,9 @@ def _create_model(self, model_provider, num_cuda_graphs): hybrid_layer_pattern="M*-", pre_process=parallel_state.is_pipeline_first_stage(), post_process=parallel_state.is_pipeline_last_stage(), - ).cuda() + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ).cuda() else: raise ValueError(f"Invalid model_provider {model_provider}") diff --git a/tests/unit_tests/inference/engines/test_hybrid_prefix_caching_e2e.py b/tests/unit_tests/inference/engines/test_hybrid_prefix_caching_e2e.py index be57f8a9e15..97ebf00e5c6 100644 --- a/tests/unit_tests/inference/engines/test_hybrid_prefix_caching_e2e.py +++ b/tests/unit_tests/inference/engines/test_hybrid_prefix_caching_e2e.py @@ -38,6 +38,7 @@ import pytest import torch +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core import parallel_state from megatron.core.inference.config import ( AsyncScheduleMode, @@ -141,7 +142,9 @@ def _create_model(self, num_cuda_graphs=None): hybrid_layer_pattern="M*-", pre_process=parallel_state.is_pipeline_first_stage(), post_process=parallel_state.is_pipeline_last_stage(), - ).cuda() + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ).cuda() for param in model.parameters(): param.data = param.data.to(transformer_config.params_dtype) model.eval() diff --git a/tests/unit_tests/inference/engines/test_prefix_caching_cuda_graphs.py b/tests/unit_tests/inference/engines/test_prefix_caching_cuda_graphs.py index a3dc42f1c71..87e220f59e8 100644 --- a/tests/unit_tests/inference/engines/test_prefix_caching_cuda_graphs.py +++ b/tests/unit_tests/inference/engines/test_prefix_caching_cuda_graphs.py @@ -18,6 +18,7 @@ import pytest import torch +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core import parallel_state from megatron.core.inference.config import ( InferenceConfig, @@ -102,7 +103,9 @@ def _create_model(self, model_type, num_cuda_graphs=None): parallel_output=True, pre_process=parallel_state.is_pipeline_first_stage(), post_process=parallel_state.is_pipeline_last_stage(), - ).cuda() + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ).cuda() mamba_config = None else: # hybrid config = TransformerConfig( @@ -129,7 +132,9 @@ def _create_model(self, model_type, num_cuda_graphs=None): hybrid_layer_pattern="M*-", pre_process=parallel_state.is_pipeline_first_stage(), post_process=parallel_state.is_pipeline_last_stage(), - ).cuda() + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ).cuda() mamba_config = MambaInferenceStateConfig.from_model(model) for param in model.parameters(): @@ -361,7 +366,9 @@ def _create_hybrid_model(self, num_cuda_graphs=None): hybrid_layer_pattern="M*-", pre_process=parallel_state.is_pipeline_first_stage(), post_process=parallel_state.is_pipeline_last_stage(), - ).cuda() + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ).cuda() for param in model.parameters(): param.data = param.data.to(config.params_dtype) model.eval() diff --git a/tests/unit_tests/inference/engines/test_static_engine.py b/tests/unit_tests/inference/engines/test_static_engine.py index d40766b6a26..4d9746ffcf7 100644 --- a/tests/unit_tests/inference/engines/test_static_engine.py +++ b/tests/unit_tests/inference/engines/test_static_engine.py @@ -9,6 +9,7 @@ import pytest import torch +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core import parallel_state from megatron.core.inference.contexts import StaticInferenceContext from megatron.core.inference.engines import StaticInferenceEngine @@ -77,7 +78,9 @@ def setup_engine( parallel_output=True, pre_process=parallel_state.is_pipeline_first_stage(), post_process=parallel_state.is_pipeline_last_stage(), - ).cuda() + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ).cuda() gpt_model.to(inference_config_params_dtype) inference_context = StaticInferenceContext( diff --git a/tests/unit_tests/inference/model_inference_wrappers/gpt/test_gpt_inference_wrapper.py b/tests/unit_tests/inference/model_inference_wrappers/gpt/test_gpt_inference_wrapper.py index e4b68995b37..bfbb92722ca 100644 --- a/tests/unit_tests/inference/model_inference_wrappers/gpt/test_gpt_inference_wrapper.py +++ b/tests/unit_tests/inference/model_inference_wrappers/gpt/test_gpt_inference_wrapper.py @@ -3,6 +3,7 @@ import pytest import torch +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core import parallel_state from megatron.core.inference.contexts import StaticInferenceContext from megatron.core.inference.model_inference_wrappers.gpt.gpt_inference_wrapper import ( @@ -44,7 +45,9 @@ def setup_model(self, tensor_parallel_size, pipeline_parallel_size): parallel_output=True, pre_process=parallel_state.is_pipeline_first_stage(), post_process=parallel_state.is_pipeline_last_stage(), - ).cuda() + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ).cuda() inference_context = StaticInferenceContext(self.batch_size, self.sequence_length) diff --git a/tests/unit_tests/inference/model_inference_wrappers/t5/test_t5_inference_wrapper.py b/tests/unit_tests/inference/model_inference_wrappers/t5/test_t5_inference_wrapper.py index e04b25ea449..6dc0cf8b842 100644 --- a/tests/unit_tests/inference/model_inference_wrappers/t5/test_t5_inference_wrapper.py +++ b/tests/unit_tests/inference/model_inference_wrappers/t5/test_t5_inference_wrapper.py @@ -7,6 +7,7 @@ import numpy as np import torch +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core import parallel_state from megatron.core.inference.contexts import StaticInferenceContext from megatron.core.inference.model_inference_wrappers.t5.t5_inference_wrapper import ( @@ -75,7 +76,9 @@ def setup_model(self, tensor_parallel_size, pipeline_parallel_size): post_process=True, add_encoder=True, add_decoder=True, - ).cuda() + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ).cuda() inference_context = StaticInferenceContext(max_batch_size=8, max_sequence_length=2560) diff --git a/tests/unit_tests/inference/test_hybrid_moe.py b/tests/unit_tests/inference/test_hybrid_moe.py index 9587cc2d285..aad07b8d406 100644 --- a/tests/unit_tests/inference/test_hybrid_moe.py +++ b/tests/unit_tests/inference/test_hybrid_moe.py @@ -21,6 +21,7 @@ import pytest import torch +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core import parallel_state from megatron.core.inference.batch_dimensions_utils import InferenceBatchDimensions from megatron.core.inference.config import InferenceConfig, MambaInferenceStateConfig @@ -163,7 +164,9 @@ def _build_model(self, inference_moe_token_dispatcher_type='nvls'): vocab_size=self.VOCAB_SIZE, max_sequence_length=self.MAX_SEQ_LEN, hybrid_layer_pattern="ME*", - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) model.cuda() model.eval() return model @@ -300,7 +303,9 @@ def test_batch_invariant_prefill_matches_full_forward(self): vocab_size=self.VOCAB_SIZE, max_sequence_length=self.MAX_SEQ_LEN, hybrid_layer_pattern="ME*", - ).cuda() + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ).cuda() model.eval() input_ids = torch.arange(64, device="cuda", dtype=torch.long).unsqueeze(0) diff --git a/tests/unit_tests/inference/test_mtp_cuda_graph_inference.py b/tests/unit_tests/inference/test_mtp_cuda_graph_inference.py index ca604e8d25d..243e1091d05 100644 --- a/tests/unit_tests/inference/test_mtp_cuda_graph_inference.py +++ b/tests/unit_tests/inference/test_mtp_cuda_graph_inference.py @@ -20,6 +20,7 @@ import torch import torch.distributed as dist +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core import parallel_state from megatron.core.inference.batch_dimensions_utils import InferenceBatchDimensions from megatron.core.inference.config import InferenceConfig @@ -127,7 +128,9 @@ def _build_model( pre_process=True, post_process=True, mtp_block_spec=mtp_block_spec, - ).cuda() + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ).cuda() for param in model.parameters(): param.data = param.data.to(config.params_dtype) model.eval() @@ -885,7 +888,9 @@ def _build_model(self, inference_moe_token_dispatcher_type='nccl'): pre_process=True, post_process=True, mtp_block_spec=mtp_block_spec, - ).cuda() + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ).cuda() for param in model.parameters(): param.data = param.data.to(config.params_dtype) model.eval() @@ -1207,7 +1212,9 @@ def _build_model(self, *, inference_cuda_graph_scope='block', model_type='hybrid pre_process=True, post_process=True, position_embedding_type='rope', - ).cuda() + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ).cuda() elif model_type == 'hybrid': hybrid_stack_spec = _build_hybrid_stack_spec() model = HybridModel( @@ -1220,7 +1227,9 @@ def _build_model(self, *, inference_cuda_graph_scope='block', model_type='hybrid post_process=True, hybrid_layer_pattern="****/*", position_embedding_type='rope', - ).cuda() + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ).cuda() else: raise ValueError(f"Unknown model_type: {model_type!r}") for param in model.parameters(): diff --git a/tests/unit_tests/inference/text_generation_controllers/test_encoder_decoder_text_generation_controller.py b/tests/unit_tests/inference/text_generation_controllers/test_encoder_decoder_text_generation_controller.py index 4ddcc66427b..a27b920a620 100644 --- a/tests/unit_tests/inference/text_generation_controllers/test_encoder_decoder_text_generation_controller.py +++ b/tests/unit_tests/inference/text_generation_controllers/test_encoder_decoder_text_generation_controller.py @@ -12,6 +12,7 @@ import pytest import torch +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core.inference.contexts import StaticInferenceContext from megatron.core.inference.inference_request import InferenceRequest, Status from megatron.core.inference.model_inference_wrappers.t5.t5_inference_wrapper import ( @@ -83,7 +84,9 @@ def setup_method(self, method): post_process=True, add_encoder=True, add_decoder=True, - ).cuda() + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ).cuda() inference_context = StaticInferenceContext(max_batch_size=8, max_sequence_length=2560) diff --git a/tests/unit_tests/inference/text_generation_controllers/test_text_generation_controller.py b/tests/unit_tests/inference/text_generation_controllers/test_text_generation_controller.py index fbc05f7ff71..25abe4c31dd 100644 --- a/tests/unit_tests/inference/text_generation_controllers/test_text_generation_controller.py +++ b/tests/unit_tests/inference/text_generation_controllers/test_text_generation_controller.py @@ -15,6 +15,7 @@ import torch from transformer_engine.pytorch.fp8 import check_fp8_support +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core import parallel_state from megatron.core.inference.config import ( AsyncScheduleMode, @@ -130,7 +131,9 @@ def setup_model( hybrid_layer_pattern=hybrid_layer_pattern, pre_process=parallel_state.is_pipeline_first_stage(), post_process=parallel_state.is_pipeline_last_stage(), - ).cuda() + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ).cuda() mamba_inference_state_config = MambaInferenceStateConfig.from_model(model) else: layer_spec = get_gpt_layer_local_spec() @@ -150,7 +153,9 @@ def setup_model( pre_process=parallel_state.is_pipeline_first_stage(), post_process=parallel_state.is_pipeline_last_stage(), mtp_block_spec=mtp_block_spec, - ).cuda() + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ).cuda() model.eval() if dtype == torch.bfloat16: diff --git a/tests/unit_tests/models/test_bert_model.py b/tests/unit_tests/models/test_bert_model.py index e878978f64e..be14a085a48 100644 --- a/tests/unit_tests/models/test_bert_model.py +++ b/tests/unit_tests/models/test_bert_model.py @@ -6,6 +6,7 @@ import torch from packaging.version import Version as PkgVersion +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core.models.bert.bert_layer_specs import ( bert_layer_local_spec, get_bert_layer_with_transformer_engine_spec, @@ -45,7 +46,9 @@ def setup_method(self, method): transformer_layer_spec=get_bert_layer_with_transformer_engine_spec(), vocab_size=100, max_sequence_length=4, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) def teardown_method(self, method): Utils.destroy_model_parallel() @@ -107,7 +110,9 @@ def test_output_layer_bias_false_disables_bias(self): max_sequence_length=self.bert_model.max_sequence_length, apply_lm_head=False, output_layer_bias=False, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) assert bert_model.output_layer.bias is None @@ -124,7 +129,9 @@ def test_apply_lm_head_false_bypasses_head(self): vocab_size=100, max_sequence_length=sequence_length, apply_lm_head=False, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) assert bert_model.lm_head is None bert_model.cuda() @@ -179,7 +186,9 @@ def test_qk_layernorm_from_config_fallback(self): transformer_layer_spec=get_bert_layer_with_transformer_engine_spec(), vocab_size=100, max_sequence_length=4, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) attention = bert_model.encoder.layers[0].self_attention assert isinstance(attention.q_layernorm, te_pytorch.LayerNorm) assert isinstance(attention.k_layernorm, te_pytorch.LayerNorm) @@ -208,7 +217,9 @@ def setup_method(self, method): transformer_layer_spec=get_bert_layer_with_transformer_engine_spec(), vocab_size=100, max_sequence_length=4, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) @pytest.mark.internal def test_local_spec(self, mocker): @@ -272,7 +283,9 @@ def test_transformer_engine_version_1_7_to_1_10_rng_error(self, mocker): transformer_layer_spec=ModuleSpec(module=TransformerLayer, submodules=submodules), vocab_size=100, max_sequence_length=4, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) assert str(exc_info.value) == ( "Linear.__init__() got an unexpected keyword argument 'rng_tracker_name' when " "instantiating TERowParallelLinear when instantiating SelfAttention when " @@ -311,7 +324,9 @@ def test_transformer_engine_version_less_than_1_7(self, mocker): transformer_layer_spec=get_bert_layer_with_transformer_engine_spec(), vocab_size=100, max_sequence_length=4, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) assert str(exc_info.value) == ( "Flash and fused attention is not supported with transformer engine version " diff --git a/tests/unit_tests/models/test_dsa_gpt_mamba_equivalence.py b/tests/unit_tests/models/test_dsa_gpt_mamba_equivalence.py index 51568243d0d..27cf1e3cb4a 100644 --- a/tests/unit_tests/models/test_dsa_gpt_mamba_equivalence.py +++ b/tests/unit_tests/models/test_dsa_gpt_mamba_equivalence.py @@ -185,7 +185,9 @@ def _build_gpt_model( post_process=post_process, parallel_output=False, # Gather logits across TP for easy comparison position_embedding_type='rope', - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) return model.cuda() @@ -210,7 +212,9 @@ def _build_mamba_model( parallel_output=False, hybrid_layer_pattern=layer_pattern, position_embedding_type='rope', - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) return model.cuda() diff --git a/tests/unit_tests/models/test_gpt_model.py b/tests/unit_tests/models/test_gpt_model.py index d2cb12841c4..86eca36fef9 100644 --- a/tests/unit_tests/models/test_gpt_model.py +++ b/tests/unit_tests/models/test_gpt_model.py @@ -54,7 +54,9 @@ def setup_method(self, method): transformer_layer_spec=get_gpt_layer_with_transformer_engine_spec(), vocab_size=100, max_sequence_length=4, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) self.mock_log_single_rank = mock_log_single_rank def teardown_method(self, method): @@ -300,7 +302,9 @@ def setup_method(self, method) -> None: transformer_layer_spec=get_gpt_layer_with_transformer_engine_spec(use_te_op_fuser=True), vocab_size=100, max_sequence_length=4, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) def teardown_method(self, method) -> None: Utils.destroy_model_parallel() @@ -362,7 +366,9 @@ def test_gpt_with_te_activation_func(num_experts, gated_linear_unit): ), vocab_size=128, max_sequence_length=128, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) # test sequence_length = gpt_model.max_sequence_length @@ -501,7 +507,9 @@ def setup_method(self, method): vocab_size=128, max_sequence_length=DynamicInferenceContext.TOKEN_ROUNDER, parallel_output=True, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) self.gpt_model = Float16Module(self.gpt_model.config, self.gpt_model) def teardown_method(self, method): diff --git a/tests/unit_tests/models/test_gpt_model_batch_invariant.py b/tests/unit_tests/models/test_gpt_model_batch_invariant.py index 1e93687bcd7..1c69ba008d4 100644 --- a/tests/unit_tests/models/test_gpt_model_batch_invariant.py +++ b/tests/unit_tests/models/test_gpt_model_batch_invariant.py @@ -5,6 +5,7 @@ import torch import torch.distributed as dist +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core.inference.config import InferenceConfig from megatron.core.inference.contexts.dynamic_context import DynamicInferenceContext from megatron.core.inference.engines.dynamic_engine import DynamicInferenceEngine @@ -120,7 +121,9 @@ def _build_flash_attn_bik_model(seq_len: int, vocab_size: int, hidden_size: int transformer_layer_spec=get_gpt_layer_with_transformer_engine_spec(), vocab_size=vocab_size, max_sequence_length=seq_len, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) return model.cuda().eval() diff --git a/tests/unit_tests/models/test_gpt_model_quantization.py b/tests/unit_tests/models/test_gpt_model_quantization.py index 6f8d7c7b63f..930799b2e76 100644 --- a/tests/unit_tests/models/test_gpt_model_quantization.py +++ b/tests/unit_tests/models/test_gpt_model_quantization.py @@ -2,6 +2,7 @@ import pytest +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core.enums import Fp8Recipe from megatron.core.extensions.transformer_engine import HAVE_TE from megatron.core.models.gpt import GPTModel @@ -81,7 +82,9 @@ def test_kitchen_config_resolution_dense(self) -> None: transformer_layer_spec=transformer_layer_spec, vocab_size=padded_vocab_size, max_sequence_length=max_position_embeddings, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) expected_types = { "decoder.layers.0.self_attention.linear_proj": KitchenRowParallelLinear, @@ -186,7 +189,9 @@ def test_kitchen_config_resolution_dense_compound_params(self) -> None: transformer_layer_spec=transformer_layer_spec, vocab_size=padded_vocab_size, max_sequence_length=max_position_embeddings, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) expected_types = { "decoder.layers.0.self_attention.linear_proj": KitchenRowParallelLinear, @@ -304,7 +309,9 @@ def test_kitchen_config_resolution_moe(self) -> None: transformer_layer_spec=transformer_layer_spec, vocab_size=padded_vocab_size, max_sequence_length=max_position_embeddings, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) expected_types = { "decoder.layers.0.self_attention.linear_proj": KitchenRowParallelLinear, @@ -421,7 +428,9 @@ def test_kitchen_flash_attention_config_resolution(self) -> None: transformer_layer_spec=transformer_layer_spec, vocab_size=padded_vocab_size, max_sequence_length=max_position_embeddings, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) expected_types = { "decoder.layers.0.self_attention.linear_proj": KitchenRowParallelLinear, @@ -535,7 +544,9 @@ def test_kitchen_flash_attention_with_compound_params(self) -> None: transformer_layer_spec=transformer_layer_spec, vocab_size=padded_vocab_size, max_sequence_length=max_position_embeddings, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) expected_types = { "decoder.layers.0.self_attention.linear_proj": KitchenRowParallelLinear, @@ -628,7 +639,9 @@ def test_te_config_resolution_dense(self) -> None: transformer_layer_spec=transformer_layer_spec, vocab_size=padded_vocab_size, max_sequence_length=max_position_embeddings, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) expected_types = { "decoder.layers.0.self_attention.linear_proj": TERowParallelLinear, diff --git a/tests/unit_tests/models/test_heterogeneous_gpt_model.py b/tests/unit_tests/models/test_heterogeneous_gpt_model.py index 56d112021c7..7eb41275f08 100644 --- a/tests/unit_tests/models/test_heterogeneous_gpt_model.py +++ b/tests/unit_tests/models/test_heterogeneous_gpt_model.py @@ -5,6 +5,7 @@ import pytest import torch +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core.models.gpt.gpt_model import GPTModel from megatron.core.models.gpt.heterogeneous.heterogeneous_layer_specs import ( get_gpt_heterogeneous_layer_spec, @@ -73,7 +74,9 @@ def heterogeneous_gpt_model(request, tmp_path): vocab_size=128256, position_embedding_type="rope", max_sequence_length=4, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) @pytest.mark.parametrize( diff --git a/tests/unit_tests/models/test_hybrid_model.py b/tests/unit_tests/models/test_hybrid_model.py index 95bcaa2d7d0..338ae8ab628 100644 --- a/tests/unit_tests/models/test_hybrid_model.py +++ b/tests/unit_tests/models/test_hybrid_model.py @@ -12,6 +12,7 @@ import torch from transformer_engine.pytorch.fp8 import check_fp8_support +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core import parallel_state from megatron.core.hyper_comm_grid import HyperCommGrid from megatron.core.inference.config import InferenceConfig, MambaInferenceStateConfig @@ -224,6 +225,7 @@ def setup_method(self, method): vocab_size=100, max_sequence_length=4, hybrid_layer_pattern="M*-", # 1 Mamba, 1 attention, 1 MLP + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), ) def teardown_method(self, method): @@ -292,6 +294,7 @@ def test_forward_packed_sequence(self): vocab_size=vocab_size, max_sequence_length=12, hybrid_layer_pattern="M*-", # 1 Mamba, 1 attention, 1 MLP + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), ) sequence_length = model.max_sequence_length @@ -424,6 +427,7 @@ def _build_model(self, spec=None, **config_overrides): vocab_size=100, max_sequence_length=4, hybrid_layer_pattern="M*-", + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), ) def _get_attention_layer(self, model): @@ -533,6 +537,7 @@ def _build_model(self, spec=None, **config_overrides): vocab_size=100, max_sequence_length=4, hybrid_layer_pattern="M+-", + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), ) def test_qk_l2_norm_from_config(self): @@ -594,6 +599,7 @@ def _build_model(self, spec=None, **config_overrides): vocab_size=100, max_sequence_length=4, hybrid_layer_pattern="MD-", + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), ) def test_qk_l2_norm_from_config(self): @@ -655,6 +661,7 @@ def _build_model(self, spec=None, **config_overrides): vocab_size=100, max_sequence_length=4, hybrid_layer_pattern=self.hybrid_layer_pattern, + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), ) def _get_mla_attention(self, model): @@ -952,6 +959,7 @@ def _build_model(self, pattern="M+-", **config_overrides): vocab_size=100, max_sequence_length=4, hybrid_layer_pattern=pattern, + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), ) def _get_layer_with_mla(self, model): @@ -1156,7 +1164,8 @@ def setup_method(self, method): hybrid_stack_spec=hybrid_stack_spec, vocab_size=128, max_sequence_length=DynamicInferenceContext.TOKEN_ROUNDER, - hybrid_layer_pattern="M*", # 1 Mamba, 1 attention + hybrid_layer_pattern="M*", + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), # 1 Mamba, 1 attention ) self.model = Float16Module(self.model.config, self.model) @@ -1268,6 +1277,7 @@ def setup_method(self, method): hybrid_layer_pattern="M*-", # 1 Mamba, 1 attention, 1 MLP position_embedding_type='yarn', rotary_base=10000, + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), ) def teardown_method(self, method): diff --git a/tests/unit_tests/models/test_hybrid_moe_model.py b/tests/unit_tests/models/test_hybrid_moe_model.py index eb871568046..2575bada7d3 100644 --- a/tests/unit_tests/models/test_hybrid_moe_model.py +++ b/tests/unit_tests/models/test_hybrid_moe_model.py @@ -10,6 +10,7 @@ import pytest # type: ignore[import] import torch +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core.models.hybrid.hybrid_layer_specs import hybrid_stack_spec from megatron.core.models.hybrid.hybrid_model import HybridModel from megatron.core.num_microbatches_calculator import destroy_num_microbatches_calculator @@ -562,7 +563,9 @@ def setup_method(self, method): position_embedding_type=args.position_embedding_type, rotary_base=args.rotary_base, rotary_percent=args.rotary_percent, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) def teardown_method(self, method): Utils.destroy_model_parallel() diff --git a/tests/unit_tests/pipeline_parallel/test_fine_grained_activation_offloading.py b/tests/unit_tests/pipeline_parallel/test_fine_grained_activation_offloading.py index dce35bdd4fa..ad253e2f157 100644 --- a/tests/unit_tests/pipeline_parallel/test_fine_grained_activation_offloading.py +++ b/tests/unit_tests/pipeline_parallel/test_fine_grained_activation_offloading.py @@ -8,6 +8,7 @@ import pytest import torch +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core.models.gpt.gpt_layer_specs import get_gpt_layer_with_transformer_engine_spec from megatron.core.models.gpt.gpt_model import GPTModel from megatron.core.pipeline_parallel.fine_grained_activation_offload import ChunkOffloadHandler @@ -120,7 +121,9 @@ def _build_gpt_model( ), vocab_size=vocab_size, max_sequence_length=seq_length, - ).bfloat16() + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ).bfloat16() return gpt_model @@ -499,6 +502,8 @@ def _build_overlap_moe_gpt( ), vocab_size=vocab_size, max_sequence_length=seq_length, + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), ) .bfloat16() .cuda() @@ -698,7 +703,9 @@ def _build_gpt_model_with_cuda_graph( ), vocab_size=vocab_size, max_sequence_length=seq_length, - ).bfloat16() + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ).bfloat16() return gpt_model diff --git a/tests/unit_tests/pipeline_parallel/test_pipeline_layout.py b/tests/unit_tests/pipeline_parallel/test_pipeline_layout.py index 1c998181b50..b7e04685913 100644 --- a/tests/unit_tests/pipeline_parallel/test_pipeline_layout.py +++ b/tests/unit_tests/pipeline_parallel/test_pipeline_layout.py @@ -8,6 +8,7 @@ import torch import torch.distributed +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core import mpu, parallel_state from megatron.core.models.gpt.gpt_layer_specs import get_gpt_decoder_block_spec from megatron.core.models.gpt.gpt_layer_specs import ( @@ -104,6 +105,8 @@ def initialize_gpt_model( vp_stage=i, mtp_block_spec=mtp_block_spec, share_embeddings_and_output_weights=False, + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), ) .bfloat16() .cuda() diff --git a/tests/unit_tests/post_training/test_freeze_base_for_mtp.py b/tests/unit_tests/post_training/test_freeze_base_for_mtp.py index 647334a28d1..9a6a2cccfd7 100644 --- a/tests/unit_tests/post_training/test_freeze_base_for_mtp.py +++ b/tests/unit_tests/post_training/test_freeze_base_for_mtp.py @@ -6,6 +6,7 @@ import torch from packaging.version import Version +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core.models.gpt.gpt_layer_specs import ( get_gpt_decoder_layer_specs, get_gpt_mtp_block_spec, @@ -46,7 +47,9 @@ def setup_method(self, method): mtp_block_spec=mtp_block_spec, vocab_size=100, max_sequence_length=8, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) def teardown_method(self, method): Utils.destroy_model_parallel() diff --git a/tests/unit_tests/post_training/test_modelopt_module_spec.py b/tests/unit_tests/post_training/test_modelopt_module_spec.py index 380c5249eb0..cc9099f250c 100644 --- a/tests/unit_tests/post_training/test_modelopt_module_spec.py +++ b/tests/unit_tests/post_training/test_modelopt_module_spec.py @@ -6,6 +6,7 @@ import torch from packaging.version import Version +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core import dist_checkpointing, parallel_state from megatron.core.inference.contexts import StaticInferenceContext from megatron.core.inference.utils import InferenceMode @@ -95,7 +96,9 @@ def setup_method(self, method): transformer_layer_spec=get_gpt_layer_with_transformer_engine_spec(), vocab_size=100, max_sequence_length=4, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) # Ensure that a GPTModel can be built with the modelopt spec. self.modelopt_model = GPTModel( config=transformer_config, @@ -104,7 +107,9 @@ def setup_method(self, method): ), vocab_size=100, max_sequence_length=4, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) def test_sharded_state_dict_restore(self, tmp_path_dist_ckpt): """Save with the default TE spec and restore using the ModelOpt spec.""" @@ -159,7 +164,9 @@ def setup_method(self, method): transformer_layer_spec=default_spec, vocab_size=100, max_sequence_length=8, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) modelopt_spec = get_gpt_modelopt_spec(transformer_config, remap_te_layernorm=True) # Ensure that a GPTModel can be built with the modelopt spec. self.modelopt_model = GPTModel( @@ -167,7 +174,9 @@ def setup_method(self, method): transformer_layer_spec=modelopt_spec, vocab_size=100, max_sequence_length=8, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) class TestModelOptLlama4MoE(TestModelOptGPTModel): @@ -201,7 +210,9 @@ def setup_method(self, method): transformer_layer_spec=default_spec, vocab_size=100, max_sequence_length=8, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) modelopt_spec = get_gpt_modelopt_spec( transformer_config, remap_te_layernorm=True, qk_l2_norm=True ) @@ -211,7 +222,9 @@ def setup_method(self, method): transformer_layer_spec=modelopt_spec, vocab_size=100, max_sequence_length=8, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) class TestModelOptHybridModel(TestModelOptGPTModel): @@ -230,7 +243,9 @@ def setup_method(self, method): vocab_size=100, max_sequence_length=4, hybrid_layer_pattern="M*-", - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) # A Hybrid HybridModel using ModelOpt spec (local + TENorm). self.modelopt_model = HybridModel( @@ -239,7 +254,9 @@ def setup_method(self, method): vocab_size=100, max_sequence_length=4, hybrid_layer_pattern="M*-", - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) def test_get_gpt_modelopt_spec_interface(): diff --git a/tests/unit_tests/rl/test_rl_utils.py b/tests/unit_tests/rl/test_rl_utils.py index c37d5ec00e1..ac10a373504 100644 --- a/tests/unit_tests/rl/test_rl_utils.py +++ b/tests/unit_tests/rl/test_rl_utils.py @@ -396,7 +396,7 @@ def test_get_logprobs(self, initialize_model_parallel, use_sequence_packing): """Test that getting logprobs at least does not crash.""" self.create_test_args(rl_use_sequence_packing=use_sequence_packing) - model = MockModel() + model = MockModel(pg_collection=ProcessGroupCollection.use_mpu_process_groups()) tokens = torch.ones((BATCH, SEQ), dtype=torch.long) logprobs = rl_utils.get_logprobs( model, tokens, position_ids=None, sequence_packing=use_sequence_packing @@ -519,7 +519,7 @@ def test_prepare_data_for_update(self, initialize_model_parallel): grpo_group_size=group_size, ) - model = MockModel() + model = MockModel(pg_collection=ProcessGroupCollection.use_mpu_process_groups()) tokenizer = MockTokenizer() # A single-turn rollout whose only turn is short and lacks eod must be rejected: @@ -600,7 +600,7 @@ def test_prepare_data_for_update_oversampling(self, initialize_model_parallel): the microbatch calculator is sized by ceil(ratio * total turns), not the full batch.""" world_size, dp, tp, pp = initialize_model_parallel tokenizer = MockTokenizer() - model = MockModel() + model = MockModel(pg_collection=ProcessGroupCollection.use_mpu_process_groups()) # ratio = global_batch_size/(prompts*group) = 2*dp/(dp*4) = 0.5. # 4*dp single-turn turns (already a multiple of 2*dp); ceil(0.5 * 4*dp) = 2*dp; @@ -809,7 +809,9 @@ def test_grad_buffer_offload(self, initialize_model_parallel): transformer_layer_spec=get_gpt_layer_with_transformer_engine_spec(), vocab_size=256, max_sequence_length=32, - ).cuda() + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ).cuda() ddp_config = DistributedDataParallelConfig( grad_reduce_in_fp32=True, @@ -868,7 +870,9 @@ def test_optimizer_offload(self, initialize_model_parallel): transformer_layer_spec=get_gpt_layer_with_transformer_engine_spec(), vocab_size=256, max_sequence_length=32, - ).cuda() + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ).cuda() ddp_config = DistributedDataParallelConfig( grad_reduce_in_fp32=True, @@ -997,7 +1001,9 @@ def test_gpt_logprobs(self, initialize_model_parallel): max_sequence_length=4192, pre_process=is_pp_first_stage(pp_group), post_process=is_pp_last_stage(pp_group), - ).cuda() + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ).cuda() sequence_length = gpt_model.max_sequence_length gpt_model = Float16Module(gpt_model.config, gpt_model) @@ -1064,7 +1070,9 @@ def test_get_logprobs_cuda_graphs(self, initialize_model_parallel): transformer_layer_spec=get_gpt_layer_with_transformer_engine_spec(), vocab_size=256, max_sequence_length=32, - ).cuda() + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ).cuda() # Wrap in Float16Module so it accepts fp32_output argument from get_logprobs wrapped_model = Float16Module(transformer_config, model) diff --git a/tests/unit_tests/test_checkpointing.py b/tests/unit_tests/test_checkpointing.py index 2e717e4424f..bda2754cc91 100644 --- a/tests/unit_tests/test_checkpointing.py +++ b/tests/unit_tests/test_checkpointing.py @@ -9,6 +9,7 @@ import torch import torch.distributed.checkpoint +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core.distributed import DistributedDataParallelConfig from megatron.core.distributed.fsdp.mcore_fsdp_adapter import FullyShardedDataParallel from megatron.core.num_microbatches_calculator import ( @@ -292,7 +293,7 @@ def test_save_checkpoint(init_model_parallel, create_args, tmp_path_dist_ckpt, c iteration = 123 config = TransformerConfig(num_layers=1, kv_channels=1) - model = MockModel(config) + model = MockModel(config, pg_collection=ProcessGroupCollection.use_mpu_process_groups()) optimizer = MockState({"optimizer": "optimizer_state"}) if ckpt_format == "fsdp_dtensor": model = FullyShardedDataParallel( @@ -349,7 +350,7 @@ def test_load_checkpoint( # Create and save a checkpoint first. iteration = 123 config = TransformerConfig(num_layers=1, kv_channels=1) - model = MockModel(config) + model = MockModel(config, pg_collection=ProcessGroupCollection.use_mpu_process_groups()) optimizer = MockState({"optimizer": "optimizer_state"}) opt_param_scheduler = MockState({"opt_param_scheduler": "scheduler_state"}) @@ -360,7 +361,7 @@ def test_load_checkpoint( ) # Create new model, optimizer, and scheduler instances to load into. - new_model = MockModel(config) + new_model = MockModel(config, pg_collection=ProcessGroupCollection.use_mpu_process_groups()) new_optimizer = MockState({"optimizer": "dummy1"}) new_opt_param_scheduler = MockState({"opt_param_scheduler": "dummy2"}) @@ -396,7 +397,7 @@ def test_dist_checkpoint_versioning(init_model_parallel, tmp_path_dist_ckpt, cre # Create and save a checkpoint first. iteration = 123 config = TransformerConfig(num_layers=1, kv_channels=1) - model = MockModel(config) + model = MockModel(config, pg_collection=ProcessGroupCollection.use_mpu_process_groups()) optimizer = MockState({"optimizer": "optimizer_state"}) opt_param_scheduler = MockState({"opt_param_scheduler": "scheduler_state"}) diff --git a/tests/unit_tests/test_fp4_param.py b/tests/unit_tests/test_fp4_param.py index 3860c40ea2f..70189a0c1ff 100644 --- a/tests/unit_tests/test_fp4_param.py +++ b/tests/unit_tests/test_fp4_param.py @@ -9,6 +9,7 @@ from transformer_engine.pytorch.fp8 import check_nvfp4_support import megatron.core.parallel_state as ps +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core.distributed import DistributedDataParallel as DDP from megatron.core.enums import ModelType from megatron.core.fp4_utils import is_nvfp4tensor @@ -100,7 +101,9 @@ def model_provider( share_embeddings_and_output_weights=not args.untie_embeddings_and_output_weights, position_embedding_type=args.position_embedding_type, rotary_percent=args.rotary_percent, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) def create_test_args( self, tp, sequence_length, micro_batch_size, inference, fp4_param_gather, **kwargs diff --git a/tests/unit_tests/test_fp8_param.py b/tests/unit_tests/test_fp8_param.py index a8bb324f373..713a565fe5e 100644 --- a/tests/unit_tests/test_fp8_param.py +++ b/tests/unit_tests/test_fp8_param.py @@ -107,7 +107,9 @@ def model_provider( share_embeddings_and_output_weights=not args.untie_embeddings_and_output_weights, position_embedding_type=args.position_embedding_type, rotary_percent=args.rotary_percent, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) def create_test_args( self, diff --git a/tests/unit_tests/training/test_param_norm.py b/tests/unit_tests/training/test_param_norm.py index 12dbfba6992..0bd0e0f7143 100644 --- a/tests/unit_tests/training/test_param_norm.py +++ b/tests/unit_tests/training/test_param_norm.py @@ -6,6 +6,7 @@ import pytest import torch +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core.distributed import DistributedDataParallel, DistributedDataParallelConfig from megatron.core.models.gpt.gpt_layer_specs import get_gpt_layer_with_transformer_engine_spec from megatron.core.models.gpt.gpt_model import GPTModel @@ -56,7 +57,9 @@ def _build_tiny_moe_gpt( vocab_size=16, max_sequence_length=8, position_embedding_type="rope", - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) if not add_bias_linear: assert any(".shared_experts." in name for name, _ in model.named_parameters()) return model.cuda() diff --git a/tests/unit_tests/transformer/moe/test_upcycling.py b/tests/unit_tests/transformer/moe/test_upcycling.py index feb9c9b9d2f..4cbd906c908 100644 --- a/tests/unit_tests/transformer/moe/test_upcycling.py +++ b/tests/unit_tests/transformer/moe/test_upcycling.py @@ -84,7 +84,9 @@ def model_provider( share_embeddings_and_output_weights=not args.untie_embeddings_and_output_weights, position_embedding_type=args.position_embedding_type, rotary_percent=args.rotary_percent, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) return model diff --git a/tests/unit_tests/transformer/test_cuda_graphs.py b/tests/unit_tests/transformer/test_cuda_graphs.py index f387e58953e..9794ddec62a 100644 --- a/tests/unit_tests/transformer/test_cuda_graphs.py +++ b/tests/unit_tests/transformer/test_cuda_graphs.py @@ -746,7 +746,9 @@ def test_cuda_graph_determine_first_last_layer_logic( max_sequence_length=1024, position_embedding_type="rope", vp_stage=i, - ).cuda() + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ).cuda() model.append(this_model) # create runner by running a fake forward pass @@ -1136,7 +1138,9 @@ def test_get_cuda_graph_input_data(self, num_microbatches, pp_size, vpp_size): parallel_output=True, position_embedding_type="rope", vp_stage=i if vpp_size else None, - ).cuda() + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ).cuda() model.append(this_model) # Initialize TECudaGraphHelper @@ -1399,7 +1403,9 @@ def model_provider( position_embedding_type=args.position_embedding_type, rotary_percent=args.rotary_percent, mtp_block_spec=mtp_block_spec, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) def create_test_args( self, cuda_graph_impl, cuda_graph_modules, cuda_graph_warmup_steps, ep_size, **kwargs diff --git a/tests/unit_tests/transformer/test_multi_latent_attention.py b/tests/unit_tests/transformer/test_multi_latent_attention.py index 95517f7176d..15786e82b90 100644 --- a/tests/unit_tests/transformer/test_multi_latent_attention.py +++ b/tests/unit_tests/transformer/test_multi_latent_attention.py @@ -8,6 +8,7 @@ import torch import megatron.core.transformer.multi_latent_attention as mla_module +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core import parallel_state from megatron.core.extensions.transformer_engine_spec_provider import TESpecProvider from megatron.core.models.common.embeddings.rope_utils import ( @@ -1488,7 +1489,9 @@ def initialize_gpt_model( pre_process=pre_process, post_process=post_process, vp_stage=vp_stage, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) return gpt_model # Initialize baseline parallel state diff --git a/tests/unit_tests/transformer/test_multi_token_prediction.py b/tests/unit_tests/transformer/test_multi_token_prediction.py index c3c3944e007..5f60596eabb 100644 --- a/tests/unit_tests/transformer/test_multi_token_prediction.py +++ b/tests/unit_tests/transformer/test_multi_token_prediction.py @@ -463,7 +463,9 @@ def model_provider( share_embeddings_and_output_weights=not args.untie_embeddings_and_output_weights, position_embedding_type=args.position_embedding_type, rotary_percent=args.rotary_percent, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) return model @@ -1299,7 +1301,9 @@ def model_provider(self, pre_process=True, post_process=True, **config_kwargs): share_embeddings_and_output_weights=not args.untie_embeddings_and_output_weights, position_embedding_type=args.position_embedding_type, rotary_percent=args.rotary_percent, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) return model def create_test_args( diff --git a/tests/unit_tests/transformer/test_transformer_block.py b/tests/unit_tests/transformer/test_transformer_block.py index 0f19bc3dc95..659e358b5c8 100644 --- a/tests/unit_tests/transformer/test_transformer_block.py +++ b/tests/unit_tests/transformer/test_transformer_block.py @@ -861,7 +861,9 @@ def test_layout_layer_number(self, pipeline_model_parallel_layout, layer_number_ pre_process=pre_process, post_process=post_process, vp_stage=i, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) this_model.model_type = ModelType.encoder_or_decoder gpt_model.append(this_model) diff --git a/tests/unit_tests/transformer/test_utils.py b/tests/unit_tests/transformer/test_utils.py index 481a6bf5d61..cc34e7c248e 100644 --- a/tests/unit_tests/transformer/test_utils.py +++ b/tests/unit_tests/transformer/test_utils.py @@ -7,6 +7,7 @@ import torch import megatron.core.transformer.utils as transformer_utils +from megatron.core.process_groups_config import ProcessGroupCollection from megatron.core.models.gpt.gpt_layer_specs import get_gpt_layer_with_transformer_engine_spec from megatron.core.models.gpt.gpt_model import GPTModel from megatron.core.tensor_parallel.random import model_parallel_cuda_manual_seed @@ -44,7 +45,9 @@ def setup_method(self, method): max_sequence_length=8, position_embedding_type="rope", parallel_output=False, - ) + + pg_collection=ProcessGroupCollection.use_mpu_process_groups(), + ) def teardown_method(self, method): Utils.destroy_model_parallel()