From 66ce2fa5b76e45cd40ac0394da8e7072f7ac36c2 Mon Sep 17 00:00:00 2001 From: Erwin Huizenga Date: Wed, 23 Jul 2025 19:15:09 +0000 Subject: [PATCH] feat: Add sample for Vertex distributed training (#4163) * feat: Add sample for Vertex distributed training * refactor: Move distributed training to community content and add job config * fix: Address review comments and update files * minor fixes in the script * updated codeowners --- community-content/CODEOWNERS | 1 + .../llama-3-8b-nemo-pretraining/README.md | 126 +++++++++ .../configs/llama3_1_8b_pretrain_a3mega.yaml | 265 ++++++++++++++++++ .../docker/cloudbuild.yml | 26 ++ .../docker/patches/24.09/gpu_stats.patch | 41 +++ .../docker/patches/24.09/local_rank.patch | 41 +++ .../docker/patches/24.09/nemo2hf.patch | 13 + .../docker/patches/24.09/sigabort.patch | 24 ++ .../patches/24.09/throughput_calc.patch | 13 + .../docker/requirements.txt | 10 + .../docker/uninstall.txt | 18 ++ .../docker/vertex-dist-recipes.Dockerfile | 66 +++++ .../job_config.json | 16 ++ .../requirements.txt | 49 ++++ .../scripts/launch.py | 173 ++++++++++++ .../scripts/run.py | 85 ++++++ .../scripts/util/__init__.py | 0 .../scripts/util/cluster_spec.py | 81 ++++++ .../scripts/util/cluster_spec_test.py | 59 ++++ 19 files changed, 1107 insertions(+) create mode 100644 community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/README.md create mode 100644 community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/configs/llama3_1_8b_pretrain_a3mega.yaml create mode 100644 community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/cloudbuild.yml create mode 100644 community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/patches/24.09/gpu_stats.patch create mode 100644 community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/patches/24.09/local_rank.patch create mode 100644 community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/patches/24.09/nemo2hf.patch create mode 100644 community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/patches/24.09/sigabort.patch create mode 100644 community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/patches/24.09/throughput_calc.patch create mode 100644 community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/requirements.txt create mode 100644 community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/uninstall.txt create mode 100644 community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/vertex-dist-recipes.Dockerfile create mode 100644 community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/job_config.json create mode 100644 community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/requirements.txt create mode 100644 community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/scripts/launch.py create mode 100644 community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/scripts/run.py create mode 100644 community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/scripts/util/__init__.py create mode 100644 community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/scripts/util/cluster_spec.py create mode 100644 community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/scripts/util/cluster_spec_test.py diff --git a/community-content/CODEOWNERS b/community-content/CODEOWNERS index e4c72f6f9..842344e4c 100644 --- a/community-content/CODEOWNERS +++ b/community-content/CODEOWNERS @@ -29,4 +29,5 @@ /vertex_model_garden/model_oss/vllm @kathyyu-google /vertex_model_garden/benchmarking_reports @lavraicse /vertex_model_garden/model_oss/autogluon @lavraicse +/vertex_distributed_training/a3mega/llama-3-8b-nemo-pretraining @mstyer-google @erwinh85 @mchrestkha diff --git a/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/README.md b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/README.md new file mode 100644 index 000000000..759410695 --- /dev/null +++ b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/README.md @@ -0,0 +1,126 @@ +# Vertex AI Training: Llama 3.1 8B pre-training using Nvidia A3 Mega VMs (H100) +This document provides a step-by-step guide for pre-training a Llama 3.1 8B model on the `en-wiki` dataset using multiple [Vertex AI Custom Training](https://cloud.google.com/vertex-ai/docs/training/overview) `a3-megagpu-8g` nodes. + +We will use a custom container based on NVIDIA's [NeMo Framework](https://docs.nvidia.com/nemo-framework/user-guide/24.07/overview.html) to demonstrate a scalable, multi-node training workflow. All required artifacts and commands are included. + +## 1. Prerequisites + +### 1.1. Google Cloud Project setup +- **Enable APIs:** Ensure the Vertex AI API is [enabled for your project](http://console.cloud.google.com/flows/enableapi?apiid=aiplatform.googleapis.com). +- **H100 Mega Quota:** A3 Mega VMs are powered by H100 GPUs. Request quota for `custom_model_training_nvidia_h100_mega_gpus` in one of the [supported regions](https://cloud.google.com/vertex-ai/docs/general/locations#accelerator_support). If using Spot VMs, request `custom_model_training_preemptible_nvidia_h100_mega_gpus` quota instead. +- **Reservations (Optional but recommended):** For guaranteed capacity, [create a reservation](https://cloud.google.com/compute/docs/instances/reservations-shared) and ensure the reservation is shared with the Vertex AI service account. This guide requires a minimum of **16 H100 GPUs** (2 full A3 Mega nodes). + +### 1.2. GCS bucket +Create a [Cloud Storage bucket](https://cloud.google.com/storage/docs/creating-buckets) in the same region where you have quota. If you're using Hierarchical Namespace for your bucket, you may need to update permissions of the Vertex AI Custom Code Service Agent . + +This bucket is used for: +- Staging the training application. +- Storing model checkpoints and logs. +- Storing data if you use your own data. + + +## 2. Setup & configuration + +### 2.1. Clone the repo +First clone the repo into your development environment. + +```bash +git clone https://github.com/GoogleCloudPlatform/vertex-ai-samples.git +``` + +Navigate to the root folder for this sample. + +### 2.2. Environment Setup +First, configure your local environment. These variables are used in subsequent commands. + +```bash +# Required: Update with your values +export PROJECT_ID="" +export REPOSITORY="" # e.g., "my-containers" +export BUCKET="" + +# Optional: Change if needed +export REGION="us-central1" + +# --- Do not change the lines below --- +export ARTIFACT_REGISTRY="${REGION}-docker.pkg.dev/${PROJECT_ID}/${REPOSITORY}" +export REPO_ROOT=$(git rev-parse --show-toplevel) +``` + +## 3. Build and push a docker container image to Artifact Registry +Normally, you can use any custom training container on Vertex AI Training. In this example you build a NeMo Docker image that is based on the [Nvidia’s NeMo 24.09](https://catalog.ngc.nvidia.com/orgs/nvidia/containers/nemo/tags) image. Use Cloud Build to build and push the container image. + +This document picked NeMo as the demonstrating container since it’s a widely adopted GPU LLM training framework providing high performance and versatile training functionalities. + +In addition to the base image, some customizations are included to form the final prebuilt image: +- Some dependencies are installed to integrate with Vertex AI Training. +- An entrypoint script that sets up required environments and calls the training job. +- Some patches are applied to the NeMo code to let it load the dataset from a GCS bucket. + +Run this command to build the container and push the container into the Google Artifact Registry. + +```bash +cd "${REPO_ROOT}/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining" +export IMAGE_NAME="vertex-nemo-llama" +gcloud builds submit . \ + --project="${PROJECT_ID}" \ + --region="${REGION}" \ + --config=docker/cloudbuild.yml \ + --substitutions="_ARTIFACT_REGISTRY=${ARTIFACT_REGISTRY},_IMAGE_NAME=${IMAGE_NAME}" \ + --timeout="2h" \ + --machine-type="e2-highcpu-32" +``` + +## 4. Launch the Training Job + + +### 4.1. Job Configuration File +Once the container is built, update the job_config.json to set up the training job. +File: job_config.json +```json +{ + "project_id": "", + "region": "", + "zone": "", + "bucket": "", + "dataset_bucket": "github-repo/data/third-party/enwiki-latest-pages-articles", + "image_uri": "", + "strategy": "spot", + "nodes": "2", + "machine_type": "a3-megagpu-8g", + "gpu_type": "NVIDIA_H100_MEGA_80GB", + "gpus_per_node": "8", + "recipe_name": "llama3_1_8b_pretrain_a3mega", + "job_prefix": "vertex-spot-", + "reservation_name": "" +} +``` + +### 4.2 Launch the Training Job + +First, create a Python virtual environment using your tool of choice, then install +the requirements specified in `requirements.txt`. Using `pip`, the command would be: +```bash +pip install -r requirements.txt +``` + +Now launch the Vertex AI training job using the provided Python script. + +```bash +python3 scripts/launch.py --config_file=job_config.json +``` + +This script reads job_config.json, defines the cluster specification (2 nodes, 8 GPUs each), and submits the custom training job to Vertex AI. + +## 5. Monitor and Clean Up + +### 5.1. Monitoring +Vertex AI Console: Track the job's status in the Google Cloud Console under Vertex AI > Training > Custom Jobs. +Logs: View detailed logs in Cloud Logging by filtering for your job name. +Checkpoints: Model checkpoints are saved to your GCS bucket at the path specified in your training script's configuration. + +### 5.2. Cleaning Up +To avoid ongoing charges, delete the resources you created: +- The Artifact Registry image. +- The contents of the GCS bucket (checkpoints, logs). +- The Vertex AI Custom Job will eventually complete or fail, incurring no further cost. diff --git a/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/configs/llama3_1_8b_pretrain_a3mega.yaml b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/configs/llama3_1_8b_pretrain_a3mega.yaml new file mode 100644 index 000000000..c048fef52 --- /dev/null +++ b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/configs/llama3_1_8b_pretrain_a3mega.yaml @@ -0,0 +1,265 @@ +# Reference: +# https://github.com/NVIDIA/NeMo-Framework-Launcher/blob/24.07/launcher_scripts/conf/training/llama/llama3_1_8b.yaml +name: llama3_1_8b_pretrain_a3mega +restore_from_path: null # used when starting from a .nemo file + +trainer: + devices: 8 + num_nodes: 1 + accelerator: gpu + precision: bf16 + logger: false # logger provided by exp_manager + enable_checkpointing: false + use_distributed_sampler: false + max_epochs: -1 # PTL default. In practice, max_steps will be reached first. + max_steps: 30 # consumed_samples = global_step * micro_batch_size * data_parallel_size * accumulate_grad_batches + log_every_n_steps: 1 + val_check_interval: null + limit_val_batches: 1 + limit_test_batches: 1 + accumulate_grad_batches: 1 # do not modify, grad acc is automatic for training megatron models + gradient_clip_val: 1.0 + benchmark: false + enable_model_summary: false # default PTL callback for this does not support model parallelism, instead we log manually + +exp_manager: + explicit_log_dir: null + exp_dir: /data + name: ${name} + create_dllogger_logger: true + dllogger_logger_kwargs: + verbose: true + stdout: true + json_file: "/data/dllogger.json" + create_wandb_logger: false + wandb_logger_kwargs: + project: null + name: null + resume_if_exists: true + resume_ignore_no_checkpoint: true + create_checkpoint_callback: false + checkpoint_callback_params: + monitor: val_loss + save_top_k: 3 + mode: min + always_save_nemo: false # saves nemo file during validation, not implemented for model parallel + save_nemo_on_train_end: false # not recommended when training large models on clusters with short time limits + filename: 'megatron_gpt--{val_loss:.2f}-{step}-{consumed_samples}' + model_parallel_size: ${multiply:${model.tensor_model_parallel_size}, ${model.pipeline_model_parallel_size}} + seconds_to_sleep: 5 # Allows node_rank!=0 to sleep and let node0 to init, like preparing data + +model: + mcore_gpt: true + # specify micro_batch_size, global_batch_size, and model parallelism + # gradient accumulation will be done automatically based on data_parallel_size + micro_batch_size: 1 # limited by GPU memory + global_batch_size: 1024 # will use more micro batches to reach global batch size + tensor_model_parallel_size: 1 # intra-layer model parallelism + pipeline_model_parallel_size: 2 # inter-layer model parallelism + context_parallel_size: 1 + virtual_pipeline_model_parallel_size: null # interleaved pipeline + ## Sequence Parallelism + # Makes tensor parallelism more memory efficient for LLMs (20B+) by parallelizing layer norms and dropout sequentially + # See Reducing Activation Recomputation in Large Transformer Models: https://arxiv.org/abs/2205.05198 for more details. + sequence_parallel: false + + fsdp: false + fsdp_cpu_offload: true + fsdp_sharding_strategy: "full" # Method to shard model states. Available options are 'full', 'hybrid', and 'grad'. + fsdp_grad_reduce_dtype: "16" # Gradient reduction data type. + fsdp_sharded_checkpoint: false # Store and load FSDP shared checkpoint. + fsdp_use_orig_params: false # Set to True to use FSDP for specific peft scheme. + + # Distributed checkpoint setup + dist_ckpt_format: "torch_dist" # Set to 'torch_dist' to use PyTorch distributed checkpoint format. + dist_ckpt_load_on_device: true # whether to load checkpoint weights directly on GPU or to CPU + dist_ckpt_parallel_save: true # if true, each worker will write its own part of the dist checkpoint + dist_ckpt_parallel_save_within_dp: false # if true, save will be parallelized only within a DP group (whole world otherwise), which might slightly reduce the save overhead + dist_ckpt_parallel_load: false # if true, each worker will load part of the dist checkpoint and exchange with NCCL. Might use some extra GPU memory + dist_ckpt_torch_dist_multiproc: 2 # number of extra processes per rank used during ckpt save with PyTorch distributed format + dist_ckpt_assume_constant_structure: false # set to True only if the state dict structure doesn't change within a single job. Allows caching some computation across checkpoint saves. + dist_ckpt_parallel_dist_opt: true # parallel save/load of a DistributedOptimizer. 'True' allows performant save and reshardable checkpoints. Set to 'False' only in order to minimize the number of checkpoint files. + dist_ckpt_load_strictness: null # defines checkpoint keys mismatch behavior (only during dist-ckpt load). Choices: assume_ok_unexpected (default - try loading without any check), log_all (log mismatches), raise_all (raise mismatches) + + # model architecture + encoder_seq_length: 8192 + max_position_embeddings: ${.encoder_seq_length} + num_layers: 32 # 8b: 32 | 70b: 80 | 405b: 126 + hidden_size: 4096 # 8b: 4096 | 70b: 8192 | 405b: 16384 + ffn_hidden_size: 14336 # 8b: 14336 | 70b: 28672 | 405b: 53248 + num_attention_heads: 32 # 8b: 32 | 70b: 64 | 405b: 128 + num_query_groups: 8 # Number of query groups for group query attention. If None, normal attention is used. 8b: 8 | 70b: 8 | 405b: 16 + init_method_std: 0.01 # Standard deviation of the zero mean normal distribution used for weight initialization. 8b: 0.01 | 70b: 0.008944 | 405b: 0.02 + use_scaled_init_method: true # use scaled residuals initialization + hidden_dropout: 0.0 # Dropout probability for hidden state transformer. + attention_dropout: 0.0 # Dropout probability for attention + ffn_dropout: 0.0 # Dropout probability in the feed-forward layer. + kv_channels: null # Projection weights dimension in multi-head attention. Set to hidden_size // num_attention_heads if null + apply_query_key_layer_scaling: true # scale Q * K^T by 1 / layer-number. + normalization: 'rmsnorm' # Normalization layer to use. Options are 'layernorm', 'rmsnorm' + layernorm_epsilon: 1e-5 + do_layer_norm_weight_decay: false # True means weight decay on all params + make_vocab_size_divisible_by: 128 # Pad the vocab size to be divisible by this value for computation efficiency. + pre_process: true # add embedding + post_process: true # add pooler + persist_layer_norm: true # Use of persistent fused layer norm kernel. + bias: false # Whether to use bias terms in all weight matrices. + activation: 'fast-swiglu' # Options ['gelu', 'geglu', 'swiglu', 'reglu', 'squared-relu', 'fast-geglu', 'fast-swiglu', 'fast-reglu'] + headscale: false # Whether to learn extra parameters that scale the output of the each self-attention head. + transformer_block_type: 'pre_ln' # Options ['pre_ln', 'post_ln', 'normformer'] + openai_gelu: false # Use OpenAI's GELU instead of the default GeLU + normalize_attention_scores: true # Whether to scale the output Q * K^T by 1 / sqrt(hidden_size_per_head). This arg is provided as a configuration option mostly for compatibility with models that have been weight-converted from HF. You almost always want to se this to True. + position_embedding_type: 'rope' # Position embedding type. Options ['learned_absolute', 'rope'] + rotary_percentage: 1.0 # If using position_embedding_type=rope, then the per head dim is multiplied by this. + attention_type: 'multihead' # Attention type. Options ['multihead'] + share_embeddings_and_output_weights: false # Share embedding and output layer weights. + scale_positional_embedding: true # This is false for llama3 models. Only used for >= llama3.1. + + # Use GPT2BPETokenizer for test, because the testing dataset is tokenized by this tokenizer. + # https://docs.nvidia.com/nemo-framework/user-guide/24.07/playbooks/singlenodepretrain.html#data-download-and-pre-processing + tokenizer: + library: megatron + type: GPT2BPETokenizer + model: null # /path/to/tokenizer.model + vocab_file: null + merge_file: null + delimiter: null # only used for tabular tokenizer + sentencepiece_legacy: false # Legacy=True allows you to add special tokens to sentencepiece tokenizers. + + # Mixed precision + native_amp_init_scale: 4294967296 # 2 ** 32 + native_amp_growth_interval: 1000 + hysteresis: 2 # Gradient scale hysteresis + fp32_residual_connection: false # Move residual connections to fp32 + fp16_lm_cross_entropy: false # Move the cross entropy unreduced loss calculation for lm head to fp16 + + # Megatron O2-style half-precision + megatron_amp_O2: true # Enable O2-level automatic mixed precision using main parameters + grad_allreduce_chunk_size_mb: 125 + + # Fusion + grad_div_ar_fusion: true # Fuse grad division into torch.distributed.all_reduce. Only used with O2 and no pipeline parallelism.. + gradient_accumulation_fusion: true # Fuse weight gradient accumulation to GEMMs. Only used with pipeline parallelism and O2. + bias_activation_fusion: true # Use a kernel that fuses the bias addition from weight matrices with the subsequent activation function. + bias_dropout_add_fusion: true # Use a kernel that fuses the bias addition, dropout and residual connection addition. + masked_softmax_fusion: true # Use a kernel that fuses the attention softmax with it's mask. + apply_rope_fusion: true # Use a kernel to add rotary positional embeddings. Only used if position_embedding_type=rope + cross_entropy_loss_fusion: true + + # Miscellaneous + seed: 1234 + resume_from_checkpoint: null # manually set the checkpoint file to load from + use_cpu_initialization: false # Init weights on the CPU (slow for large models) + onnx_safe: false # Use work-arounds for known problems with Torch ONNX exporter. + apex_transformer_log_level: 30 # Python logging level displays logs with severity greater than or equal to this + gradient_as_bucket_view: true # PyTorch DDP argument. Allocate gradients in a contiguous bucket to save memory (less fragmentation and buffer memory) + sync_batch_comm: false # Enable stream synchronization after each p2p communication between pipeline stages + + ## Activation Checkpointing + # NeMo Megatron supports 'selective' activation checkpointing where only the memory intensive part of attention is checkpointed. + # These memory intensive activations are also less compute intensive which makes activation checkpointing more efficient for LLMs (20B+). + # See Reducing Activation Recomputation in Large Transformer Models: https://arxiv.org/abs/2205.05198 for more details. + # 'full' will checkpoint the entire transformer layer. + activations_checkpoint_granularity: null # 'selective' or 'full' + activations_checkpoint_method: null # 'uniform', 'block' + # 'uniform' divides the total number of transformer layers and checkpoints the input activation + # of each chunk at the specified granularity. When used with 'selective', 'uniform' checkpoints all attention blocks in the model. + # 'block' checkpoints the specified number of layers per pipeline stage at the specified granularity + activations_checkpoint_num_layers: null + # when using 'uniform' this creates groups of transformer layers to checkpoint. Usually set to 1. Increase to save more memory. + # when using 'block' this this will checkpoint the first activations_checkpoint_num_layers per pipeline stage. + num_micro_batches_with_partial_activation_checkpoints: null + # This feature is valid only when used with pipeline-model-parallelism. + # When an integer value is provided, it sets the number of micro-batches where only a partial number of Transformer layers get checkpointed + # and recomputed within a window of micro-batches. The rest of micro-batches in the window checkpoint all Transformer layers. The size of window is + # set by the maximum outstanding micro-batch backpropagations, which varies at different pipeline stages. The number of partial layers to checkpoint + # per micro-batch is set by 'activations_checkpoint_num_layers' with 'activations_checkpoint_method' of 'block'. + # This feature enables using activation checkpoint at a fraction of micro-batches up to the point of full GPU memory usage. + activations_checkpoint_layers_per_pipeline: null + # This feature is valid only when used with pipeline-model-parallelism. + # When an integer value (rounded down when float is given) is provided, it sets the number of Transformer layers to skip checkpointing at later + # pipeline stages. For example, 'activations_checkpoint_layers_per_pipeline' of 3 makes pipeline stage 1 to checkpoint 3 layers less than + # stage 0 and stage 2 to checkpoint 6 layers less stage 0, and so on. This is possible because later pipeline stage + # uses less GPU memory with fewer outstanding micro-batch backpropagations. Used with 'num_micro_batches_with_partial_activation_checkpoints', + # this feature removes most of activation checkpoints at the last pipeline stage, which is the critical execution path. + + ## Transformer Engine + transformer_engine: true + fp8: false # enables fp8 in TransformerLayer forward + fp8_e4m3: false # sets fp8_format = recipe.Format.E4M3 + fp8_hybrid: false # sets fp8_format = recipe.Format.HYBRID + fp8_margin: 0 # scaling margin + fp8_interval: 1 # scaling update interval + fp8_amax_history_len: 1024 # Number of steps for which amax history is recorded per tensor + fp8_amax_compute_algo: 'max' # 'most_recent' or 'max'. Algorithm for computing amax from history + ub_tp_comm_overlap: false # do not turn on because of b/397797926 + use_flash_attention: true + gc_interval: 100 + + ## Offloading Activations/Weights to CPU + cpu_offloading: false + cpu_offloading_num_layers: ${sum:${.num_layers},-1} # This value should be between [1,num_layers-1] as we don't want to offload the final layer's activations and expose any offloading duration for the final layer + cpu_offloading_activations: true + cpu_offloading_weights: true + + data: + # Path to data must be specified by the user. + # Supports List, String and Dictionary + # List : can override from the CLI: "model.data.data_prefix=[.5,/raid/data/pile/my-gpt3_00_text_document,.5,/raid/data/pile/my-gpt3_01_text_document]", + # Or see example below: + # data_prefix: + # - .5 + # - /raid/data/pile/my-gpt3_00_text_document + # - .5 + # - /raid/data/pile/my-gpt3_01_text_document + # Dictionary: can override from CLI "model.data.data_prefix"={"train":[1.0, /path/to/data], "validation":/path/to/data, "test":/path/to/test} + # Or see example below: + # "model.data.data_prefix: {train:[1.0,/path/to/data], validation:[/path/to/data], test:[/path/to/test]}" + data_prefix: [1.0, /data/hfbpe_gpt_training_data_text_document] + index_mapping_dir: null # path to save index mapping .npy files, by default will save in the same location as data_prefix + data_impl: mmap + splits_string: 900,50,50 + seq_length: ${model.encoder_seq_length} + skip_warmup: true + num_workers: 2 + dataloader_type: single # cyclic + reset_position_ids: false # Reset position ids after end-of-document token + reset_attention_mask: false # Reset attention mask after end-of-document token + eod_mask_loss: false # Mask loss for the end of document tokens + validation_drop_last: true # Set to false if the last partial validation samples is to be consumed + no_seqlen_plus_one_input_tokens: false # Set to True to disable fetching (sequence length + 1) input tokens, instead get (sequence length) input tokens and mask the last token + pad_samples_to_global_batch_size: false # Set to True if you want to pad the last partial batch with -1's to equal global batch size + shuffle_documents: true # Set to False to disable documents shuffling. Sample index will still be shuffled + + # Nsys profiling options + nsys_profile: + enabled: false + start_step: 0 # Global batch to start profiling + end_step: 1 # Global batch to end profiling + ranks: [0] # Global rank IDs to profile + gen_shape: false # Generate model and kernel details including input shapes + + memory_profile: + enabled: false + start_step: 0 + end_step: 1 + ranks: [0] + output_path: /data # Must be a dir + + optim: + name: distributed_fused_adam # E.g., fused_adam or set _target_: torch.optim.AdamW field + lr: 2e-5 + weight_decay: 0.01 + betas: + - 0.9 + - 0.98 + bucket_cap_mb: 125 + overlap_grad_sync: true + overlap_param_sync: true + contiguous_grad_buffer: true + contiguous_param_buffer: true + sched: + name: CosineAnnealing + warmup_steps: 400 + constant_steps: 0 + min_lr: 2e-6 diff --git a/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/cloudbuild.yml b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/cloudbuild.yml new file mode 100644 index 000000000..ba5023960 --- /dev/null +++ b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/cloudbuild.yml @@ -0,0 +1,26 @@ +# Copyright 2024 Google LLC +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +steps: +- name: 'gcr.io/cloud-builders/docker' + args: + - 'build' + - '--tag=${_ARTIFACT_REGISTRY}/${_IMAGE_NAME}' + - '--file=docker/vertex-dist-recipes.Dockerfile' + - '.' + automapSubstitutions: true + env: + - 'DOCKER_BUILDKIT=1' +images: +- '${_ARTIFACT_REGISTRY}/${_IMAGE_NAME}' diff --git a/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/patches/24.09/gpu_stats.patch b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/patches/24.09/gpu_stats.patch new file mode 100644 index 000000000..531a3d6da --- /dev/null +++ b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/patches/24.09/gpu_stats.patch @@ -0,0 +1,41 @@ +diff --git a/nemo/collections/nlp/parts/megatron_trainer_builder.py b/nemo/collections/nlp/parts/megatron_trainer_builder.py +index b2c85cde4..a3a9670c3 100644 +--- a/nemo/collections/nlp/parts/megatron_trainer_builder.py ++++ b/nemo/collections/nlp/parts/megatron_trainer_builder.py +@@ -19,6 +19,7 @@ from lightning_fabric.utilities.exceptions import MisconfigurationException + from omegaconf import DictConfig + from pytorch_lightning import Trainer + from pytorch_lightning.callbacks import ModelSummary ++from pytorch_lightning.callbacks import Callback + from pytorch_lightning.plugins.environments import TorchElasticEnvironment + + from nemo.collections.common.metrics.perf_metrics import FLOPsMeasurementCallback +@@ -38,6 +39,23 @@ from nemo.utils.callbacks.dist_ckpt_io import ( + AsyncFinalizerCallback, + DistributedCheckpointIO, + ) ++from vmg.util.device_stats import gpu_stats_str ++ ++class GpuStatsMon(Callback): ++ def on_train_start(self, trainer, pl_module) -> None: ++ rank=pl_module.global_rank ++ print(f'train_start: {rank=} {gpu_stats_str()}', flush=True) ++ ++ def on_train_batch_start(self, trainer, pl_module, batch, batch_idx) -> None: ++ rank=pl_module.global_rank ++ print(f'batch_start: {rank=} {gpu_stats_str()}', flush=True) ++ ++ def on_train_batch_end(self, trainer, pl_module, outputs, batch, batch_idx) -> None: ++ rank=pl_module.global_rank ++ print(f'batch_end: {rank=} {gpu_stats_str()}', flush=True) + + + class MegatronTrainerBuilder: +@@ -178,6 +196,7 @@ class MegatronTrainerBuilder: + if self.cfg.get('exp_manager', {}).get('log_tflops_per_sec_per_gpu', True): + callbacks.append(FLOPsMeasurementCallback(self.cfg)) + ++ callbacks.append(GpuStatsMon()) + return callbacks + + def create_trainer(self, callbacks=None) -> Trainer: diff --git a/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/patches/24.09/local_rank.patch b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/patches/24.09/local_rank.patch new file mode 100644 index 000000000..2ac2e8a20 --- /dev/null +++ b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/patches/24.09/local_rank.patch @@ -0,0 +1,41 @@ +diff -ruN old-datasets/blended_megatron_dataset_builder.py datasets/blended_megatron_dataset_builder.py +--- old-datasets/blended_megatron_dataset_builder.py 2025-05-02 04:08:45.369199665 +0000 ++++ datasets/blended_megatron_dataset_builder.py 2025-05-02 04:10:47.369119891 +0000 +@@ -2,6 +2,7 @@ + + import logging + import math ++import os + from concurrent.futures import ThreadPoolExecutor + from typing import Any, Callable, Iterable, List, Optional, Type, Union + +@@ -353,7 +354,7 @@ + num_dataset_builder_threads = self.config.num_dataset_builder_threads + + if torch.distributed.is_initialized(): +- rank = torch.distributed.get_rank() ++ rank = int(os.getenv("LOCAL_RANK", "0")) + # First, build on rank 0 + if rank == 0: + num_workers = num_dataset_builder_threads +@@ -475,7 +476,7 @@ + Optional[Union[DistributedDataset, Iterable]]: The DistributedDataset instantion, the Iterable instantiation, or None + """ + if torch.distributed.is_initialized(): +- rank = torch.distributed.get_rank() ++ rank = int(os.getenv("LOCAL_RANK", "0")) + + dataset = None + +diff -ruN old-datasets/gpt_dataset.py datasets/gpt_dataset.py +--- old-datasets/gpt_dataset.py 2025-05-02 04:08:45.369199665 +0000 ++++ datasets/gpt_dataset.py 2025-05-02 04:09:30.309170278 +0000 +@@ -351,7 +351,7 @@ + + if not path_to_cache or ( + not cache_hit +- and (not torch.distributed.is_initialized() or torch.distributed.get_rank() == 0) ++ and (not torch.distributed.is_initialized() or int(os.getenv("LOCAL_RANK", "0")) == 0) + ): + + log_single_rank( diff --git a/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/patches/24.09/nemo2hf.patch b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/patches/24.09/nemo2hf.patch new file mode 100644 index 000000000..569c7e181 --- /dev/null +++ b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/patches/24.09/nemo2hf.patch @@ -0,0 +1,13 @@ +diff --git a/scripts/checkpoint_converters/convert_llama_nemo_to_hf.py b/scripts/checkpoint_converters/convert_llama_nemo_to_hf.py +index 8da15148d..005cae6c9 100644 +--- a/scripts/checkpoint_converters/convert_llama_nemo_to_hf.py ++++ b/scripts/checkpoint_converters/convert_llama_nemo_to_hf.py +@@ -104,6 +104,8 @@ def convert(input_nemo_file, output_hf_file, precision=None, cpu_only=False) -> + dummy_trainer = Trainer(devices=1, accelerator='cpu', strategy=NLPDDPStrategy()) + model_config = MegatronGPTModel.restore_from(input_nemo_file, trainer=dummy_trainer, return_config=True) + model_config.tensor_model_parallel_size = 1 ++ model_config.virtual_pipeline_model_parallel_size = None ++ model_config.sequence_parallel = False + model_config.pipeline_model_parallel_size = 1 + if cpu_only: + map_location = torch.device('cpu') diff --git a/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/patches/24.09/sigabort.patch b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/patches/24.09/sigabort.patch new file mode 100644 index 000000000..7d12080da --- /dev/null +++ b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/patches/24.09/sigabort.patch @@ -0,0 +1,24 @@ +diff --git a/examples/nlp/language_modeling/tuning/megatron_gpt_finetuning.py b/examples/nlp/language_modeling/tuning/megatron_gpt_finetuning.py +index bfe8ea359..dfeaf93b5 100644 +--- a/examples/nlp/language_modeling/tuning/megatron_gpt_finetuning.py ++++ b/examples/nlp/language_modeling/tuning/megatron_gpt_finetuning.py +@@ -13,6 +13,8 @@ + # limitations under the License. + + import torch.multiprocessing as mp ++import torch.distributed as dist ++ + from omegaconf.omegaconf import OmegaConf + + from nemo.collections.nlp.models.language_modeling.megatron_gpt_sft_model import MegatronGPTSFTModel +@@ -76,6 +78,10 @@ def main(cfg) -> None: + + trainer.fit(model) + ++ if dist.is_available() and dist.is_initialized(): ++ dist.barrier() ++ dist.destroy_process_group() ++ + + if __name__ == '__main__': + main() diff --git a/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/patches/24.09/throughput_calc.patch b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/patches/24.09/throughput_calc.patch new file mode 100644 index 000000000..ea8a620ae --- /dev/null +++ b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/patches/24.09/throughput_calc.patch @@ -0,0 +1,13 @@ +diff --git a/src/utils/training_metrics/process_training_results.py b/src/utils/training_metrics/process_training_results.py +index 3e82a66..e61e1d8 100644 +--- a/src/utils/training_metrics/process_training_results.py ++++ b/src/utils/training_metrics/process_training_results.py +@@ -134,7 +134,7 @@ def get_average_step_time(file: str, start_step: int, end_step: int) -> float: + for line in datajson: + if line.get("step") != "PARAMETER": + step = line.get("step") +- if step >= start_step and step <= end_step: ++ if step >= start_step and step <= end_step and "train_step_timing in s" in line["data"]: + time_step_accumulator += line["data"].get("train_step_timing in s") + num_steps += 1 + if num_steps == 0: diff --git a/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/requirements.txt b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/requirements.txt new file mode 100644 index 000000000..c178cabbc --- /dev/null +++ b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/requirements.txt @@ -0,0 +1,10 @@ +dllogger@git+https://github.com/NVIDIA/dllogger@v1.0.0 + +# Fixing these libraries versions to avoid conflicting or broken packages. +immutabledict==4.2.1 +protobuf==3.20.3 +opencv-python-headless==4.11.0.86 +docutils==0.16 +urllib3==2.0.7 +google-cloud-storage==3.0.0 +retrying diff --git a/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/uninstall.txt b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/uninstall.txt new file mode 100644 index 000000000..5ef99361a --- /dev/null +++ b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/uninstall.txt @@ -0,0 +1,18 @@ +# cuml-cu12==24.8.0 was installed in nemo:24.09 +# Removing cuml=24.4.0 to avoid conflicting packages. +cudf==24.4.0 +cugraph==24.4.0 +cugraph-service-server==24.4.0 +cuml==24.4.0 +dask-cudf==24.4.0 +raft-dask==24.4.0 +cugraph-dgl==24.4.0 +cugraph-pyg==24.4.0 +# The following packages are removed temporarily to avoid conflicting packages +# and can be brought back if needed. +tensorrt-llm==0.12.0 +img2dataset==1.45.0 +Sphinx==8.1.3 +sphinxcontrib-bibtex==2.6.3 +torchx==0.7.0 +nemo-run diff --git a/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/vertex-dist-recipes.Dockerfile b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/vertex-dist-recipes.Dockerfile new file mode 100644 index 000000000..1d7c0bce2 --- /dev/null +++ b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/docker/vertex-dist-recipes.Dockerfile @@ -0,0 +1,66 @@ +# Dockerfile wrapping NeMo. +# +# To workaround base nemo docker image using too many layers, we use Multi-stage +# build to first collect the additional files we'll need. +FROM alpine:latest AS prep_files +WORKDIR /workspace +RUN mkdir -p configs vdt vdt/util +COPY scripts/*.py vdt/ +COPY scripts/util/*.py vdt/util/ +COPY configs/* configs/ +COPY docker/patches/24.09/* vdt/patches/ +RUN chmod a+rwX -R vdt +# Copy license. +RUN wget https://raw.githubusercontent.com/GoogleCloudPlatform/vertex-ai-samples/main/LICENSE + +# Available tags +# https://catalog.ngc.nvidia.com/orgs/nvidia/containers/nemo/tags +# It installs NeMo source code in /opt/NeMo folder, with tag=r2.0.0 +FROM nvcr.io/nvidia/nemo:24.09 + +RUN apt-get update && apt-get install -y sudo zsh tmux && \ + rm -rf /var/lib/apt/lists* + +RUN echo "deb [signed-by=/usr/share/keyrings/cloud.google.gpg] http://packages.cloud.google.com/apt cloud-sdk main" | \ + tee -a /etc/apt/sources.list.d/google-cloud-sdk.list && \ + curl https://packages.cloud.google.com/apt/doc/apt-key.gpg | \ + apt-key --keyring /usr/share/keyrings/cloud.google.gpg add - && \ + apt-get update -y && apt-get install google-cloud-sdk -y && \ + rm -rf /var/lib/apt/lists* + +# Install libraries with pip +ENV PIP_ROOT_USER_ACTION=ignore + +# We expect this will be run in the root directory of the vertex-dist-recipes repo +ARG HOST_SRC_DIR="." + +# The pre-installed NeMo introduces a lot of deps conflicts. +# We uninstall the confilicting libs and reinstall some of them as needed. +COPY ${HOST_SRC_DIR}/docker/uninstall.txt /tmp/uninstall.txt +RUN cat /tmp/uninstall.txt | grep -v '#' | xargs pip uninstall -y +COPY ${HOST_SRC_DIR}/docker/requirements.txt /tmp/requirements.txt +RUN pip install -r /tmp/requirements.txt + +# Make sure there's no inconsistent pip libraries. +RUN pip check + +WORKDIR /workspace + +# Copy configs +COPY ${HOST_SRC_DIR}/configs/* /opt/NeMo/examples/nlp/language_modeling/conf/ + +# Copy all additional files we need from `prep_files` image. +COPY --from=prep_files /workspace/ . + +# Install for `src/utils/training_metrics/process_training_results.py` to report +# throughput and MFU numbers. +RUN git clone https://github.com/AI-Hypercomputer/gpu-recipes.git + +# This hack is needed for multi-node training while not using a sharing file system. +RUN patch --verbose -l -d /opt/megatron-lm/megatron/core/datasets -p1 -i /workspace/vdt/patches/local_rank.patch; \ + git -C /workspace/gpu-recipes apply /workspace/vdt/patches/throughput_calc.patch; \ + git -C /opt/NeMo apply /workspace/vdt/patches/nemo2hf.patch; \ + git -C /opt/NeMo apply /workspace/vdt/patches/sigabort.patch; +# git -C /opt/NeMo apply /workspace/vdt/patches/gpu_stats.patch; + +# Do not put an entrypoint here. Specify the entrypoint in the docker run script. diff --git a/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/job_config.json b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/job_config.json new file mode 100644 index 000000000..f02cf8eb9 --- /dev/null +++ b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/job_config.json @@ -0,0 +1,16 @@ +{ + "project_id": "", + "region": "us-central1", + "zone": "us-central1-c", + "bucket": "", + "strategy": "spot", + "nodes": "2", + "machine_type": "a3-megagpu-8g", + "gpu_type": "NVIDIA_H100_MEGA_80GB", + "gpus_per_node": "8", + "recipe_name": "llama3_1_8b_pretrain_a3mega", + "job_prefix": "vertex-ai", + "reservation_name": "" +} diff --git a/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/requirements.txt b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/requirements.txt new file mode 100644 index 000000000..f3ecc9fe7 --- /dev/null +++ b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/requirements.txt @@ -0,0 +1,49 @@ +absl-py==2.2.2 +annotated-types==0.7.0 +anyio==4.9.0 +black==25.1.0 +cachetools==5.5.2 +certifi==2025.4.26 +charset-normalizer==3.4.2 +click==8.1.8 +docstring_parser==0.16 +google-api-core==2.24.2 +google-auth==2.40.1 +google-cloud-aiplatform==1.92.0 +google-cloud-bigquery==3.31.0 +google-cloud-core==2.4.3 +google-cloud-resource-manager==1.14.2 +google-cloud-storage==2.19.0 +google-crc32c==1.7.1 +google-genai==1.14.0 +google-resumable-media==2.7.2 +googleapis-common-protos==1.70.0 +grpc-google-iam-v1==0.14.2 +grpcio==1.71.0 +grpcio-status==1.71.0 +h11==0.16.0 +httpcore==1.0.9 +httpx==0.28.1 +idna==3.10 +mypy_extensions==1.1.0 +numpy==2.2.5 +packaging==25.0 +pathspec==0.12.1 +platformdirs==4.3.8 +proto-plus==1.26.1 +protobuf==5.29.4 +pyasn1==0.6.1 +pyasn1_modules==0.4.2 +pydantic==2.11.4 +pydantic_core==2.33.2 +python-dateutil==2.9.0.post0 +pytz==2025.2 +requests==2.32.3 +rsa==4.9.1 +shapely==2.1.0 +six==1.17.0 +sniffio==1.3.1 +typing-inspection==0.4.0 +typing_extensions==4.13.2 +urllib3==2.4.0 +websockets==15.0.1 \ No newline at end of file diff --git a/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/scripts/launch.py b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/scripts/launch.py new file mode 100644 index 000000000..7f049c914 --- /dev/null +++ b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/scripts/launch.py @@ -0,0 +1,173 @@ +"""Launch script for Vertex distributed training""" + +# Copy the sample_job_config.json file to job_config.json +# to define the job parameters. +# +# Run like this: +# +# python3 vertex_dist_train/launch.py --config_file=job_config.json +# + +import datetime +import json +import os +import pprint +from collections.abc import Sequence +from typing import Any, List + +from absl import app, flags +from google.cloud import aiplatform +from google.cloud.aiplatform_v1.types.custom_job import Scheduling +from pytz import timezone + +FLAGS = flags.FLAGS +flags.DEFINE_string("config_file", None, "Path to JSON config file") +flags.DEFINE_boolean( + "debug", False, "Debug mode: just print the command, don't run it." +) + + +def launch_job( + job_name: str, + project: str, + region: str, + gcs_bucket: str, + image_uri: str, + entrypoint_cmd: List[str], + trainer_args: List[Any], + num_nodes: int, + machine_type: str, + num_gpus_per_node: int, + gpu_type: str, + strategy: str, + reservation_name: str = "", +): + assert strategy in ("dws", "spot", "reservation") + aiplatform.init( + project=project, location=region, staging_bucket=gcs_bucket + ) + + train_job = aiplatform.CustomContainerTrainingJob( + display_name=job_name, + container_uri=image_uri, + command=entrypoint_cmd, + ) + + job_args = dict( + args=trainer_args, + enable_web_access=True, + replica_count=num_nodes, + machine_type=machine_type, + accelerator_type=gpu_type, + accelerator_count=num_gpus_per_node, + boot_disk_size_gb=1000, + restart_job_on_worker_restart=True, + #restart_job_on_worker_restart=False, + ) + + if strategy == "spot": + job_args.update({"scheduling_strategy": Scheduling.Strategy.SPOT.name}) + elif strategy == "dws": + job_args.update( + {"scheduling_strategy": Scheduling.Strategy.FLEX_START.name} + ) + elif strategy == "reservation": + assert reservation_name != "", ( + "If using a reservation, provide the reservation_name in the " + "format `projects/{project_id_or_number}/zones/{zone}/" + "reservations/{reservation_name}`" + ) + job_args.update( + { + "reservation_affinity_type": "SPECIFIC_RESERVATION", + "reservation_affinity_key": "compute.googleapis.com/reservation-name", + "reservation_affinity_values": [reservation_name], + } + ) + + pprint.pprint(job_args) + if not FLAGS.debug: + train_job.submit(**job_args) + + +def main(argv: Sequence[str]) -> None: + config_file_path = FLAGS.config_file + print(f"Reading job config from {config_file_path}") + with open(config_file_path, encoding="utf-8") as config_file: + config = json.load(config_file) + + project_id = config["project_id"] + region = config["region"] + zone = config["zone"] + bucket = config["bucket"] + dataset_bucket = config["dataset_bucket"] + n_nodes = int(config["nodes"]) + machine_type = config["machine_type"] + num_gpus_per_node = int(config["gpus_per_node"]) + gpu_type = config["gpu_type"] + reservation_name = config.get("reservation_name") + reservation_full_name = ( + f"projects/{project_id}/zones/{zone}/reservations/{reservation_name}" + if "reservation_name" in config + else "" + ) + + strategy = config["strategy"] + recipe_name = config["recipe_name"] + job_prefix = config["job_prefix"] + image_uri = config["image_uri"] + + # Job name + timestamp = ( + datetime.datetime.now() + .astimezone(timezone("US/Pacific")) + .strftime("%Y%m%d_%H%M%S") + ) + job_name = f"{recipe_name}-{timestamp}" + if job_prefix: + job_name = f"{job_prefix}-{job_name}" + + base_output_dir = os.path.join("/gcs", bucket, job_name) + + # Training command and args + entrypoint_cmd = ["python3", "vdt/run.py"] + + dataset_bucket = f"gs://{config['dataset_bucket']}" + + trainer_args = [ + f"--train_data_gcs={dataset_bucket}", + "/opt/NeMo/examples/nlp/language_modeling/megatron_gpt_pretraining.py", + "--config-path=conf/", + f"--config-name={recipe_name}.yaml", + f"exp_manager.explicit_log_dir={base_output_dir}", + f"exp_manager.dllogger_logger_kwargs.json_file={base_output_dir}/dllogger.json", + "+exp_manager.create_tensorboard_logger=true", + "exp_manager.create_checkpoint_callback=false", + f"trainer.num_nodes={n_nodes}", + f"trainer.devices={num_gpus_per_node}", + "trainer.max_steps=10", + "trainer.log_every_n_steps=1", + "model.tokenizer.vocab_file=/data/gpt2-vocab.json", + "model.tokenizer.merge_file=/data/gpt2-merges.txt", + "model.data.data_prefix=[1.0,/data/hfbpe_gpt_training_data_text_document]", + ] + + launch_job( + job_name=job_name, + project=project_id, + region=region, + gcs_bucket=bucket, + image_uri=image_uri, + entrypoint_cmd=entrypoint_cmd, + trainer_args=trainer_args, + num_nodes=n_nodes, + machine_type=machine_type, + num_gpus_per_node=num_gpus_per_node, + gpu_type=gpu_type, + strategy=strategy, + reservation_name=reservation_full_name, + ) + + +if __name__ == "__main__": + app.run(main) diff --git a/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/scripts/run.py b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/scripts/run.py new file mode 100644 index 000000000..513f5b95a --- /dev/null +++ b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/scripts/run.py @@ -0,0 +1,85 @@ +"""Entrypoint for Vertex Distributed Training container.""" + +import argparse +import os +import sys +from collections.abc import Sequence +from subprocess import STDOUT, check_output, run + +from absl import app, flags, logging +from util import cluster_spec + +from retrying import retry + +# PyTorch barrier call which synchronizes all of the nodes before launching the training process. +# This makes sure that processes will block until all processes are ready. +# Improves the reliability of spot VM usage for multi-node training jobs + +@retry(stop_max_attempt_number=100, wait_exponential_multiplier=1000) +def barrier_with_retry() -> None: + import torch + logging.info("Starting barrier on RANK {}".format(os.environ["RANK"])) + torch.distributed.init_process_group() + torch.distributed.barrier() + torch.distributed.destroy_process_group() + logging.info("Finished barrier on RANK {}".format(os.environ["RANK"])) + +def main(unused_argv: Sequence[str]) -> None: + parser = argparse.ArgumentParser() + parser.add_argument( + "--train_data_gcs", + type=str, + help="Download training data from gcs path", + ) + args, unknown = parser.parse_known_args() + + for key, val in os.environ.items(): + logging.info("ENV %s=%s", key, val) + + if args.train_data_gcs: + local_dir = "/data" + if not os.path.exists(local_dir): + os.mkdir(local_dir) + logging.info("downloading %s to %s...", args.train_data_gcs, local_dir) + check_output( + [ + "gcloud", + "storage", + "cp", + "-r", + f"{args.train_data_gcs}/*", + local_dir, + ], + stderr=STDOUT, + ) + logging.info("%s downloaded.", args.train_data_gcs) + + primary_node_addr, primary_node_port, node_rank, num_nodes = ( + cluster_spec.get_cluster_spec() + ) + + cmd = [ + "torchrun", + "--nproc-per-node=8", + f"--nnodes={num_nodes}", + f"--node_rank={node_rank}", + ] + if num_nodes > 1: + cmd += [ + "--max-restarts=3", + "--rdzv-backend=static", + f'--rdzv_id={os.getenv("CLOUD_ML_JOB_ID", primary_node_port)}', + f"--rdzv-endpoint={primary_node_addr}:{primary_node_port}", + ] + cmd += unknown + + logging.info("launching with cmd: \n%s", " \\\n".join(cmd)) + barrier_with_retry() + run(cmd, stdout=sys.stdout, stderr=sys.stdout, check=True) + + +if __name__ == "__main__": + logging.get_absl_handler().python_handler.stream = sys.stdout + app.run( + main, flags_parser=lambda _args: flags.FLAGS(_args, known_only=True) + ) diff --git a/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/scripts/util/__init__.py b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/scripts/util/__init__.py new file mode 100644 index 000000000..e69de29bb diff --git a/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/scripts/util/cluster_spec.py b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/scripts/util/cluster_spec.py new file mode 100644 index 000000000..4ef1e2577 --- /dev/null +++ b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/scripts/util/cluster_spec.py @@ -0,0 +1,81 @@ +"""Get cluster info from environment variables.""" + +import dataclasses +import json +import os + +from absl import logging + + +@dataclasses.dataclass +class ClusterInfo: + """Contains information about the cluster. + + Attributes: + primary_node_addr: The address of the primary node. + primary_node_port: The port of the primary node. + node_rank: The rank of the node. + num_nodes: The number of nodes in the cluster. + """ + + primary_node_addr: str | None = None + primary_node_port: str | None = None + node_rank: int = 0 + num_nodes: int = 1 + + # Allows unpacking operation like + # primary_node_addr, primary_node_port, _, _ = ClusterInfo() + # See https://stackoverflow.com/a/70753113 + def __iter__(self): + return iter(dataclasses.astuple(self)) + + +def get_cluster_spec() -> ClusterInfo: + """Parses CLUSTER_SPEC environment variable and returns the cluster info. + + Returns: + A ClusterInfo object. + """ + cluster_spec = os.getenv("CLUSTER_SPEC", None) + + # If CLUSTER_SPEC is not set, use individual vars to construct cluster info. + if not cluster_spec: + cluster_info = ClusterInfo( + primary_node_addr=os.getenv("MASTER_ADDR", None), + primary_node_port=os.getenv("MASTER_PORT", None), + node_rank=int(os.getenv("RANK", "0")), + num_nodes=int(os.getenv("NNODES", "1")), + ) + return cluster_info + + cluster_data = json.loads(cluster_spec) + # Get primary node info + primary_node = cluster_data["cluster"]["workerpool0"][0] + logging.info("primary node: %s", primary_node) + primary_node_addr, primary_node_port = primary_node.split(":") + logging.info("primary node address: %s", primary_node_addr) + logging.info("primary node port: %s", primary_node_port) + + # Determine node rank of this machine + workerpool = cluster_data["task"]["type"] + if workerpool == "workerpool0": + node_rank = 0 + elif workerpool == "workerpool1": + # Add 1 for the primary node, since `index` is the index of workerpool1. + node_rank = cluster_data["task"]["index"] + 1 + else: + raise ValueError( + "Only workerpool0 and workerpool1 are supported. Unknown workerpool:" + f" {workerpool}" + ) + logging.info("node rank: %s", node_rank) + + # Calculate total nodes. + num_nodes = 1 # For the primary node. + if "workerpool1" in cluster_data["cluster"]: + num_nodes += len(cluster_data["cluster"]["workerpool1"]) + logging.info("num nodes: %s", num_nodes) + + return ClusterInfo( + primary_node_addr, primary_node_port, node_rank, num_nodes + ) diff --git a/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/scripts/util/cluster_spec_test.py b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/scripts/util/cluster_spec_test.py new file mode 100644 index 000000000..2839c69a4 --- /dev/null +++ b/community-content/vertex-distributed-training/a3mega/llama-3-8b-nemo-pretraining/scripts/util/cluster_spec_test.py @@ -0,0 +1,59 @@ +"""Add tests for cluster_spec.py.""" + +import os + +from . import cluster_spec + + +# TODO(styer): Use pytest instead +class ClusterSpecTest(googletest.TestCase): + + def setUp(self): + super().setUp() + self.curr_env_var = os.environ.copy() + + def tearDown(self): + super().tearDown() + os.environ = self.curr_env_var + + def test_get_cluster_spec_from_env_vars(self): + os.environ["CLUSTER_SPEC"] = "" + os.environ["MASTER_ADDR"] = "127.0.0.1" + os.environ["MASTER_PORT"] = "8080" + os.environ["RANK"] = "0" + os.environ["NNODES"] = "2" + cluster_info = cluster_spec.get_cluster_spec() + self.assertEqual(cluster_info.primary_node_addr, "127.0.0.1") + self.assertEqual(cluster_info.primary_node_port, "8080") + self.assertEqual(cluster_info.node_rank, 0) + self.assertEqual(cluster_info.num_nodes, 2) + + def test_get_cluster_spec_from_cluster_spec(self): + os.environ[ + "CLUSTER_SPEC" + ] = """ + { + "cluster": { + "workerpool0": [ + "127.0.0.1:8080" + ], + "workerpool1": [ + "127.0.0.2:8080", + "127.0.0.3:8080" + ] + }, + "task": { + "type": "workerpool1", + "index": 0 + } + } + """ + cluster_info = cluster_spec.get_cluster_spec() + self.assertEqual(cluster_info.primary_node_addr, "127.0.0.1") + self.assertEqual(cluster_info.primary_node_port, "8080") + self.assertEqual(cluster_info.node_rank, 1) + self.assertEqual(cluster_info.num_nodes, 3) + + +if __name__ == "__main__": + googletest.main()