diff --git a/deepspeed/monitor/comet.py b/deepspeed/monitor/comet.py index d8bc4017800f..9e161ac7675a 100644 --- a/deepspeed/monitor/comet.py +++ b/deepspeed/monitor/comet.py @@ -70,6 +70,10 @@ def write_events(self, event_list: List[Event]) -> None: step=engine_global_samples, ) + def update_config(self, config_dict): + if self._experiment is not None: + self._experiment.log_parameters(config_dict) + class EventsLogScheduler: diff --git a/deepspeed/monitor/csv_monitor.py b/deepspeed/monitor/csv_monitor.py index c7a19b14ad82..e452efc57def 100644 --- a/deepspeed/monitor/csv_monitor.py +++ b/deepspeed/monitor/csv_monitor.py @@ -65,3 +65,14 @@ def write_events(self, event_list): self.filenames.append(filename) csv_monitor_writer.writerow(['step', header]) csv_monitor_writer.writerow([step, value]) + + def update_config(self, config_dict): + if not self.enabled or dist.get_rank() != 0 or self.log_dir is None: + return + import csv + fname = os.path.join(self.log_dir, 'deepspeed_config.csv') + with open(fname, 'w') as csv_monitor_file: + csv_monitor_writer = csv.writer(csv_monitor_file) + csv_monitor_writer.writerow(['key', 'value']) + for key, value in config_dict.items(): + csv_monitor_writer.writerow([key, value]) diff --git a/deepspeed/monitor/monitor.py b/deepspeed/monitor/monitor.py index e7e26dc483d9..f9ae6f957b6b 100644 --- a/deepspeed/monitor/monitor.py +++ b/deepspeed/monitor/monitor.py @@ -57,3 +57,14 @@ def write_events(self, event_list): self.csv_monitor.write_events(event_list) if self.comet_monitor is not None: self.comet_monitor.write_events(event_list) + + def update_config(self, config_dict): + if dist.get_rank() == 0: + if self.tb_monitor is not None: + self.tb_monitor.update_config(config_dict) + if self.wandb_monitor is not None: + self.wandb_monitor.update_config(config_dict) + if self.csv_monitor is not None: + self.csv_monitor.update_config(config_dict) + if self.comet_monitor is not None: + self.comet_monitor.update_config(config_dict) diff --git a/deepspeed/monitor/tensorboard.py b/deepspeed/monitor/tensorboard.py index 985c9ed44b6f..45526e27dc30 100644 --- a/deepspeed/monitor/tensorboard.py +++ b/deepspeed/monitor/tensorboard.py @@ -5,6 +5,7 @@ from .utils import check_tb_availability from .monitor import Monitor +import json import os import deepspeed.comm as dist @@ -54,3 +55,7 @@ def write_events(self, event_list, flush=True): def flush(self): if self.enabled and self.summary_writer is not None and dist.get_rank() == 0: self.summary_writer.flush() + + def update_config(self, config_dict): + if self.summary_writer is not None: + self.summary_writer.add_text("DeepSpeed/config", json.dumps(config_dict, indent=2), global_step=0) diff --git a/deepspeed/monitor/utils.py b/deepspeed/monitor/utils.py index f5530e8532e1..57e1f58c3882 100644 --- a/deepspeed/monitor/utils.py +++ b/deepspeed/monitor/utils.py @@ -35,3 +35,111 @@ def check_comet_availability(): except ImportError: print('If you want to use comet logging, please `pip install "comet_ml>=3.41.0"`') raise + + +def _json_safe_int(value, default): + if value is None: + return default + return int(value) + + +def _parallel_int(getter, default): + # TP/DP groups are only created when those parallelisms are configured. + try: + return _json_safe_int(getter(), default) + except (AssertionError, RuntimeError, ValueError, TypeError): + return default + + +def _offload_device(offload_cfg): + if offload_cfg is None: + return None + device = getattr(offload_cfg, "device", None) + if device is None: + return None + device_name = device.value if hasattr(device, "value") else str(device) + if device_name == 'none': + return None + return device_name + + +def _is_sequence_only_mpu(mpu): + # Ulysses parallel_state_sp aliases MP getters to SP. It is not a TP MPU. + return mpu is not None and hasattr(mpu, "initialize_sequence_parallel") + + +def _pipeline_parallel_rank(engine): + if not getattr(engine, "pipeline_parallelism", False) and getattr(engine, "mpu", None) is None: + return 0 + mpu = getattr(engine, "mpu", None) + if mpu is None: + return 0 + if hasattr(mpu, "get_pipeline_model_parallel_rank"): + return int(mpu.get_pipeline_model_parallel_rank()) + if hasattr(mpu, "get_pipe_parallel_rank"): + return int(mpu.get_pipe_parallel_rank()) + return 0 + + +def _data_parallel_world_size(engine, default): + mpu = getattr(engine, "mpu", None) + if mpu is not None and not _is_sequence_only_mpu(mpu) and hasattr(mpu, "get_data_parallel_world_size"): + return int(mpu.get_data_parallel_world_size()) + from deepspeed.utils import groups + size = _parallel_int(groups._get_data_parallel_world_size, default) + if size is None: + return default + return size + + +def _data_parallel_rank(engine, default): + mpu = getattr(engine, "mpu", None) + if mpu is not None and not _is_sequence_only_mpu(mpu) and hasattr(mpu, "get_data_parallel_rank"): + return int(mpu.get_data_parallel_rank()) + from deepspeed.utils import groups + return _parallel_int(groups._get_data_parallel_rank, default) + + +def _tensor_parallel_world_size(engine): + mpu = getattr(engine, "mpu", None) + if _is_sequence_only_mpu(mpu): + return 1 + from deepspeed.utils.bwc import bwc_tensor_model_parallel_world_size + return int(bwc_tensor_model_parallel_world_size(mpu)) + + +def _tensor_parallel_rank(engine): + mpu = getattr(engine, "mpu", None) + if _is_sequence_only_mpu(mpu): + return 0 + from deepspeed.utils.bwc import bwc_tensor_model_parallel_rank + return int(bwc_tensor_model_parallel_rank(mpu)) + + +def collect_monitor_config(engine): + """Build a JSON-safe snapshot of ZeRO / precision / parallelism settings.""" + from deepspeed.utils import groups + from deepspeed.utils.bwc import bwc_pipeline_parallel_world_size + import deepspeed.comm as dist + + world_size_default = dist.get_world_size() if dist.is_initialized() else 1 + rank_default = dist.get_rank() if dist.is_initialized() else 0 + + return { + 'zero_stage': int(engine.zero_optimization_stage()), + 'offload_optimizer': _offload_device(engine.zero_offload_optimizer()), + 'offload_param': _offload_device(engine.zero_offload_param()), + 'fp16': bool(engine.fp16_enabled()), + 'bf16': bool(engine.bfloat16_enabled()), + 'train_batch_size': int(engine.train_batch_size()), + 'train_micro_batch_size_per_gpu': int(engine.train_micro_batch_size_per_gpu()), + 'gradient_accumulation_steps': int(engine.gradient_accumulation_steps()), + 'data_parallel_world_size': _data_parallel_world_size(engine, world_size_default), + 'tensor_parallel_world_size': _tensor_parallel_world_size(engine), + 'pipeline_parallel_world_size': int(bwc_pipeline_parallel_world_size(engine.mpu)), + 'sequence_parallel_world_size': int(groups._get_sequence_parallel_world_size()), + 'data_parallel_rank': _data_parallel_rank(engine, rank_default), + 'tensor_parallel_rank': _tensor_parallel_rank(engine), + 'pipeline_parallel_rank': _pipeline_parallel_rank(engine), + 'sequence_parallel_rank': int(groups._get_sequence_parallel_rank()), + } diff --git a/deepspeed/monitor/wandb.py b/deepspeed/monitor/wandb.py index 174a2eb2d3b7..fb3cbd446aab 100644 --- a/deepspeed/monitor/wandb.py +++ b/deepspeed/monitor/wandb.py @@ -36,3 +36,8 @@ def write_events(self, event_list): value = event[1] step = event[2] self.log({label: value}, step=step) + + def update_config(self, config_dict): + if self.enabled and dist.get_rank() == 0: + import wandb + wandb.config.update(config_dict, allow_val_change=True) diff --git a/deepspeed/runtime/engine.py b/deepspeed/runtime/engine.py index 86918bd71c5a..0242c9b77973 100644 --- a/deepspeed/runtime/engine.py +++ b/deepspeed/runtime/engine.py @@ -113,6 +113,7 @@ STEP_GLOBAL_TIMER from deepspeed.utils.debug import debug_extract_module_and_param_names, debug_clear_module_and_param_names from deepspeed.monitor.monitor import MonitorMaster +from deepspeed.monitor.utils import collect_monitor_config from deepspeed.runtime.progressive_layer_drop import ProgressiveLayerDrop from deepspeed.runtime.utils import clip_grad_norm_, compare_tensors_in_structures, maybe_loss_for_backward from deepspeed.runtime.eigenvalue import Eigenvalue @@ -613,6 +614,9 @@ def __init__(self, if self.dist_backend is None: self.enable_backward_allreduce = False + if self.monitor.enabled: + self.monitor.update_config(collect_monitor_config(self)) + def _optimized_linear_offload_setup(self): self.optimized_linear_base_weight_sharding = False self.optimized_linear_lora_enabled = False diff --git a/docs/_tutorials/monitor.md b/docs/_tutorials/monitor.md index 5e7a6fc4e834..281a6237913f 100644 --- a/docs/_tutorials/monitor.md +++ b/docs/_tutorials/monitor.md @@ -42,18 +42,18 @@ When using DeepSpeed for model training, the Monitor can be configured in the De "enabled": true, "output_path": "output/ds_logs/", "job_name": "train_bert" - } + }, "wandb": { "enabled": true, "team": "my_team", "group": "my_group", "project": "my_project" - } + }, "comet": { "enabled": true, "project": "my_project", "experiment_name": "my_experiment" - } + }, "csv_monitor": { "enabled": true, "output_path": "output/ds_logs/", @@ -62,7 +62,7 @@ When using DeepSpeed for model training, the Monitor can be configured in the De } ``` -DeepSpeed will automatically log to all available and enabled monitoring backends listed in the config, and will generate live monitoring views such as those listed above. +DeepSpeed will automatically log to all available and enabled monitoring backends listed in the config, and will generate live monitoring views such as those listed above. When the engine is initialized, DeepSpeed also records static parallelism and ZeRO settings such as ZeRO stage, offload devices, precision, and DP/TP/PP/SP sizes on the enabled monitor backends. ### Custom Monitoring diff --git a/tests/unit/monitor/test_monitor.py b/tests/unit/monitor/test_monitor.py index d4b3cf43921d..7af915e93230 100644 --- a/tests/unit/monitor/test_monitor.py +++ b/tests/unit/monitor/test_monitor.py @@ -3,18 +3,46 @@ # DeepSpeed Team +import json +import os +import sys +from types import SimpleNamespace + +import deepspeed from deepspeed.monitor.tensorboard import TensorBoardMonitor from deepspeed.monitor.wandb import WandbMonitor from deepspeed.monitor.csv_monitor import csvMonitor from deepspeed.monitor.config import DeepSpeedMonitorConfig from deepspeed.monitor.comet import CometMonitor +from deepspeed.monitor.monitor import MonitorMaster +from deepspeed.monitor.utils import collect_monitor_config from unit.common import DistributedTest +from unit.simple_model import SimpleModel from unittest.mock import Mock, patch from deepspeed.runtime.config import DeepSpeedConfig import deepspeed.comm as dist +EXPECTED_MONITOR_CONFIG_KEYS = { + "zero_stage", + "offload_optimizer", + "offload_param", + "fp16", + "bf16", + "train_batch_size", + "train_micro_batch_size_per_gpu", + "gradient_accumulation_steps", + "data_parallel_world_size", + "tensor_parallel_world_size", + "pipeline_parallel_world_size", + "sequence_parallel_world_size", + "data_parallel_rank", + "tensor_parallel_rank", + "pipeline_parallel_rank", + "sequence_parallel_rank", +} + class TestTensorBoard(DistributedTest): world_size = 2 @@ -164,3 +192,272 @@ def test_empty_comet(self): assert comet_monitor.enabled == defaults.enabled assert comet_monitor.samples_log_interval == defaults.samples_log_interval mock_start.assert_not_called() + + +class TestMonitorConfigUpdate(DistributedTest): + world_size = 2 + + def test_wandb_update_config_rank0_only(self): + config_dict = { + "train_batch_size": 2, + "wandb": { + "enabled": True, + "group": "my_group", + "team": "my_team", + "project": "my_project" + } + } + ds_config = DeepSpeedConfig(config_dict) + payload = { + "zero_stage": 2, + "data_parallel_world_size": 2, + "fp16": False, + } + mock_wandb = Mock() + with patch.dict(sys.modules, {"wandb": mock_wandb}): + with patch("deepspeed.monitor.wandb.check_wandb_availability"): + wandb_monitor = WandbMonitor(ds_config.monitor_config.wandb) + wandb_monitor.update_config(payload) + + if dist.get_rank() == 0: + mock_wandb.init.assert_called_once() + mock_wandb.config.update.assert_called_once_with(payload, allow_val_change=True) + called_cfg = mock_wandb.config.update.call_args[0][0] + assert "zero_stage" in called_cfg + assert "data_parallel_world_size" in called_cfg + else: + mock_wandb.config.update.assert_not_called() + + def test_monitor_master_update_config_rank0_only(self): + config_dict = {"train_batch_size": 2} + ds_config = DeepSpeedConfig(config_dict) + master = MonitorMaster(ds_config.monitor_config) + master.tb_monitor = Mock() + master.wandb_monitor = Mock() + master.csv_monitor = Mock() + master.comet_monitor = Mock() + + payload = {"zero_stage": 1, "data_parallel_world_size": 2} + master.update_config(payload) + + if dist.get_rank() == 0: + master.tb_monitor.update_config.assert_called_once_with(payload) + master.wandb_monitor.update_config.assert_called_once_with(payload) + master.csv_monitor.update_config.assert_called_once_with(payload) + master.comet_monitor.update_config.assert_called_once_with(payload) + else: + master.tb_monitor.update_config.assert_not_called() + master.wandb_monitor.update_config.assert_not_called() + master.csv_monitor.update_config.assert_not_called() + master.comet_monitor.update_config.assert_not_called() + + def test_comet_update_config(self): + import comet_ml + mock_experiment = Mock() + mock_start = Mock(return_value=mock_experiment) + + config_dict = { + "train_batch_size": 2, + "comet": { + "enabled": True, + "project": "some-project", + "api_key": "some-api-key" + } + } + ds_config = DeepSpeedConfig(config_dict) + payload = {"zero_stage": 1, "bf16": True} + + with patch.object(comet_ml, "start", mock_start): + comet_monitor = CometMonitor(ds_config.monitor_config.comet) + comet_monitor.update_config(payload) + + if dist.get_rank() == 0: + mock_experiment.log_parameters.assert_called_once_with(payload) + else: + mock_experiment.log_parameters.assert_not_called() + + def test_tensorboard_update_config(self): + config_dict = { + "train_batch_size": 2, + "tensorboard": { + "enabled": True, + "output_path": "test_output/ds_logs/", + "job_name": "test" + } + } + ds_config = DeepSpeedConfig(config_dict) + tb_monitor = TensorBoardMonitor(ds_config.monitor_config.tensorboard) + writer = Mock() + if dist.get_rank() == 0: + tb_monitor.summary_writer = writer + payload = {"zero_stage": 0, "fp16": True} + tb_monitor.update_config(payload) + + if dist.get_rank() == 0: + writer.add_text.assert_called_once() + args, kwargs = writer.add_text.call_args + assert args[0] == "DeepSpeed/config" + logged = json.loads(args[1]) + assert logged["zero_stage"] == 0 + assert logged["fp16"] is True + assert kwargs.get("global_step", args[2] if len(args) > 2 else None) == 0 + else: + assert tb_monitor.summary_writer is None + writer.add_text.assert_not_called() + + def test_csv_update_config(self, tmpdir): + config_dict = { + "train_batch_size": 2, + "csv_monitor": { + "enabled": True, + "output_path": str(tmpdir), + "job_name": "cfg_test" + } + } + ds_config = DeepSpeedConfig(config_dict) + csv_monitor = csvMonitor(ds_config.monitor_config.csv_monitor) + payload = {"zero_stage": 2, "fp16": False, "offload_optimizer": None} + csv_monitor.update_config(payload) + + cfg_path = os.path.join(str(tmpdir), "cfg_test", "deepspeed_config.csv") + if dist.get_rank() == 0: + assert os.path.isfile(cfg_path) + with open(cfg_path, "r") as f: + contents = f.read() + assert "key,value" in contents.replace(" ", "") + assert "zero_stage" in contents + assert "2" in contents + else: + assert not os.path.isfile(cfg_path) + + +class TestCollectMonitorConfig(DistributedTest): + world_size = 2 + + def test_collect_monitor_config_from_engine_methods(self): + engine = Mock() + engine.zero_optimization_stage.return_value = 3 + engine.zero_offload_optimizer.return_value = Mock(device="cpu") + engine.zero_offload_param.return_value = None + engine.fp16_enabled.return_value = False + engine.bfloat16_enabled.return_value = True + engine.train_batch_size.return_value = 16 + engine.train_micro_batch_size_per_gpu.return_value = 2 + engine.gradient_accumulation_steps.return_value = 4 + engine.pipeline_parallelism = False + engine.mpu = None + + cfg = collect_monitor_config(engine) + assert set(cfg.keys()) == EXPECTED_MONITOR_CONFIG_KEYS + assert cfg["zero_stage"] == 3 + assert cfg["offload_optimizer"] == "cpu" + assert cfg["offload_param"] is None + assert cfg["fp16"] is False + assert cfg["bf16"] is True + assert cfg["train_batch_size"] == 16 + assert cfg["train_micro_batch_size_per_gpu"] == 2 + assert cfg["gradient_accumulation_steps"] == 4 + assert cfg["data_parallel_world_size"] == 2 + assert cfg["tensor_parallel_world_size"] == 1 + assert cfg["pipeline_parallel_world_size"] == 1 + assert cfg["sequence_parallel_world_size"] == 1 + assert cfg["pipeline_parallel_rank"] == 0 + assert cfg["sequence_parallel_rank"] == 0 + assert cfg["data_parallel_rank"] == dist.get_rank() + for value in cfg.values(): + assert value is None or isinstance(value, (bool, int, str)) + + def test_collect_monitor_config_reads_pipeline_mpu(self): + engine = Mock() + engine.zero_optimization_stage.return_value = 0 + engine.zero_offload_optimizer.return_value = None + engine.zero_offload_param.return_value = None + engine.fp16_enabled.return_value = False + engine.bfloat16_enabled.return_value = False + engine.train_batch_size.return_value = 8 + engine.train_micro_batch_size_per_gpu.return_value = 1 + engine.gradient_accumulation_steps.return_value = 1 + engine.pipeline_parallelism = True + engine.mpu = SimpleNamespace( + get_data_parallel_rank=lambda: 1, + get_data_parallel_world_size=lambda: 2, + get_slice_parallel_rank=lambda: 1, + get_slice_parallel_world_size=lambda: 2, + get_pipe_parallel_rank=lambda: 0, + get_pipe_parallel_world_size=lambda: 2, + ) + + cfg = collect_monitor_config(engine) + assert cfg["data_parallel_world_size"] == 2 + assert cfg["data_parallel_rank"] == 1 + assert cfg["tensor_parallel_world_size"] == 2 + assert cfg["tensor_parallel_rank"] == 1 + assert cfg["pipeline_parallel_world_size"] == 2 + assert cfg["pipeline_parallel_rank"] == 0 + + def test_collect_monitor_config_ignores_sequence_only_mpu_tp_alias(self): + engine = Mock() + engine.zero_optimization_stage.return_value = 0 + engine.zero_offload_optimizer.return_value = None + engine.zero_offload_param.return_value = None + engine.fp16_enabled.return_value = False + engine.bfloat16_enabled.return_value = False + engine.train_batch_size.return_value = 4 + engine.train_micro_batch_size_per_gpu.return_value = 1 + engine.gradient_accumulation_steps.return_value = 1 + engine.pipeline_parallelism = False + engine.mpu = SimpleNamespace( + initialize_sequence_parallel=lambda *args, **kwargs: None, + get_model_parallel_rank=lambda: 3, + get_model_parallel_world_size=lambda: 4, + ) + + cfg = collect_monitor_config(engine) + assert cfg["tensor_parallel_world_size"] == 1 + assert cfg["tensor_parallel_rank"] == 0 + + def test_engine_init_collects_config(self, tmpdir): + hidden_dim = 4 + model = SimpleModel(hidden_dim) + config_dict = { + "train_batch_size": 2, + "train_micro_batch_size_per_gpu": 1, + "optimizer": { + "type": "Adam", + "params": { + "lr": 0.0001 + } + }, + "zero_optimization": { + "stage": 2 + }, + "csv_monitor": { + "enabled": True, + "output_path": str(tmpdir), + "job_name": "engine_cfg" + } + } + engine, _, _, _ = deepspeed.initialize(config=config_dict, model=model, model_parameters=model.parameters()) + cfg = collect_monitor_config(engine) + assert set(cfg.keys()) == EXPECTED_MONITOR_CONFIG_KEYS + assert cfg["zero_stage"] == 2 + assert cfg["offload_optimizer"] is None + assert cfg["offload_param"] is None + assert cfg["fp16"] is False + assert cfg["bf16"] is False + assert cfg["train_batch_size"] == 2 + assert cfg["train_micro_batch_size_per_gpu"] == 1 + assert cfg["gradient_accumulation_steps"] == 1 + assert cfg["data_parallel_world_size"] == 2 + assert cfg["tensor_parallel_world_size"] == 1 + assert cfg["pipeline_parallel_world_size"] == 1 + assert cfg["sequence_parallel_world_size"] == 1 + assert cfg["pipeline_parallel_rank"] == 0 + + cfg_path = os.path.join(str(tmpdir), "engine_cfg", "deepspeed_config.csv") + if dist.get_rank() == 0: + assert os.path.isfile(cfg_path) + with open(cfg_path, "r") as f: + contents = f.read() + assert "zero_stage" in contents + assert "data_parallel_world_size" in contents