Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions deepspeed/monitor/comet.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:

Expand Down
11 changes: 11 additions & 0 deletions deepspeed/monitor/csv_monitor.py
Original file line number Diff line number Diff line change
Expand Up @@ -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])
11 changes: 11 additions & 0 deletions deepspeed/monitor/monitor.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
5 changes: 5 additions & 0 deletions deepspeed/monitor/tensorboard.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@

from .utils import check_tb_availability
from .monitor import Monitor
import json
import os

import deepspeed.comm as dist
Expand Down Expand Up @@ -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)
108 changes: 108 additions & 0 deletions deepspeed/monitor/utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()),
}
5 changes: 5 additions & 0 deletions deepspeed/monitor/wandb.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
4 changes: 4 additions & 0 deletions deepspeed/runtime/engine.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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))
Comment thread
YeonwooSung marked this conversation as resolved.

def _optimized_linear_offload_setup(self):
self.optimized_linear_base_weight_sharding = False
self.optimized_linear_lora_enabled = False
Expand Down
8 changes: 4 additions & 4 deletions docs/_tutorials/monitor.md
Original file line number Diff line number Diff line change
Expand Up @@ -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/",
Expand All @@ -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

Expand Down
Loading
Loading