Skip to content

Commit f33ad0a

Browse files
committed
Log ZeRO and parallelism settings to monitor backends
Record static training topology on WandB/Comet/TensorBoard/CSV at engine init so runs are comparable without digging through logs. Fixes #7494 Signed-off-by: YeonwooSung <neos960518@gmail.com>
1 parent 56de570 commit f33ad0a

9 files changed

Lines changed: 359 additions & 4 deletions

File tree

deepspeed/monitor/comet.py

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -70,6 +70,10 @@ def write_events(self, event_list: List[Event]) -> None:
7070
step=engine_global_samples,
7171
)
7272

73+
def update_config(self, config_dict):
74+
if self._experiment is not None:
75+
self._experiment.log_parameters(config_dict)
76+
7377

7478
class EventsLogScheduler:
7579

deepspeed/monitor/csv_monitor.py

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -65,3 +65,14 @@ def write_events(self, event_list):
6565
self.filenames.append(filename)
6666
csv_monitor_writer.writerow(['step', header])
6767
csv_monitor_writer.writerow([step, value])
68+
69+
def update_config(self, config_dict):
70+
if not self.enabled or dist.get_rank() != 0 or self.log_dir is None:
71+
return
72+
import csv
73+
fname = os.path.join(self.log_dir, 'deepspeed_config.csv')
74+
with open(fname, 'w') as csv_monitor_file:
75+
csv_monitor_writer = csv.writer(csv_monitor_file)
76+
csv_monitor_writer.writerow(['key', 'value'])
77+
for key, value in config_dict.items():
78+
csv_monitor_writer.writerow([key, value])

deepspeed/monitor/monitor.py

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -57,3 +57,14 @@ def write_events(self, event_list):
5757
self.csv_monitor.write_events(event_list)
5858
if self.comet_monitor is not None:
5959
self.comet_monitor.write_events(event_list)
60+
61+
def update_config(self, config_dict):
62+
if dist.get_rank() == 0:
63+
if self.tb_monitor is not None:
64+
self.tb_monitor.update_config(config_dict)
65+
if self.wandb_monitor is not None:
66+
self.wandb_monitor.update_config(config_dict)
67+
if self.csv_monitor is not None:
68+
self.csv_monitor.update_config(config_dict)
69+
if self.comet_monitor is not None:
70+
self.comet_monitor.update_config(config_dict)

deepspeed/monitor/tensorboard.py

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@
55

66
from .utils import check_tb_availability
77
from .monitor import Monitor
8+
import json
89
import os
910

1011
import deepspeed.comm as dist
@@ -54,3 +55,7 @@ def write_events(self, event_list, flush=True):
5455
def flush(self):
5556
if self.enabled and self.summary_writer is not None and dist.get_rank() == 0:
5657
self.summary_writer.flush()
58+
59+
def update_config(self, config_dict):
60+
if self.summary_writer is not None:
61+
self.summary_writer.add_text("DeepSpeed/config", json.dumps(config_dict, indent=2), global_step=0)

deepspeed/monitor/utils.py

Lines changed: 68 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -35,3 +35,71 @@ def check_comet_availability():
3535
except ImportError:
3636
print('If you want to use comet logging, please `pip install "comet_ml>=3.41.0"`')
3737
raise
38+
39+
40+
def _json_safe_int(value, default):
41+
if value is None:
42+
return default
43+
return int(value)
44+
45+
46+
def _parallel_int(getter, default):
47+
# TP/DP groups are only created when those parallelisms are configured.
48+
try:
49+
return _json_safe_int(getter(), default)
50+
except (AssertionError, RuntimeError, ValueError, TypeError):
51+
return default
52+
53+
54+
def _offload_device(offload_cfg):
55+
if offload_cfg is None:
56+
return None
57+
device = getattr(offload_cfg, "device", None)
58+
if device is None:
59+
return None
60+
device_name = device.value if hasattr(device, "value") else str(device)
61+
if device_name == 'none':
62+
return None
63+
return device_name
64+
65+
66+
def _pipeline_parallel_rank(engine):
67+
if not getattr(engine, "pipeline_parallelism", False) and getattr(engine, "mpu", None) is None:
68+
return 0
69+
mpu = getattr(engine, "mpu", None)
70+
if mpu is None:
71+
return 0
72+
if hasattr(mpu, "get_pipeline_model_parallel_rank"):
73+
return int(mpu.get_pipeline_model_parallel_rank())
74+
if hasattr(mpu, "get_pipe_parallel_rank"):
75+
return int(mpu.get_pipe_parallel_rank())
76+
return 0
77+
78+
79+
def collect_monitor_config(engine):
80+
"""Build a JSON-safe snapshot of ZeRO / precision / parallelism settings."""
81+
from deepspeed.utils import groups
82+
from deepspeed.utils.bwc import bwc_pipeline_parallel_world_size
83+
import deepspeed.comm as dist
84+
85+
world_size_default = dist.get_world_size() if dist.is_initialized() else 1
86+
rank_default = dist.get_rank() if dist.is_initialized() else 0
87+
88+
return {
89+
'zero_stage': int(engine.zero_optimization_stage()),
90+
'offload_optimizer': _offload_device(engine.zero_offload_optimizer()),
91+
'offload_param': _offload_device(engine.zero_offload_param()),
92+
'fp16': bool(engine.fp16_enabled()),
93+
'bf16': bool(engine.bfloat16_enabled()),
94+
'train_batch_size': int(engine.train_batch_size()),
95+
'train_micro_batch_size_per_gpu': int(engine.train_micro_batch_size_per_gpu()),
96+
'gradient_accumulation_steps': int(engine.gradient_accumulation_steps()),
97+
'data_parallel_world_size': _parallel_int(groups.get_data_parallel_world_size, world_size_default),
98+
'tensor_parallel_world_size': _parallel_int(groups.get_tensor_model_parallel_world_size, 1),
99+
'pipeline_parallel_world_size': int(bwc_pipeline_parallel_world_size(engine.mpu)),
100+
'sequence_parallel_world_size': int(groups._get_sequence_parallel_world_size()),
101+
'data_parallel_rank': _parallel_int(groups.get_data_parallel_rank, rank_default),
102+
'tensor_parallel_rank': _parallel_int(groups.get_tensor_model_parallel_rank, 0),
103+
'pipeline_parallel_rank': _pipeline_parallel_rank(engine),
104+
'sequence_parallel_rank': int(groups._get_sequence_parallel_rank()),
105+
}

deepspeed/monitor/wandb.py

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -36,3 +36,8 @@ def write_events(self, event_list):
3636
value = event[1]
3737
step = event[2]
3838
self.log({label: value}, step=step)
39+
40+
def update_config(self, config_dict):
41+
if self.enabled and dist.get_rank() == 0:
42+
import wandb
43+
wandb.config.update(config_dict, allow_val_change=True)

deepspeed/runtime/engine.py

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -113,6 +113,7 @@
113113
STEP_GLOBAL_TIMER
114114
from deepspeed.utils.debug import debug_extract_module_and_param_names, debug_clear_module_and_param_names
115115
from deepspeed.monitor.monitor import MonitorMaster
116+
from deepspeed.monitor.utils import collect_monitor_config
116117
from deepspeed.runtime.progressive_layer_drop import ProgressiveLayerDrop
117118
from deepspeed.runtime.utils import clip_grad_norm_, compare_tensors_in_structures, maybe_loss_for_backward
118119
from deepspeed.runtime.eigenvalue import Eigenvalue
@@ -613,6 +614,9 @@ def __init__(self,
613614
if self.dist_backend is None:
614615
self.enable_backward_allreduce = False
615616

617+
if self.monitor.enabled:
618+
self.monitor.update_config(collect_monitor_config(self))
619+
616620
def _optimized_linear_offload_setup(self):
617621
self.optimized_linear_base_weight_sharding = False
618622
self.optimized_linear_lora_enabled = False

docs/_tutorials/monitor.md

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -42,18 +42,18 @@ When using DeepSpeed for model training, the Monitor can be configured in the De
4242
"enabled": true,
4343
"output_path": "output/ds_logs/",
4444
"job_name": "train_bert"
45-
}
45+
},
4646
"wandb": {
4747
"enabled": true,
4848
"team": "my_team",
4949
"group": "my_group",
5050
"project": "my_project"
51-
}
51+
},
5252
"comet": {
5353
"enabled": true,
5454
"project": "my_project",
5555
"experiment_name": "my_experiment"
56-
}
56+
},
5757
"csv_monitor": {
5858
"enabled": true,
5959
"output_path": "output/ds_logs/",
@@ -62,7 +62,7 @@ When using DeepSpeed for model training, the Monitor can be configured in the De
6262
}
6363
```
6464

65-
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.
65+
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.
6666

6767
### Custom Monitoring
6868

0 commit comments

Comments
 (0)