From 47e170e38b68ffe4f9f3a263cbf6c260f2572235 Mon Sep 17 00:00:00 2001 From: Jeff Rasley Date: Thu, 25 Jun 2020 07:37:07 +0000 Subject: [PATCH 1/7] deepspeed distributed with mpi discovery --- deepspeed/__init__.py | 1 + deepspeed/pt/deepspeed_constants.py | 4 ++ deepspeed/pt/deepspeed_distributed.py | 60 +++++++++++++++++ deepspeed/pt/deepspeed_launch.py | 7 +- deepspeed/pt/deepspeed_light.py | 95 ++++++--------------------- 5 files changed, 90 insertions(+), 77 deletions(-) create mode 100644 deepspeed/pt/deepspeed_distributed.py diff --git a/deepspeed/__init__.py b/deepspeed/__init__.py index b1970ac4ebbe..7a6a51390847 100755 --- a/deepspeed/__init__.py +++ b/deepspeed/__init__.py @@ -8,6 +8,7 @@ from deepspeed.pt.log_utils import logger from deepspeed.pt.deepspeed_cuda import DeepSpeedTransformerLayer, DeepSpeedTransformerConfig from deepspeed.pt.deepspeed_config import DeepSpeedConfig +from deepspeed.pt.deepspeed_distributed import distributed_init, mpi_discovery import deepspeed.pt.deepspeed_checkpointing as checkpointing diff --git a/deepspeed/pt/deepspeed_constants.py b/deepspeed/pt/deepspeed_constants.py index f6afcedec92b..db5750b51125 100755 --- a/deepspeed/pt/deepspeed_constants.py +++ b/deepspeed/pt/deepspeed_constants.py @@ -266,3 +266,7 @@ # Tensorboard job name TENSORBOARD_JOB_NAME = "job_name" TENSORBOARD_JOB_NAME_DEFAULT = "DeepSpeedJobName" + +# DeepSpeed launcher +DEEPSPEED_LAUNCHER = "DEEPSPEED_LAUNCHER" +DEEPSPEED_LAUNCHER_DEFAULT = "TRUE" diff --git a/deepspeed/pt/deepspeed_distributed.py b/deepspeed/pt/deepspeed_distributed.py new file mode 100644 index 000000000000..5a39c2e05e26 --- /dev/null +++ b/deepspeed/pt/deepspeed_distributed.py @@ -0,0 +1,60 @@ +import torch +from deepspeed.pt.log_utils import logger + + +def distributed_init(dist_backend="nccl"): + """ + Initialize torch.distributed backend, potentially performing MPI discovery if needed + + Arguments: + dist_backend: torch distributed backend + """ + required_env = ["RANK", "WORLD_SIZE", "MASTER_ADDR", "MASTER_PORT", "LOCAL_RANK"] + if not all(map(lambda v: v in os.environ, required_env)): + logger.info( + "Not using the DeepSpeed or torch.distributed launchers, attempting to detect MPI environment..." + ) + mpi_discovery() + + if not dist.is_initialized(): + logger.info("Initializing torch distributed with backend: {}".format( + self.dist_backend)) + torch.distributed.init_process_group(backend=self.dist_backend) + + +def mpi_discovery(): + """ + Discovery MPI environment and map to relevant torch.distributed state + """ + from mpi4py import MPI + import subprocess + + comm = MPI.COMM_WORLD + rank = comm.Get_rank() + world_size = comm.Get_size() + + master_addr = None + if rank == 0: + hostname_cmd = ["hostname -I"] + result = subprocess.check_output(hostname_cmd, shell=True) + master_addr = result.decode('utf-8').split()[0] + master_addr = comm.bcast(master_addr, root=0) + + # Determine local rank by assuming hostnames are unique + proc_name = MPI.Get_processor_name() + all_procs = comm.allgather(proc_name) + local_rank = sum([i == proc_name for i in all_procs[:rank]]) + + os.environ['RANK'] = str(rank) + os.environ['WORLD_SIZE'] = str(world_size) + os.environ['LOCAL_RANK'] = local_rank + os.environ['MASTER_ADDR'] = master_addr + os.environ['MASTER_PORT'] = TORCH_DISTRIBUTED_DEFAULT_PORT + + logger.info( + "Discovered MPI settings of world_rank={}, local_rank={}, world_size={}, master_addr={}, master_port={}" + .format(os.environ['RANK'], + os.environ['LOCAL_RANK'], + os.environ['WORLD_SIZE'], + os.environ['MASTER_ADDR'], + os.environ['MASTER_PORT'])) diff --git a/deepspeed/pt/deepspeed_launch.py b/deepspeed/pt/deepspeed_launch.py index 55399194d23d..a3ce867b0319 100755 --- a/deepspeed/pt/deepspeed_launch.py +++ b/deepspeed/pt/deepspeed_launch.py @@ -11,6 +11,7 @@ from argparse import ArgumentParser, REMAINDER from deepspeed.pt.log_utils import logger +from deepspeed.deepspeed_constants import DEEPSPEED_LAUNCHER, DEEPSPEED_LAUNCHER_DEFAULT def parse_args(): @@ -96,16 +97,20 @@ def main(): current_env["CUDA_VISIBLE_DEVICES"])) exclusion_counts_per_node = None - # set PyTorch distributed related environmental variables + # Set torch distributed related environmental variables current_env["MASTER_ADDR"] = args.master_addr current_env["MASTER_PORT"] = str(args.master_port) current_env["WORLD_SIZE"] = str(dist_world_size) + # Set deepspeed launcher environment + current_env[DEEPSPEED_LAUNCHER] = DEEPSPEED_LAUNCHER_DEFAULT + processes = [] for local_rank in range(0, num_local_procs): # each process's rank dist_rank = global_rank_mapping[local_node][local_rank] current_env["RANK"] = str(dist_rank) + current_env["LOCAL_RANK"] = str(local_rank) # spawn the processes cmd = [ diff --git a/deepspeed/pt/deepspeed_light.py b/deepspeed/pt/deepspeed_light.py index 4ae66db8daaf..856a6e39eb6e 100755 --- a/deepspeed/pt/deepspeed_light.py +++ b/deepspeed/pt/deepspeed_light.py @@ -32,6 +32,8 @@ import deepspeed.pt.deepspeed_lr_schedules as lr_schedules from deepspeed.pt.deepspeed_csr_tensor import CSRTensor +from deepspeed.pt.deepspeed_distributed import distributed_init + MEMORY_OPT_ALLREDUCE_SIZE = 500000000 SUMMARY_WRITER_DIR_NAME = "JobId" @@ -102,7 +104,6 @@ def __init__(self, training_data=None, lr_scheduler=None, mpu=None, - dist_init_required=None, collate_fn=None, config_params=None): super(DeepSpeedLight, self).__init__() @@ -121,21 +122,8 @@ def __init__(self, self.warn_unscaled_loss = True self.config_params = config_params - if dist_init_required is None: - dist_init_required = not dist.is_initialized() - - self._mpi_check(args, dist_init_required) - - self.dist_backend = "nccl" - if dist_init_required: - if not dist.is_initialized(): - logger.info("Initializing torch distributed with backend: {}".format( - self.dist_backend)) - dist.init_process_group(backend=self.dist_backend) - else: - logger.warning( - "Was given dist_init_required=True but detected that torch" - "distributed was already initialized, cannot initialize twice.") + distributed_init() + self.local_rank = int(os.environ["LOCAL_RANK"]) self._do_args_sanity_check(args) self._configure_with_arguments(args, mpu) @@ -145,8 +133,6 @@ def __init__(self, if self.tensorboard_enabled(): self.summary_writer = self.get_summary_writer() - self._init_distributed(dist_init_required) - # Configure distributed model self._configure_distributed_model(model) @@ -181,51 +167,13 @@ def __init__(self, self.save_non_zero_checkpoint = False self.save_zero_checkpoint = False - self._configure_checkpointing(dist_init_required) + self._configure_checkpointing() if self.global_rank == 0: self._config.print('DeepSpeedLight configuration') if self.dump_state(): print_configuration(self, 'DeepSpeedLight') - def _mpi_check(self, args, dist_init_required): - if hasattr(args, 'deepspeed_mpi') and args.deepspeed_mpi: - from mpi4py import MPI - import subprocess - comm = MPI.COMM_WORLD - rank = comm.Get_rank() - world_size = comm.Get_size() - - master_addr = None - if rank == 0: - hostname_cmd = ["hostname -I"] - result = subprocess.check_output(hostname_cmd, shell=True) - master_addr = result.decode('utf-8').split()[0] - master_addr = comm.bcast(master_addr, root=0) - - # Determine local rank by assuming hostnames are unique - proc_name = MPI.Get_processor_name() - all_procs = comm.allgather(proc_name) - local_rank = sum([i == proc_name for i in all_procs[:rank]]) - - os.environ['RANK'] = str(rank) - os.environ['WORLD_SIZE'] = str(world_size) - args.local_rank = local_rank - os.environ['MASTER_ADDR'] = master_addr - os.environ['MASTER_PORT'] = TORCH_DISTRIBUTED_DEFAULT_PORT - - logger.info( - "Discovered MPI settings of world_rank={}, local_rank={}, world_size={}, master_addr={}, master_port={}" - .format(os.environ['RANK'], - args.local_rank, - os.environ['WORLD_SIZE'], - os.environ['MASTER_ADDR'], - os.environ['MASTER_PORT'])) - - if not dist_init_required and dist.is_initialized(): - assert dist.get_rank() == rank, "MPI rank {} does not match torch rank {}".format(rank, dist.get_rank()) - assert dist.get_world_size() == world_size, "MPI world size {} does not match torch world size {}".format(world_size, dist.get_world_size()) - def tensorboard_enabled(self): return self._config.tensorboard_enabled @@ -360,7 +308,7 @@ def _configure_lr_scheduler(self, client_lr_scheduler): self.lr_scheduler = client_lr_scheduler logger.info(f'DeepSpeed LR Scheduler = {self.lr_scheduler}') - def _configure_checkpointing(self, dist_init_required): + def _configure_checkpointing(self): dp_rank = self.global_rank if self.mpu: @@ -393,19 +341,6 @@ def _scheduler_from_config(self, optimizer): else: return None - def _init_distributed(self, dist_init_required): - if self.local_rank >= 0: - torch.cuda.set_device(self.local_rank) - self.device = torch.device("cuda", self.local_rank) - self.world_size = dist.get_world_size() - self.global_rank = dist.get_rank() - logger.info("Set device to local rank {} within node.".format( - self.local_rank)) - else: - self.world_size = 1 - self.global_rank = 0 - self.device = torch.device("cuda") - # Configure based on command line arguments def _configure_with_arguments(self, args, mpu): self.local_rank = args.local_rank if hasattr(args, 'local_rank') else 0 @@ -450,10 +385,23 @@ def _do_sanity_check(self): 'DeepSpeed {} optimizer requires dynamic loss scaling'.format(self.optimizer_name()) def _configure_distributed_model(self, model): + if self.local_rank >= 0: + torch.cuda.set_device(self.local_rank) + self.device = torch.device("cuda", self.local_rank) + self.world_size = dist.get_world_size() + self.global_rank = dist.get_rank() + logger.info("Set device to local rank {} within node.".format( + self.local_rank)) + else: + self.world_size = 1 + self.global_rank = 0 + self.device = torch.device("cuda") + self.module = model if self.fp16_enabled(): self.module.half() self.module.to(self.device) + if self.mpu is None: self.data_parallel_group = _initialize_parameter_parallel_groups() self.dp_world_size = dist.get_world_size() @@ -467,11 +415,6 @@ def _configure_distributed_model(self, model): if torch.is_tensor(p): dist.broadcast(p, src_rank, group=self.data_parallel_group) - # TODO: support new AMP optimizer - # self.module.half() - # self.module.to(self.local_rank) - #self.module, self.optimizer = amp.initialize(self.module, self.optimizer, opt_level="O2") - # Configure optimizer def _configure_optimizer(self, client_optimizer, model_parameters): if client_optimizer is not None: From 4e36fc3e818eccb341286157010fe06bfd5554cf Mon Sep 17 00:00:00 2001 From: Jeff Rasley Date: Thu, 25 Jun 2020 07:39:11 +0000 Subject: [PATCH 2/7] remove env var --- deepspeed/pt/deepspeed_constants.py | 4 ---- deepspeed/pt/deepspeed_launch.py | 3 --- 2 files changed, 7 deletions(-) diff --git a/deepspeed/pt/deepspeed_constants.py b/deepspeed/pt/deepspeed_constants.py index db5750b51125..f6afcedec92b 100755 --- a/deepspeed/pt/deepspeed_constants.py +++ b/deepspeed/pt/deepspeed_constants.py @@ -266,7 +266,3 @@ # Tensorboard job name TENSORBOARD_JOB_NAME = "job_name" TENSORBOARD_JOB_NAME_DEFAULT = "DeepSpeedJobName" - -# DeepSpeed launcher -DEEPSPEED_LAUNCHER = "DEEPSPEED_LAUNCHER" -DEEPSPEED_LAUNCHER_DEFAULT = "TRUE" diff --git a/deepspeed/pt/deepspeed_launch.py b/deepspeed/pt/deepspeed_launch.py index a3ce867b0319..b9449d228b2d 100755 --- a/deepspeed/pt/deepspeed_launch.py +++ b/deepspeed/pt/deepspeed_launch.py @@ -102,9 +102,6 @@ def main(): current_env["MASTER_PORT"] = str(args.master_port) current_env["WORLD_SIZE"] = str(dist_world_size) - # Set deepspeed launcher environment - current_env[DEEPSPEED_LAUNCHER] = DEEPSPEED_LAUNCHER_DEFAULT - processes = [] for local_rank in range(0, num_local_procs): # each process's rank From 434d61698eb56ebc6a5d886bc418333097909596 Mon Sep 17 00:00:00 2001 From: Jeff Rasley Date: Thu, 25 Jun 2020 07:40:16 +0000 Subject: [PATCH 3/7] remove import --- deepspeed/pt/deepspeed_launch.py | 1 - 1 file changed, 1 deletion(-) diff --git a/deepspeed/pt/deepspeed_launch.py b/deepspeed/pt/deepspeed_launch.py index b9449d228b2d..2ea3ebdc6317 100755 --- a/deepspeed/pt/deepspeed_launch.py +++ b/deepspeed/pt/deepspeed_launch.py @@ -11,7 +11,6 @@ from argparse import ArgumentParser, REMAINDER from deepspeed.pt.log_utils import logger -from deepspeed.deepspeed_constants import DEEPSPEED_LAUNCHER, DEEPSPEED_LAUNCHER_DEFAULT def parse_args(): From 51b88b07d703f8be6b86ed4229811cf818056f20 Mon Sep 17 00:00:00 2001 From: Jeff Rasley Date: Thu, 25 Jun 2020 07:43:49 +0000 Subject: [PATCH 4/7] add warning about dist_init_required --- deepspeed/__init__.py | 9 ++++++--- 1 file changed, 6 insertions(+), 3 deletions(-) diff --git a/deepspeed/__init__.py b/deepspeed/__init__.py index 7a6a51390847..199d7c0b60ef 100755 --- a/deepspeed/__init__.py +++ b/deepspeed/__init__.py @@ -63,8 +63,7 @@ def initialize(args, mpu: Optional: A model parallelism unit object that implements get_{model,data}_parallel_{rank,group,world_size}() - dist_init_required: Optional: None will auto-initialize torch.distributed if needed, - otherwise the user can force it to be initialized or not via boolean. + dist_init_required: deprecated argument, torch.distributed will be auto-initialized if needed collate_fn: Optional: Merges a list of samples to form a mini-batch of Tensor(s). Used when using batched loading from a @@ -91,6 +90,11 @@ def initialize(args, __git_branch__), ) + if dist_init_required is not None: + logger.warning( + "deepspeed.initialize argument of 'dist_init_required' is deprecated and not used, torch.distributed will be auto-initialized if needed." + ) + engine = DeepSpeedLight(args=args, model=model, optimizer=optimizer, @@ -98,7 +102,6 @@ def initialize(args, training_data=training_data, lr_scheduler=lr_scheduler, mpu=mpu, - dist_init_required=dist_init_required, collate_fn=collate_fn, config_params=config_params) From 309bbb8c1c09f8b6f1e55ad782006ff5f1327e8a Mon Sep 17 00:00:00 2001 From: Jeff Rasley Date: Thu, 25 Jun 2020 08:02:42 +0000 Subject: [PATCH 5/7] fixes --- deepspeed/pt/deepspeed_distributed.py | 6 ++++-- deepspeed/pt/deepspeed_light.py | 1 - 2 files changed, 4 insertions(+), 3 deletions(-) diff --git a/deepspeed/pt/deepspeed_distributed.py b/deepspeed/pt/deepspeed_distributed.py index 5a39c2e05e26..443b207f5c91 100644 --- a/deepspeed/pt/deepspeed_distributed.py +++ b/deepspeed/pt/deepspeed_distributed.py @@ -1,5 +1,7 @@ +import os import torch from deepspeed.pt.log_utils import logger +from deepspeed.pt.deepspeed_constants import TORCH_DISTRIBUTED_DEFAULT_PORT def distributed_init(dist_backend="nccl"): @@ -16,7 +18,7 @@ def distributed_init(dist_backend="nccl"): ) mpi_discovery() - if not dist.is_initialized(): + if not torch.distributed.is_initialized(): logger.info("Initializing torch distributed with backend: {}".format( self.dist_backend)) torch.distributed.init_process_group(backend=self.dist_backend) @@ -47,7 +49,7 @@ def mpi_discovery(): os.environ['RANK'] = str(rank) os.environ['WORLD_SIZE'] = str(world_size) - os.environ['LOCAL_RANK'] = local_rank + os.environ['LOCAL_RANK'] = str(local_rank) os.environ['MASTER_ADDR'] = master_addr os.environ['MASTER_PORT'] = TORCH_DISTRIBUTED_DEFAULT_PORT diff --git a/deepspeed/pt/deepspeed_light.py b/deepspeed/pt/deepspeed_light.py index 856a6e39eb6e..dbf2b76f4a14 100755 --- a/deepspeed/pt/deepspeed_light.py +++ b/deepspeed/pt/deepspeed_light.py @@ -26,7 +26,6 @@ from deepspeed.pt.deepspeed_dataloader import DeepSpeedDataLoader from deepspeed.pt.deepspeed_constants import \ ROUTE_TRAIN, ROUTE_PREDICT, ROUTE_EVAL, \ - TORCH_DISTRIBUTED_DEFAULT_PORT, \ ZERO_OPTIMIZATION_OPTIMIZER_STATES, ZERO_OPTIMIZATION_GRADIENTS import deepspeed.pt.deepspeed_lr_schedules as lr_schedules From b510f11e1bb3b6676209207404ace5af87268e6e Mon Sep 17 00:00:00 2001 From: Jeff Rasley Date: Thu, 25 Jun 2020 08:15:29 +0000 Subject: [PATCH 6/7] add all torch distributed environ vars --- tests/unit/common.py | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/tests/unit/common.py b/tests/unit/common.py index 5cea6d2f0f76..9b5ede9c76bd 100644 --- a/tests/unit/common.py +++ b/tests/unit/common.py @@ -33,6 +33,10 @@ def dist_init(local_rank, num_procs, *func_args, **func_kwargs): """Initialize torch.distributed and execute the user function. """ os.environ['MASTER_ADDR'] = '127.0.0.1' os.environ['MASTER_PORT'] = '29500' + os.environ['LOCAL_RANK'] = str(local_rank) + # multi-node tests are not supported, local_rank == rank + os.environ['RANK'] = str(local_rank) + os.environ['WORLD_SIZE'] = str(num_procs) dist.init_process_group(backend=backend, init_method='env://', rank=local_rank, From e0e616f6de9d197b81d900744cedc806b3dbd209 Mon Sep 17 00:00:00 2001 From: Jeff Rasley Date: Thu, 25 Jun 2020 17:30:11 +0000 Subject: [PATCH 7/7] allow ability to turn off mpi discovery and change name --- deepspeed/__init__.py | 2 +- deepspeed/pt/deepspeed_distributed.py | 7 ++++--- deepspeed/pt/deepspeed_light.py | 5 +++-- 3 files changed, 8 insertions(+), 6 deletions(-) diff --git a/deepspeed/__init__.py b/deepspeed/__init__.py index 199d7c0b60ef..16b2536026cd 100755 --- a/deepspeed/__init__.py +++ b/deepspeed/__init__.py @@ -8,7 +8,7 @@ from deepspeed.pt.log_utils import logger from deepspeed.pt.deepspeed_cuda import DeepSpeedTransformerLayer, DeepSpeedTransformerConfig from deepspeed.pt.deepspeed_config import DeepSpeedConfig -from deepspeed.pt.deepspeed_distributed import distributed_init, mpi_discovery +from deepspeed.pt.deepspeed_distributed import init_distributed, mpi_discovery import deepspeed.pt.deepspeed_checkpointing as checkpointing diff --git a/deepspeed/pt/deepspeed_distributed.py b/deepspeed/pt/deepspeed_distributed.py index 443b207f5c91..ce2e4092639d 100644 --- a/deepspeed/pt/deepspeed_distributed.py +++ b/deepspeed/pt/deepspeed_distributed.py @@ -4,15 +4,16 @@ from deepspeed.pt.deepspeed_constants import TORCH_DISTRIBUTED_DEFAULT_PORT -def distributed_init(dist_backend="nccl"): +def init_distributed(dist_backend="nccl", auto_mpi_discovery=True): """ Initialize torch.distributed backend, potentially performing MPI discovery if needed Arguments: - dist_backend: torch distributed backend + dist_backend: torch distributed backend, e.g., nccl, mpi, gloo + auto_mpi_discovery: if distributed environment variables are not set, attempt to discover them from MPI """ required_env = ["RANK", "WORLD_SIZE", "MASTER_ADDR", "MASTER_PORT", "LOCAL_RANK"] - if not all(map(lambda v: v in os.environ, required_env)): + if auto_mpi_discovery and not all(map(lambda v: v in os.environ, required_env)): logger.info( "Not using the DeepSpeed or torch.distributed launchers, attempting to detect MPI environment..." ) diff --git a/deepspeed/pt/deepspeed_light.py b/deepspeed/pt/deepspeed_light.py index dbf2b76f4a14..d2d3badafbde 100755 --- a/deepspeed/pt/deepspeed_light.py +++ b/deepspeed/pt/deepspeed_light.py @@ -31,7 +31,7 @@ import deepspeed.pt.deepspeed_lr_schedules as lr_schedules from deepspeed.pt.deepspeed_csr_tensor import CSRTensor -from deepspeed.pt.deepspeed_distributed import distributed_init +from deepspeed.pt.deepspeed_distributed import init_distributed MEMORY_OPT_ALLREDUCE_SIZE = 500000000 SUMMARY_WRITER_DIR_NAME = "JobId" @@ -121,7 +121,8 @@ def __init__(self, self.warn_unscaled_loss = True self.config_params = config_params - distributed_init() + # Initialize torch distributed backend + init_distributed() self.local_rank = int(os.environ["LOCAL_RANK"]) self._do_args_sanity_check(args)