From e34446317dbdaef874ef98f8c7937ca082362455 Mon Sep 17 00:00:00 2001 From: Jeff Rasley Date: Sat, 24 Apr 2021 00:33:28 +0000 Subject: [PATCH 1/6] add zero-1 mode to zero-2 codebase --- deepspeed/__init__.py | 6 +-- deepspeed/runtime/config.py | 5 -- deepspeed/runtime/engine.py | 62 +++++++++++++----------- deepspeed/runtime/zero/config.py | 20 ++++++++ deepspeed/runtime/zero/constants.py | 16 +++++- deepspeed/runtime/zero/stage2.py | 75 ++++++++++++++++++++++++----- requirements/requirements.txt | 1 + 7 files changed, 137 insertions(+), 48 deletions(-) diff --git a/deepspeed/__init__.py b/deepspeed/__init__.py index 0ace7b23ae27..1a9dd5bed73c 100755 --- a/deepspeed/__init__.py +++ b/deepspeed/__init__.py @@ -3,6 +3,7 @@ ''' import sys import types +import packaging from . import ops @@ -25,9 +26,8 @@ def _parse_version(version_str): '''Parse a version string and extract the major, minor, and patch versions.''' - import re - matched = re.search('^(\d+)\.(\d+)\.(\d+)', version_str) - return int(matched.group(1)), int(matched.group(2)), int(matched.group(3)) + ver = packaging.version.parse(version_str) + return ver.major, ver.minor, ver.micro # Export version information diff --git a/deepspeed/runtime/config.py b/deepspeed/runtime/config.py index 727c0810290e..513b1504b6ed 100755 --- a/deepspeed/runtime/config.py +++ b/deepspeed/runtime/config.py @@ -765,12 +765,7 @@ def _do_error_check(self): GRADIENT_ACCUMULATION_STEPS) if self.zero_enabled: - if self.zero_optimization_stage < ZERO_OPTIMIZATION_GRADIENTS: - assert self.fp16_enabled, "DeepSpeedConfig: ZeRO is only supported if fp16 is enabled" assert self.zero_optimization_stage <= MAX_STAGE_ZERO_OPTIMIZATION, "DeepSpeedConfig: Maximum supported ZeRO stage is {}".format(MAX_STAGE_ZERO_OPTIMIZATION) - #if self.zero_config.cpu_offload is True: - # assert self.zero_optimization_stage == ZERO_OPTIMIZATION_GRADIENTS, "DeepSpeedConfig: cpu-offload supported ZeRO stage is {}".format(ZERO_OPTIMIZATION_GRADIENTS) - #assert self.gradient_accumulation_steps == 1, "DeepSpeedConfig: {}is not supported for {}".format(GRADIENT_ACCUMULATION_STEPS, ZERO_OPTIMIZATION_CPU_OFFLOAD) def _do_warning_check(self): fp16_enabled = self.fp16_enabled or self.zero_enabled diff --git a/deepspeed/runtime/engine.py b/deepspeed/runtime/engine.py index f803592efb53..eaf498f1345e 100755 --- a/deepspeed/runtime/engine.py +++ b/deepspeed/runtime/engine.py @@ -44,6 +44,7 @@ from ..ops.op_builder import UtilsBuilder from ..ops.adam import DeepSpeedCPUAdam from ..ops.adam import FusedAdam +from ..git_version_info import version from deepspeed.profiling.flops_profiler.profiler import FlopsProfiler @@ -380,6 +381,12 @@ def zero_gather_fp16_weights_on_model_save(self): def zero_find_unused_parameters(self): return self._config.zero_config.find_unused_parameters + def zero_grad_hooks(self): + return self._config.zero_config.grad_hooks + + def zero_legacy_stage1(self): + return self._config.zero_config.legacy_stage1 + def fp16_enabled(self): return self._config.fp16_enabled @@ -754,7 +761,8 @@ def _configure_zero_optimizer(self, optimizer): assert not self.allreduce_always_fp32(), "ZeRO does not support 'fp32_allreduce': true" timers = self.timers if self.wall_clock_breakdown() else None - if zero_stage == ZERO_OPTIMIZATION_OPTIMIZER_STATES: + if self.zero_legacy_stage1( + ) and zero_stage == ZERO_OPTIMIZATION_OPTIMIZER_STATES: optimizer = FP16_DeepSpeedZeroOptimizer_Stage1( optimizer, static_loss_scale=self.loss_scale(), @@ -767,7 +775,7 @@ def _configure_zero_optimizer(self, optimizer): dp_process_group=self.data_parallel_group, elastic_checkpoint=self.zero_elastic_checkpoint(), mpu=self.mpu) - elif zero_stage == ZERO_OPTIMIZATION_GRADIENTS: + elif zero_stage <= ZERO_OPTIMIZATION_GRADIENTS: optimizer = FP16_DeepSpeedZeroOptimizer( optimizer, timers=timers, @@ -786,7 +794,9 @@ def _configure_zero_optimizer(self, optimizer): postscale_gradients=self.postscale_gradients(), gradient_predivide_factor=self.gradient_predivide_factor(), gradient_accumulation_steps=self.gradient_accumulation_steps(), - find_unused_parameters=self.zero_find_unused_parameters()) + find_unused_parameters=self.zero_find_unused_parameters(), + partition_grads=zero_stage == ZERO_OPTIMIZATION_GRADIENTS, + grad_hooks=self.zero_grad_hooks()) elif zero_stage == ZERO_OPTIMIZATION_WEIGHTS: print("Initializing ZeRO Stage 3") if dist.get_rank() == 0 else None from deepspeed.runtime.zero.stage3 import FP16_DeepSpeedZeroOptimizer_Stage3 @@ -971,18 +981,14 @@ def forward(self, *inputs, **kwargs): return loss def allreduce_gradients(self, bucket_size=MEMORY_OPT_ALLREDUCE_SIZE): - #Zero stage 2 communicates during non gradient accumulation boundaries as well + # ZeRO stage 2 communicates during non gradient accumulation boundaries as well if self.zero_optimization_partition_gradients(): self.optimizer.overlapping_partition_gradients_reduce_epilogue() - #Communicate only at gradient accumulation boundaries + # Communicate only at gradient accumulation boundaries elif self.is_gradient_accumulation_boundary(): - if self.zero_optimization_stage( - ) == ZERO_OPTIMIZATION_OPTIMIZER_STATES and self.zero_reduce_scatter(): - self.optimizer.reduce_scatter_gradients( - postscale_gradients=self.postscale_gradients(), - gradient_predivide_factor=self.gradient_predivide_factor(), - gradient_average=self.gradient_average) + if self.zero_optimization_stage() == ZERO_OPTIMIZATION_OPTIMIZER_STATES: + self.optimizer.reduce_gradients() else: self.buffered_allreduce_fallback(elements_per_buffer=bucket_size) @@ -1700,19 +1706,19 @@ def _save_checkpoint(self, save_dir, tag, client_state={}): # then instead just returns None. self._curr_ckpt_path = os.path.join(save_dir, tag) - state = dict( - module=self.module_state_dict(), - optimizer=self.optimizer.state_dict() - if self.optimizer and not self.zero_optimization() else None, - lr_scheduler=self.lr_scheduler.state_dict() - if self.lr_scheduler is not None else None, - csr_tensor_module_names=self.csr_tensor_module_names, - skipped_steps=self.skipped_steps, - global_steps=self.global_steps, - global_samples=self.global_samples, - dp_world_size=self.dp_world_size, - mp_world_size=self.mp_world_size, - ) + state = dict(module=self.module_state_dict(), + optimizer=self.optimizer.state_dict() + if self.optimizer and not self.zero_optimization() else None, + lr_scheduler=self.lr_scheduler.state_dict() + if self.lr_scheduler is not None else None, + csr_tensor_module_names=self.csr_tensor_module_names, + skipped_steps=self.skipped_steps, + global_steps=self.global_steps, + global_samples=self.global_samples, + dp_world_size=self.dp_world_size, + mp_world_size=self.mp_world_size, + ds_config=self.config, + ds_version=version) state.update(client_state) log_dist(message=f'Saving model checkpoint: {save_path}', ranks=[0]) @@ -1740,10 +1746,10 @@ def _copy_recovery_script(self, save_path): def _save_zero_checkpoint(self, save_path, tag): zero_checkpoint_name = self._get_zero_ckpt_name(save_path, tag) - zero_sd = dict( - optimizer_state_dict=self.optimizer.state_dict(), - param_shapes=self._get_param_shapes(), - ) + zero_sd = dict(optimizer_state_dict=self.optimizer.state_dict(), + param_shapes=self._get_param_shapes(), + ds_config=self.config, + ds_version=version) torch.save(zero_sd, zero_checkpoint_name) self._copy_recovery_script(save_path) logger.info('zero checkpoint saved {}'.format(zero_checkpoint_name)) diff --git a/deepspeed/runtime/zero/config.py b/deepspeed/runtime/zero/config.py index 9944116a2967..6251200aa227 100755 --- a/deepspeed/runtime/zero/config.py +++ b/deepspeed/runtime/zero/config.py @@ -182,3 +182,23 @@ def _initialize(self, zero_config_dict): zero_config_dict, ZERO_OPTIMIZATION_FIND_UNUSED_PARAMETERS, ZERO_OPTIMIZATION_FIND_UNUSED_PARAMETERS_DEFAULT) + + if self.stage >= ZERO_OPTIMIZATION_GRADIENTS: + # grad hooks are always enabled for stage 2 and above + self.grad_hooks = ZERO_OPTIMIZATION_GRAD_HOOKS_DEFAULT + + config_value = get_scalar_param(zero_config_dict, + ZERO_OPTIMIZATION_GRAD_HOOKS, + ZERO_OPTIMIZATION_GRAD_HOOKS_DEFAULT) + if config_value != self.grad_hooks: + logger.warning(f"ZeRO {ZERO_OPTIMIZATION_GRAD_HOOKS} is \ + always {ZERO_OPTIMIZATION_GRAD_HOOKS_DEFAULT} for \ + stage {ZERO_OPTIMIZATION_GRADIENTS} and above.") + else: + self.grad_hooks = get_scalar_param(zero_config_dict, + ZERO_OPTIMIZATION_GRAD_HOOKS, + ZERO_OPTIMIZATION_GRAD_HOOKS_DEFAULT) + + self.legacy_stage1 = get_scalar_param(zero_config_dict, + ZERO_OPTIMIZATION_LEGACY_STAGE1, + ZERO_OPTIMIZATION_LEGACY_STAGE1_DEFAULT) diff --git a/deepspeed/runtime/zero/constants.py b/deepspeed/runtime/zero/constants.py index 641e37754812..8f1fcbf2aede 100755 --- a/deepspeed/runtime/zero/constants.py +++ b/deepspeed/runtime/zero/constants.py @@ -122,6 +122,16 @@ ZERO_OPTIMIZATION_FIND_UNUSED_PARAMETERS = 'find_unused_parameters' ZERO_OPTIMIZATION_FIND_UNUSED_PARAMETERS_DEFAULT = False +# Enable grad hooks to reduce grads during backward pass, must be enabled for +# grad partitioning (ZeRO-2) but optional with optimizer partitioning (ZeRO-1). +ZERO_OPTIMIZATION_GRAD_HOOKS = "grad_hooks" +ZERO_OPTIMIZATION_GRAD_HOOKS_DEFAULT = True + +# Use deepspeed < v0.3.17 zero stage 1, kept for backwards compatability reasons +ZERO_OPTIMIZATION_LEGACY_STAGE1 = "legacy_stage1" +ZERO_OPTIMIZATION_LEGACY_STAGE1_DEFAULT = False + +#yapf: disable ZERO_OPTIMIZATION_DEFAULT = { ZERO_OPTIMIZATION_STAGE: ZERO_OPTIMIZATION_STAGE_DEFAULT, @@ -156,5 +166,9 @@ ZERO_OPTIMIZATION_GATHER_FP16_WEIGHTS_ON_MODEL_SAVE: ZERO_OPTIMIZATION_GATHER_FP16_WEIGHTS_ON_MODEL_SAVE_DEFAULT, ZERO_OPTIMIZATION_FIND_UNUSED_PARAMETERS: - ZERO_OPTIMIZATION_FIND_UNUSED_PARAMETERS_DEFAULT + ZERO_OPTIMIZATION_FIND_UNUSED_PARAMETERS_DEFAULT, + ZERO_OPTIMIZATION_GRAD_HOOKS: + ZERO_OPTIMIZATION_GRAD_HOOKS_DEFAULT, + ZERO_OPTIMIZATION_LEGACY_STAGE1: + ZERO_OPTIMIZATION_LEGACY_STAGE1_DEFAULT } diff --git a/deepspeed/runtime/zero/stage2.py b/deepspeed/runtime/zero/stage2.py index 9bf06a585bf1..aa0379d810bc 100755 --- a/deepspeed/runtime/zero/stage2.py +++ b/deepspeed/runtime/zero/stage2.py @@ -8,6 +8,7 @@ import math from torch._six import inf from torch.autograd import Variable +from packaging import version as pkg_version import collections @@ -17,6 +18,7 @@ from deepspeed.ops.adam import DeepSpeedCPUAdam from deepspeed.ops.op_builder import UtilsBuilder from deepspeed.utils import logger +from deepspeed.git_version_info import version #Toggle this to true to enable correctness test #with gradient partitioning and without @@ -96,7 +98,9 @@ def __init__(self, postscale_gradients=True, gradient_predivide_factor=1.0, gradient_accumulation_steps=1, - find_unused_parameters=False): + find_unused_parameters=False, + partition_grads=True, + grad_hooks=True): if dist.get_rank() == 0: logger.info(f"Reduce bucket size {reduce_bucket_size}") @@ -120,6 +124,12 @@ def __init__(self, self.flatten = util_ops.flatten self.unflatten = util_ops.unflatten + # ZeRO stage 1 (False) or 2 (True) + self.partition_gradients = partition_grads + + # Use backward hooks to reduce gradients + self.grad_hooks = grad_hooks + self.timers = timers self.reduce_scatter = reduce_scatter @@ -152,6 +162,8 @@ def __init__(self, self.micro_step_id = 0 self.find_unused_parameters = find_unused_parameters + self.extra_large_param_to_reduce = None + if self.reduce_scatter: assert not self.allreduce_always_fp32, "allreduce_always_fp32 is not yet supported with ZeRO-2 with reduce scatter enabled" assert self.gradient_predivide_factor == 1.0, "gradient_predivide_factor != 1.0 is not yet supported with ZeRO-2 with reduce scatter enabled" @@ -249,9 +261,9 @@ def __init__(self, partition_size = len(self.fp16_groups_flat[i]) / dist.get_world_size( group=self.dp_process_group) - params_in_partition, params_not_in_partition, first_offset = self.get_partition_info(self.fp16_groups[i], partition_size, partition_id) - self.partition_size.append(partition_size) + + params_in_partition, params_not_in_partition, first_offset = self.get_partition_info(self.fp16_groups[i], partition_size, partition_id) self.params_in_partition.append(params_in_partition) self.params_not_in_partition.append(params_not_in_partition) self.first_offset.append(first_offset) @@ -356,7 +368,8 @@ def __init__(self, self.reset_partition_gradient_structures() #creates backward hooks for gradient partitioning - self.create_reduce_and_remove_grad_hooks() + if self.grad_hooks: + self.create_reduce_and_remove_grad_hooks() # we may have a way of fusing dynamic scale. Do not support for now if self.dtype == torch.float or not dynamic_loss_scale: @@ -408,6 +421,22 @@ def initialize_optimizer_states(self): return + ######################################################################### + #################### ZeRO Stage 1 - reduce gradients #################### + ######################################################################### + + def reduce_gradients(self): + world_size = dist.get_world_size(self.dp_process_group) + my_rank = dist.get_rank(self.dp_process_group) + + if not self.grad_hooks: + for i, group in enumerate(self.fp16_groups): + for param in group: + self.reduce_ready_partitions_and_remove_grads(param, i) + + # reduce any pending grads in either hook/non-hook case + self.overlapping_partition_gradients_reduce_epilogue() + ######################################################################### #########################ZeRO Partition Gradients######################## ######################################################################### @@ -633,8 +662,11 @@ def reduce_independent_p_g_buckets_and_remove_grads(self, param, i): Gradient computed twice for this partition. \ Multiple gradient reduction is currently not supported" - #keeping the gradients contiguous to prevent memory fragmentation, and avoid flattening - if self.contiguous_gradients: + if param.numel() > self.reduce_bucket_size: + self.extra_large_param_to_reduce = param + + elif self.contiguous_gradients: + #keeping the gradients contiguous to prevent memory fragmentation, and avoid flattening new_grad_tensor = self.ipg_buffer[self.ipg_index].narrow( 0, self.elements_in_ipg_bucket, @@ -971,7 +1003,14 @@ def reduce_ipg_grads(self): stream = torch.cuda.current_stream() if self.contiguous_gradients: - self.average_tensor(self.ipg_buffer[self.ipg_index]) + if self.extra_large_param_to_reduce is not None: + assert len(self.params_in_ipg_bucket) == 1, "more than 1 param in ipg bucket, this shouldn't happen" + _, _, param_id = self.params_in_ipg_bucket[0] + assert self.get_param_id(self.extra_large_param_to_reduce) == param_id, "param in ipg bucket does not match extra-large param" + self.average_tensor(self.extra_large_param_to_reduce.grad.view(-1)) + self.extra_large_param_to_reduce = None + else: + self.average_tensor(self.ipg_buffer[self.ipg_index]) else: self.buffered_reduce_fallback( None, @@ -1006,7 +1045,8 @@ def reduce_ipg_grads(self): ##################################################################### def reduce_ready_partitions_and_remove_grads(self, param, i): - self.reduce_independent_p_g_buckets_and_remove_grads(param, i) + if self.partition_gradients or self.is_gradient_accumulation_boundary: + self.reduce_independent_p_g_buckets_and_remove_grads(param, i) def zero_reduced_gradients(self, partition_id, i): def are_all_related_partitions_reduced(params_id): @@ -1632,17 +1672,16 @@ def backward(self, loss, retain_graph=False): if self.cpu_offload: torch.cuda.current_stream().wait_stream(self.migration_stream) - #TODO: we need to revist this and remove the magic 4.5x multiplier here if self.contiguous_gradients: self.ipg_buffer = [] - buf_0 = torch.empty(int(self.reduce_bucket_size * 4.5), + buf_0 = torch.empty(int(self.reduce_bucket_size), dtype=self.dtype, device=torch.cuda.current_device()) self.ipg_buffer.append(buf_0) # Use double buffers to avoid data access conflict when overlap_comm is enabled. if self.overlap_comm: - buf_1 = torch.empty(int(self.reduce_bucket_size * 4.5), + buf_1 = torch.empty(int(self.reduce_bucket_size), dtype=self.dtype, device=torch.cuda.current_device()) self.ipg_buffer.append(buf_1) @@ -1740,6 +1779,8 @@ def state_dict(self): state_dict['zero_stage'] = ZERO_OPTIMIZATION_GRADIENTS state_dict['partition_count'] = self.partition_count + state_dict['ds_version'] = version + # Remove paddings for DP alignment to enable loading for other alignment values fp32_groups_without_padding = self._get_groups_without_padding( self.single_partition_of_fp32_groups) @@ -1859,6 +1900,18 @@ def load_state_dict(self, self.dynamic_loss_scale = state_dict_list[0]['dynamic_loss_scale'] self.overflow = state_dict_list[0]['overflow'] + # zero stage 1 mode + if not self.partition_gradients: + required_version = pkg_version.parse("0.3.16") + ckpt_version = state_dict_list[0].get("ds_version", False) + error_str = f"ZeRO stage 1 changed in {required_version} and is not backwards compatible " \ + "with older stage 1 checkpoints. If you'd like to load an old ZeRO-1 checkpoint " \ + "please set 'legacy_stage1': true in your zero config json. This old version of " \ + "stage 1 will be removed in v0.4.0." + + assert ckpt_version, f"Empty ds_version! {error_str}" + assert required_version <= pkg_version.parse(ckpt_version), f"Old version: {ckpt_version} {error_str}" + if load_optimizer_states: self._restore_base_optimizer_state(state_dict_list) diff --git a/requirements/requirements.txt b/requirements/requirements.txt index 43e488386866..afe7d231d4bf 100644 --- a/requirements/requirements.txt +++ b/requirements/requirements.txt @@ -5,3 +5,4 @@ tensorboardX==1.8 ninja numpy psutil +packaging From a2fb18aaafef0edbddcc256a8da4abb0e3e92c7e Mon Sep 17 00:00:00 2001 From: Jeff Rasley Date: Tue, 4 May 2021 23:36:02 +0000 Subject: [PATCH 2/6] fixes to support PP --- deepspeed/runtime/engine.py | 10 +++++++++- deepspeed/runtime/pipe/engine.py | 5 ++--- deepspeed/runtime/zero/stage2.py | 4 +++- 3 files changed, 14 insertions(+), 5 deletions(-) diff --git a/deepspeed/runtime/engine.py b/deepspeed/runtime/engine.py index eaf498f1345e..219bcfdf5090 100755 --- a/deepspeed/runtime/engine.py +++ b/deepspeed/runtime/engine.py @@ -776,6 +776,14 @@ def _configure_zero_optimizer(self, optimizer): elastic_checkpoint=self.zero_elastic_checkpoint(), mpu=self.mpu) elif zero_stage <= ZERO_OPTIMIZATION_GRADIENTS: + grad_hooks = self.zero_grad_hooks() + if isinstance(self.module, PipelineModule): + if grad_hooks: + logger.warning( + "Pipeline parallelism does not support backward grad reduction hooks, will be disabled." + ) + grad_hooks = False + optimizer = FP16_DeepSpeedZeroOptimizer( optimizer, timers=timers, @@ -796,7 +804,7 @@ def _configure_zero_optimizer(self, optimizer): gradient_accumulation_steps=self.gradient_accumulation_steps(), find_unused_parameters=self.zero_find_unused_parameters(), partition_grads=zero_stage == ZERO_OPTIMIZATION_GRADIENTS, - grad_hooks=self.zero_grad_hooks()) + grad_hooks=grad_hooks) elif zero_stage == ZERO_OPTIMIZATION_WEIGHTS: print("Initializing ZeRO Stage 3") if dist.get_rank() == 0 else None from deepspeed.runtime.zero.stage3 import FP16_DeepSpeedZeroOptimizer_Stage3 diff --git a/deepspeed/runtime/pipe/engine.py b/deepspeed/runtime/pipe/engine.py index d4e5e5edfe71..caad8c84d582 100644 --- a/deepspeed/runtime/pipe/engine.py +++ b/deepspeed/runtime/pipe/engine.py @@ -226,9 +226,8 @@ def _exec_reduce_tied_grads(self): def _exec_reduce_grads(self): self._force_grad_boundary = True - if self.is_data_parallel and self.pipeline_enable_backward_allreduce: - self.buffered_allreduce_fallback( - elements_per_buffer=MEMORY_OPT_ALLREDUCE_SIZE) + if self.pipeline_enable_backward_allreduce: + self.allreduce_gradients(bucket_size=MEMORY_OPT_ALLREDUCE_SIZE) self._force_grad_boundary = False def _reserve_pipe_buffers(self, num_buffers): diff --git a/deepspeed/runtime/zero/stage2.py b/deepspeed/runtime/zero/stage2.py index aa0379d810bc..f0ada14aac09 100755 --- a/deepspeed/runtime/zero/stage2.py +++ b/deepspeed/runtime/zero/stage2.py @@ -146,6 +146,8 @@ def __init__(self, self.partition_count = dist.get_world_size(group=self.dp_process_group) + self.is_gradient_accumulation_boundary = True + if mpu is None: self.model_parallel_group = None self.model_parallel_rank = 0 @@ -1902,7 +1904,7 @@ def load_state_dict(self, # zero stage 1 mode if not self.partition_gradients: - required_version = pkg_version.parse("0.3.16") + required_version = pkg_version.parse("0.3.17") ckpt_version = state_dict_list[0].get("ds_version", False) error_str = f"ZeRO stage 1 changed in {required_version} and is not backwards compatible " \ "with older stage 1 checkpoints. If you'd like to load an old ZeRO-1 checkpoint " \ From 10eda146c02a0fdb74622c098330c3f37623006f Mon Sep 17 00:00:00 2001 From: Jeff Rasley Date: Wed, 5 May 2021 17:29:56 +0000 Subject: [PATCH 3/6] update json docs to add grad_hooks param --- docs/_pages/config-json.md | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/docs/_pages/config-json.md b/docs/_pages/config-json.md index 3a3ee7f49b30..5303ab4d1207 100755 --- a/docs/_pages/config-json.md +++ b/docs/_pages/config-json.md @@ -352,6 +352,11 @@ Enabling and configuring ZeRO memory optimizations | --------------------------------------------------------------------------------------------------------------------------------------------------------------- | ------- | | Copies the gradients to a contiguous buffer as they are produced. Avoids memory fragmentation during backward pass. Only useful when running very large models. | `False` | +**grad_hooks**: [boolean] + +| Description | Default | +| ------------------------------------------------------------------------------------------------------------------------------------------ | ------- | +| For use with ZeRO stage 1, enable backward hooks to reduce gradients during the backward pass or wait until the end of the backward pass. | `True` | ***offload_param***: [dictionary] From df795cade1029dff80bb7323d0cb5ef4539a534e Mon Sep 17 00:00:00 2001 From: Jeff Rasley Date: Thu, 13 May 2021 20:57:57 +0000 Subject: [PATCH 4/6] add missing comma in merge confict fix --- deepspeed/runtime/zero/constants.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/deepspeed/runtime/zero/constants.py b/deepspeed/runtime/zero/constants.py index d631c85c68e4..a3f931e78381 100755 --- a/deepspeed/runtime/zero/constants.py +++ b/deepspeed/runtime/zero/constants.py @@ -164,7 +164,7 @@ ZERO_OPTIMIZATION_GATHER_FP16_WEIGHTS_ON_MODEL_SAVE: ZERO_OPTIMIZATION_GATHER_FP16_WEIGHTS_ON_MODEL_SAVE_DEFAULT, ZERO_OPTIMIZATION_IGNORE_UNUSED_PARAMETERS: - ZERO_OPTIMIZATION_IGNORE_UNUSED_PARAMETERS_DEFAULT + ZERO_OPTIMIZATION_IGNORE_UNUSED_PARAMETERS_DEFAULT, ZERO_OPTIMIZATION_GRAD_HOOKS: ZERO_OPTIMIZATION_GRAD_HOOKS_DEFAULT, ZERO_OPTIMIZATION_LEGACY_STAGE1: From a9ef5994eaee00ef29a1d636081194c2eab73877 Mon Sep 17 00:00:00 2001 From: Jeff Rasley Date: Fri, 14 May 2021 00:27:34 +0000 Subject: [PATCH 5/6] remove grad hooks, unify grad reduce from engine --- deepspeed/runtime/engine.py | 18 ++++++++++-------- deepspeed/runtime/zero/config.py | 16 ---------------- deepspeed/runtime/zero/constants.py | 7 ------- deepspeed/runtime/zero/stage1.py | 18 +++++++++++++----- deepspeed/runtime/zero/stage2.py | 10 +++------- 5 files changed, 26 insertions(+), 43 deletions(-) diff --git a/deepspeed/runtime/engine.py b/deepspeed/runtime/engine.py index 5b2f0b06084b..1ccd4f2b0298 100755 --- a/deepspeed/runtime/engine.py +++ b/deepspeed/runtime/engine.py @@ -800,15 +800,18 @@ def _configure_zero_optimizer(self, optimizer): max_elements_per_comm=self.zero_reduce_bucket_size(), dp_process_group=self.data_parallel_group, elastic_checkpoint=self.zero_elastic_checkpoint(), - mpu=self.mpu) + mpu=self.mpu, + postscale_gradients=self.postscale_gradients(), + gradient_predivide_factor=self.gradient_predivide_factor(), + gradient_predivide=self.gradient_predivide) elif zero_stage <= ZERO_OPTIMIZATION_GRADIENTS: - grad_hooks = self.zero_grad_hooks() + overlap_comm = self.zero_overlap_comm() if isinstance(self.module, PipelineModule): - if grad_hooks: + if overlap_comm: logger.warning( - "Pipeline parallelism does not support backward grad reduction hooks, will be disabled." + "Pipeline parallelism does not support overlapped communication, will be disabled." ) - grad_hooks = False + overlap_comm = False optimizer = FP16_DeepSpeedZeroOptimizer( optimizer, @@ -822,15 +825,14 @@ def _configure_zero_optimizer(self, optimizer): allgather_bucket_size=self.zero_allgather_bucket_size(), dp_process_group=self.data_parallel_group, reduce_scatter=self.zero_reduce_scatter(), - overlap_comm=self.zero_overlap_comm(), + overlap_comm=overlap_comm, cpu_offload=self.zero_cpu_offload(), mpu=self.mpu, postscale_gradients=self.postscale_gradients(), gradient_predivide_factor=self.gradient_predivide_factor(), gradient_accumulation_steps=self.gradient_accumulation_steps(), ignore_unused_parameters=self.zero_ignore_unused_parameters(), - partition_grads=zero_stage == ZERO_OPTIMIZATION_GRADIENTS, - grad_hooks=grad_hooks) + partition_grads=zero_stage == ZERO_OPTIMIZATION_GRADIENTS) elif zero_stage == ZERO_OPTIMIZATION_WEIGHTS: print("Initializing ZeRO Stage 3") if dist.get_rank() == 0 else None from deepspeed.runtime.zero.stage3 import FP16_DeepSpeedZeroOptimizer_Stage3 diff --git a/deepspeed/runtime/zero/config.py b/deepspeed/runtime/zero/config.py index 60132fea2f59..5fd5cd5109d3 100755 --- a/deepspeed/runtime/zero/config.py +++ b/deepspeed/runtime/zero/config.py @@ -183,22 +183,6 @@ def _initialize(self, zero_config_dict): ZERO_OPTIMIZATION_IGNORE_UNUSED_PARAMETERS, ZERO_OPTIMIZATION_IGNORE_UNUSED_PARAMETERS_DEFAULT) - if self.stage >= ZERO_OPTIMIZATION_GRADIENTS: - # grad hooks are always enabled for stage 2 and above - self.grad_hooks = ZERO_OPTIMIZATION_GRAD_HOOKS_DEFAULT - - config_value = get_scalar_param(zero_config_dict, - ZERO_OPTIMIZATION_GRAD_HOOKS, - ZERO_OPTIMIZATION_GRAD_HOOKS_DEFAULT) - if config_value != self.grad_hooks: - logger.warning(f"ZeRO {ZERO_OPTIMIZATION_GRAD_HOOKS} is \ - always {ZERO_OPTIMIZATION_GRAD_HOOKS_DEFAULT} for \ - stage {ZERO_OPTIMIZATION_GRADIENTS} and above.") - else: - self.grad_hooks = get_scalar_param(zero_config_dict, - ZERO_OPTIMIZATION_GRAD_HOOKS, - ZERO_OPTIMIZATION_GRAD_HOOKS_DEFAULT) - self.legacy_stage1 = get_scalar_param(zero_config_dict, ZERO_OPTIMIZATION_LEGACY_STAGE1, ZERO_OPTIMIZATION_LEGACY_STAGE1_DEFAULT) diff --git a/deepspeed/runtime/zero/constants.py b/deepspeed/runtime/zero/constants.py index a3f931e78381..8c15b4be6453 100755 --- a/deepspeed/runtime/zero/constants.py +++ b/deepspeed/runtime/zero/constants.py @@ -120,11 +120,6 @@ ZERO_OPTIMIZATION_IGNORE_UNUSED_PARAMETERS = 'ignore_unused_parameters' ZERO_OPTIMIZATION_IGNORE_UNUSED_PARAMETERS_DEFAULT = True -# Enable grad hooks to reduce grads during backward pass, must be enabled for -# grad partitioning (ZeRO-2) but optional with optimizer partitioning (ZeRO-1). -ZERO_OPTIMIZATION_GRAD_HOOKS = "grad_hooks" -ZERO_OPTIMIZATION_GRAD_HOOKS_DEFAULT = True - # Use deepspeed < v0.3.17 zero stage 1, kept for backwards compatability reasons ZERO_OPTIMIZATION_LEGACY_STAGE1 = "legacy_stage1" ZERO_OPTIMIZATION_LEGACY_STAGE1_DEFAULT = False @@ -165,8 +160,6 @@ ZERO_OPTIMIZATION_GATHER_FP16_WEIGHTS_ON_MODEL_SAVE_DEFAULT, ZERO_OPTIMIZATION_IGNORE_UNUSED_PARAMETERS: ZERO_OPTIMIZATION_IGNORE_UNUSED_PARAMETERS_DEFAULT, - ZERO_OPTIMIZATION_GRAD_HOOKS: - ZERO_OPTIMIZATION_GRAD_HOOKS_DEFAULT, ZERO_OPTIMIZATION_LEGACY_STAGE1: ZERO_OPTIMIZATION_LEGACY_STAGE1_DEFAULT } diff --git a/deepspeed/runtime/zero/stage1.py b/deepspeed/runtime/zero/stage1.py index dde8424ceaad..4cab2bc0501f 100755 --- a/deepspeed/runtime/zero/stage1.py +++ b/deepspeed/runtime/zero/stage1.py @@ -77,7 +77,10 @@ def __init__(self, allgather_size=500000000, clip_grad=0.0, max_elements_per_comm=5e8, - elastic_checkpoint=True): + elastic_checkpoint=True, + postscale_gradients=True, + gradient_predivide_factor=1.0, + gradient_average=True): # Load pre-built or JIT compile (un)flatten ops util_ops = UtilsBuilder().load() @@ -98,6 +101,10 @@ def __init__(self, self.verbose = verbose self.dp_process_group = dp_process_group + self.postscale_gradients = postscale_gradients + self.gradient_predivide_factor = gradient_predivide_factor + self.gradient_average = gradient_average + # TODO: automatically turn off if #params > some_limit self.all_gather_partitions = all_gather_partitions self.allgather_size = allgather_size @@ -575,10 +582,11 @@ def flatten_dense_tensors_sub_partition_aligned(self, flat_tensors = self.flatten(aligned_tensor_list) return flat_tensors - def reduce_scatter_gradients(self, - postscale_gradients, - gradient_predivide_factor, - gradient_average): + def reduce_gradients(self): + postscale_gradients = self.postscale_gradients + gradient_predivide_factor = self.gradient_predivide_factor + gradient_average = self.gradient_average + world_size = dist.get_world_size(group=self.dp_process_group) local_rank = dist.get_rank(group=self.dp_process_group) diff --git a/deepspeed/runtime/zero/stage2.py b/deepspeed/runtime/zero/stage2.py index 1d7c819ac44e..2b2b82f7258d 100755 --- a/deepspeed/runtime/zero/stage2.py +++ b/deepspeed/runtime/zero/stage2.py @@ -99,8 +99,7 @@ def __init__(self, gradient_predivide_factor=1.0, gradient_accumulation_steps=1, ignore_unused_parameters=True, - partition_grads=True, - grad_hooks=True): + partition_grads=True): if dist.get_rank() == 0: logger.info(f"Reduce bucket size {reduce_bucket_size}") @@ -127,9 +126,6 @@ def __init__(self, # ZeRO stage 1 (False) or 2 (True) self.partition_gradients = partition_grads - # Use backward hooks to reduce gradients - self.grad_hooks = grad_hooks - self.timers = timers self.reduce_scatter = reduce_scatter @@ -370,7 +366,7 @@ def __init__(self, self.reset_partition_gradient_structures() #creates backward hooks for gradient partitioning - if self.grad_hooks: + if self.partition_gradients or self.overlap_comm: self.create_reduce_and_remove_grad_hooks() # we may have a way of fusing dynamic scale. Do not support for now @@ -431,7 +427,7 @@ def reduce_gradients(self): world_size = dist.get_world_size(self.dp_process_group) my_rank = dist.get_rank(self.dp_process_group) - if not self.grad_hooks: + if not self.overlap_comm: for i, group in enumerate(self.fp16_groups): for param in group: self.reduce_ready_partitions_and_remove_grads(param, i) From 567339eb70f307e62f2f4ba3a2d9d9205e20c510 Mon Sep 17 00:00:00 2001 From: Jeff Rasley Date: Mon, 17 May 2021 19:58:44 +0000 Subject: [PATCH 6/6] fixes to support PP --- deepspeed/runtime/engine.py | 5 ++++- deepspeed/runtime/zero/stage1.py | 2 +- deepspeed/runtime/zero/stage2.py | 11 ++++++++++- 3 files changed, 15 insertions(+), 3 deletions(-) diff --git a/deepspeed/runtime/engine.py b/deepspeed/runtime/engine.py index 1ccd4f2b0298..5611e3d7fb4a 100755 --- a/deepspeed/runtime/engine.py +++ b/deepspeed/runtime/engine.py @@ -149,6 +149,8 @@ def __init__(self, # Configure distributed model self._configure_distributed_model(model) + self.pipeline_parallelism = isinstance(self.module, PipelineModule) + see_memory_usage(f"DeepSpeed Engine: After configure distributed model") # Configure wall clock timer @@ -1026,7 +1028,8 @@ def allreduce_gradients(self, bucket_size=MEMORY_OPT_ALLREDUCE_SIZE): # Communicate only at gradient accumulation boundaries elif self.is_gradient_accumulation_boundary(): if self.zero_optimization_stage() == ZERO_OPTIMIZATION_OPTIMIZER_STATES: - self.optimizer.reduce_gradients() + self.optimizer.reduce_gradients( + pipeline_parallel=self.pipeline_parallelism) else: self.buffered_allreduce_fallback(elements_per_buffer=bucket_size) diff --git a/deepspeed/runtime/zero/stage1.py b/deepspeed/runtime/zero/stage1.py index 4cab2bc0501f..7660c9917b84 100755 --- a/deepspeed/runtime/zero/stage1.py +++ b/deepspeed/runtime/zero/stage1.py @@ -582,7 +582,7 @@ def flatten_dense_tensors_sub_partition_aligned(self, flat_tensors = self.flatten(aligned_tensor_list) return flat_tensors - def reduce_gradients(self): + def reduce_gradients(self, pipeline_parallel=False): postscale_gradients = self.postscale_gradients gradient_predivide_factor = self.gradient_predivide_factor gradient_average = self.gradient_average diff --git a/deepspeed/runtime/zero/stage2.py b/deepspeed/runtime/zero/stage2.py index 2b2b82f7258d..e7601ef05b2a 100755 --- a/deepspeed/runtime/zero/stage2.py +++ b/deepspeed/runtime/zero/stage2.py @@ -423,10 +423,19 @@ def initialize_optimizer_states(self): #################### ZeRO Stage 1 - reduce gradients #################### ######################################################################### - def reduce_gradients(self): + def reduce_gradients(self, pipeline_parallel=False): world_size = dist.get_world_size(self.dp_process_group) my_rank = dist.get_rank(self.dp_process_group) + # with PP we must create ipg buffer, since backward is handled outside zero + if pipeline_parallel and self.contiguous_gradients: + self.ipg_buffer = [] + buf_0 = torch.empty(int(self.reduce_bucket_size), + dtype=self.dtype, + device=torch.cuda.current_device()) + self.ipg_buffer.append(buf_0) + self.ipg_index = 0 + if not self.overlap_comm: for i, group in enumerate(self.fp16_groups): for param in group: