From 647f42383c30a5803d29cbc8ffcf860d6ac03967 Mon Sep 17 00:00:00 2001 From: Saeed Maleki Date: Fri, 6 Dec 2019 21:22:31 +0000 Subject: [PATCH] adasum integration --- PyTorch/LanguageModeling/BERT/optimization.py | 359 +++++++++++++++++- .../LanguageModeling/BERT/run_pretraining.py | 65 +++- PyTorch/LanguageModeling/BERT/schedulers.py | 3 +- PyTorch/LanguageModeling/BERT/utils.py | 11 +- 4 files changed, 410 insertions(+), 28 deletions(-) diff --git a/PyTorch/LanguageModeling/BERT/optimization.py b/PyTorch/LanguageModeling/BERT/optimization.py index ac5b64f94..21ccd6750 100755 --- a/PyTorch/LanguageModeling/BERT/optimization.py +++ b/PyTorch/LanguageModeling/BERT/optimization.py @@ -18,6 +18,7 @@ import math import torch +import horovod.torch as hvd from torch.optim import Optimizer from torch.optim.optimizer import required from torch.nn.utils import clip_grad_norm_ @@ -255,7 +256,7 @@ def step(self, closure=None): # if self.max_steps != -1: schedule_fct = SCHEDULES[self.schedule] lr_scheduled = self.learning_rate * schedule_fct(self.step_count / self.max_steps, self.warmup) - if torch.distributed.get_rank() == 0: + if hvd.rank() == 0 :#torch.distributed.get_rank() == 0: print("Step {} LR {}".format(self.step_count, lr_scheduled)) # else: # lr_scheduled = self.learning_rate @@ -393,3 +394,359 @@ def step(self, closure=None): return loss +from contextlib import contextmanager + +def find_duplicates(lst): + seen = set() + dups = set() + for el in lst: + if el in seen: + dups.add(el) + seen.add(el) + return dups + +class DistributedOptimizer(torch.optim.Optimizer): + def __init__(self, optimizer, compression, + backward_passes_per_step=1): + params = [p + for param_group in optimizer.param_groups + for p in param_group['params']] + super(self.__class__, self).__init__(params, dict(DistributedOptimizer.__dict__)) + self.optimizer = optimizer + self._compression = compression + + named_parameters = [('allreduce.noname.%s' % i, v) + for param_group in self.param_groups + for i, v in enumerate(param_group['params'])] + + # make sure that named_parameters are tuples + if any([not isinstance(p, tuple) for p in named_parameters]): + raise ValueError('named_parameters should be a sequence of ' + 'tuples (name, parameter), usually produced by ' + 'model.named_parameters().') + + dups = find_duplicates([k for k, _ in named_parameters]) + if len(dups) > 0: + raise ValueError('Parameter names in named_parameters must be unique. ' + 'Found duplicates: %s' % ', '.join(dups)) + + all_param_ids = {id(v) + for param_group in self.param_groups + for v in param_group['params']} + named_param_ids = {id(v) for k, v in named_parameters} + unnamed_param_ids = all_param_ids - named_param_ids + if len(unnamed_param_ids): + raise ValueError('named_parameters was specified, but one or more model ' + 'parameters were not named. Python object ids: ' + '%s' % ', '.join(str(id) for id in unnamed_param_ids)) + + self._parameter_names = {v: k for k, v in sorted(named_parameters)} + self.backward_passes_per_step = backward_passes_per_step + self._allreduce_delay = {v: self.backward_passes_per_step + for _, v in sorted(named_parameters)} + self._handles = {} + self._grad_accs = [] + self._requires_update = set() + self._synchronized = False + self._should_synchronize = True + if hvd.size() > 1: + self._register_hooks() + + def set_backward_passes_per_step(self, passes): + self.backward_passes_per_step = passes + for p in self._allreduce_delay: + self._allreduce_delay[p] = self.backward_passes_per_step + + def _register_hooks(self): + for param_group in self.param_groups: + for p in param_group['params']: + if p.requires_grad: + p.grad = p.data.new(p.size()).zero_() + self._requires_update.add(p) + p_tmp = p.expand_as(p) + grad_acc = p_tmp.grad_fn.next_functions[0][0] + grad_acc.register_hook(self._make_hook(p)) + self._grad_accs.append(grad_acc) + + def _allreduce_grad_async(self, p): + name = self._parameter_names.get(p) + tensor = p.grad + tensor_compressed, ctx = self._compression.compress(tensor) + handle = hvd.allreduce_async_(tensor_compressed, op=hvd.Adasum, name=name) + return handle, ctx + + def _make_hook(self, p): + def hook(*ignore): + if p in self._handles and self._handles[p][0] is not None: + if self._allreduce_delay[p] <= 0: + raise AssertionError( + "Gradients were computed more than " + "backward_passes_per_step times before call " + "to step(). Increase backward_passes_per_step to " + "accumulate gradients locally.") + assert not p.grad.requires_grad + assert self._allreduce_delay[p] > 0 + handle, ctx = None, None + self._allreduce_delay[p] -= 1 + if self._allreduce_delay[p] == 0: + handle, ctx = self._allreduce_grad_async(p) + self._handles[p] = (handle, ctx) + return hook + + def synchronize(self): + missing_p = self._requires_update - set(self._handles.keys()) + for p in missing_p: + handle, ctx = self._allreduce_grad_async(p) + self._handles[p] = (handle, ctx) + + for p, value in self._handles.items(): + handle, ctx = value + if handle is None: + handle, ctx = self._allreduce_grad_async(p) + self._handles[p] = (handle, ctx) + for p, (handle, _) in self._handles.items(): + output = hvd.synchronize(handle) + self._allreduce_delay[p] = self.backward_passes_per_step + p.grad.copy_(self._compression.decompress(output, ctx).data) + self._handles.clear() + + self._synchronized = True + + @contextmanager + def skip_synchronize(self): + """ + A context manager used to specify that optimizer.step() should + not perform synchronization. + + It's typically used in a following pattern: + + .. code-block:: python + + optimizer.synchronize() + with optimizer.skip_synchronize(): + optimizer.step() + """ + self._should_synchronize = False + try: + yield + finally: + self._should_synchronize = True + + def step(self, closure=None): + if self._should_synchronize: + if self._synchronized: + warnings.warn("optimizer.step() called without " + "optimizer.skip_synchronize() context after " + "optimizer.synchronize(). This can cause training " + "slowdown. You may want to consider using " + "optimizer.skip_synchronize() context if you use " + "optimizer.synchronize() in your code.") + self.synchronize() + self._synchronized = False + return self.optimizer.step(closure) + #return super(self.__class__, self).step(closure) + + def zero_grad(self): + if self._handles: + raise AssertionError("optimizer.zero_grad() was called after loss.backward() " + "but before optimizer.step() or optimizer.synchronize(). " + "This is prohibited as it can cause a race condition.") + return self.optimizer.zero_grad() + #return super(self.__class__, self).zero_grad() + +from apex.fp16_utils.loss_scaler import DynamicLossScaler +class DistributedAdasumOptimizer(torch.optim.Optimizer): + def __init__(self, optimizer, compression, + backward_passes_per_step=1): + params = [p + for param_group in optimizer.param_groups + for p in param_group['params']] + super(self.__class__, self).__init__(params, dict(DistributedAdasumOptimizer.__dict__)) + + self._compression = compression + self.optimizer = optimizer + + named_parameters = [('allreduce.noname.%s' % i, v) + for param_group in self.param_groups + for i, v in enumerate(param_group['params'])] + + # make sure that named_parameters are tuples + if any([not isinstance(p, tuple) for p in named_parameters]): + raise ValueError('named_parameters should be a sequence of ' + 'tuples (name, parameter), usually produced by ' + 'model.named_parameters().') + + dups = find_duplicates([k for k, _ in named_parameters]) + if len(dups) > 0: + raise ValueError('Parameter names in named_parameters must be unique. ' + 'Found duplicates: %s' % ', '.join(dups)) + + all_param_ids = {id(v) + for param_group in self.param_groups + for v in param_group['params']} + named_param_ids = {id(v) for k, v in named_parameters} + unnamed_param_ids = all_param_ids - named_param_ids + if len(unnamed_param_ids): + raise ValueError('named_parameters was specified, but one or more model ' + 'parameters were not named. Python object ids: ' + '%s' % ', '.join(str(id) for id in unnamed_param_ids)) + + self._parameter_names = {v: k for k, v in sorted(named_parameters)} + self.backward_passes_per_step = backward_passes_per_step + self._allreduce_delay = {v: self.backward_passes_per_step + for _, v in sorted(named_parameters)} + self._handles = {} + self._grad_accs = [] + self._requires_update = set() + self._synchronized = False + self._should_synchronize = True + + self._starting_models = { + p : torch.zeros_like(p, requires_grad=False) + for _, p in named_parameters + } + + self._scalers = { + p : DynamicLossScaler() + for _, p in named_parameters + } + + self._is_first = True + + #self._register_hooks() + + def set_backward_passes_per_step(self, passes): + self.backward_passes_per_step = passes + for p in self._allreduce_delay: + self._allreduce_delay[p] = self.backward_passes_per_step + + def _register_hooks(self): + for param_group in self.param_groups: + for p in param_group['params']: + if p.requires_grad: + p.grad = p.data.new(p.size()).zero_() + self._requires_update.add(p) + p_tmp = p.expand_as(p) + grad_acc = p_tmp.grad_fn.next_functions[0][0] + grad_acc.register_hook(self._make_hook(p)) + self._grad_accs.append(grad_acc) + + def _allreduce_grad_async(self, p): + # Delta optimizer implements this logic: + # start = current.copy() + # step() -> computes 'current - \alpha.f(g)' where f is + # optimizer logic and g is the gradient + # delta = current-start + # allreduce_(delta) + # start += delta + # current = start + # In order to suppport this logic using function hook to improve performance, + # we do: + # delta = (start - \alpha.f(g)) - start + # = -\alpha.f(g) + # set start to zero and step computes -\alpha.f(g) + # where f is the underlying optimizer logic + + name = self._parameter_names.get(p) + start = self._starting_models[p] + + stashed_params = [] + for group in self.param_groups: + stashed_params.append(group['params']) + # only want to step on p + if any([p is v for v in group['params']]): + group['params'] = [p] + else: + group['params'] = [] + + start.data.copy_(p) + + self.optimizer.step() + + # compute delta = curr - start + p.data.sub_(start) + + # allreduce as before + scaler = self._scalers[p] + p.data.mul_(scaler.loss_scale) + tensor_compressed, ctx = self._compression.compress(p) + handle = hvd.allreduce_async_(tensor_compressed.data, name=name, op=hvd.Adasum) + + # reset stashed parameters + for stashed, group in zip(stashed_params, self.param_groups): + group['params'] = stashed + + return handle, ctx + + def _make_hook(self, p): + def hook(*ignore): + if p in self._handles and self._handles[p][0] is not None: + if self._allreduce_delay[p] <= 0: + raise AssertionError( + "Gradients were computed more than " + "backward_passes_per_step times before call " + "to step(). Increase backward_passes_per_step to " + "accumulate gradients locally.") + assert not p.grad.requires_grad + assert self._allreduce_delay[p] > 0 + handle, ctx = None, None + self._allreduce_delay[p] -= 1 + if self._allreduce_delay[p] == 0: + handle, ctx = self._allreduce_grad_async(p) + self._handles[p] = (handle, ctx) + return hook + + def synchronize(self): + pass + + @contextmanager + def skip_synchronize(self): + raise AssertionError("Skipping synchronization is not supported when using Adasum optimizer.") + + def step(self, closure=None): + loss = None + if closure is not None: + loss = closure() + + if self._is_first: + self._is_first = False + for group in self.param_groups: + for p in group['params']: + self._starting_models[p].data.copy_(p.data) + + self.optimizer.step() + + handles = [] + for group in self.param_groups: + for p in group['params']: + name = self._parameter_names.get(p) + #scaler = self._scalers[p] + start = self._starting_models[p] + p.data.sub_(start) + #p.data.mul_(scaler.loss_scale) + #tensor_compressed, ctx = self._compression.compress(p) + tensor_compressed, ctx = p, None + handle = hvd.allreduce_async_(tensor_compressed.data, name=name, op=hvd.Adasum) + handles.append((handle, p, ctx)) + + for handle, p, ctx in handles: + start = self._starting_models[p] + #scaler = self._scalers[p] + delta = hvd.synchronize(handle) + #has_overflow = not (torch.isfinite(delta.data).all().item()) + #if not has_overflow: + # delta = self._compression.decompress(delta, ctx) + # delta.data.div_(scaler.loss_scale) + # start.data.add_(delta.data) + start.data.add_(delta.data) + p.data.copy_(start) + #scaler.update_scale(has_overflow) + + return loss + + def zero_grad(self): + if self._handles: + raise AssertionError("optimizer.zero_grad() was called after loss.backward() " + "but before optimizer.step() or optimizer.synchronize(). " + "This is prohibited as it can cause a race condition.") + return self.optimizer.zero_grad() diff --git a/PyTorch/LanguageModeling/BERT/run_pretraining.py b/PyTorch/LanguageModeling/BERT/run_pretraining.py index fe926f3fe..cef6bbb11 100755 --- a/PyTorch/LanguageModeling/BERT/run_pretraining.py +++ b/PyTorch/LanguageModeling/BERT/run_pretraining.py @@ -30,6 +30,10 @@ import h5py from tqdm import tqdm, trange import os +import horovod.torch as hvd +hvd.init() +os.environ['CUDA_VISIBLE_DEVICES'] = str(hvd.local_rank()) + import numpy as np import torch from torch.utils.data import DataLoader, RandomSampler, SequentialSampler, Dataset @@ -41,6 +45,7 @@ from tokenization import BertTokenizer from modeling import BertForPreTraining, BertConfig from optimization import BertLAMB +from apex.optimizers import FusedAdam from file_utils import PYTORCH_PRETRAINED_BERT_CACHE from utils import is_main_process @@ -304,10 +309,16 @@ def prepare_model_and_optimizer(args, device): optimizer_grouped_parameters.append({'params': [p], 'weight_decay': 0.00, 'name': n}) names.append({'params': [n], 'weight_decay': 0.00}) - optimizer = BertLAMB(optimizer_grouped_parameters, - lr=args.learning_rate, - warmup=args.warmup_proportion, - t_total=args.max_steps) + #optimizer = BertLAMB(optimizer_grouped_parameters, + # lr=args.learning_rate, + # warmup=args.warmup_proportion, + # t_total=args.max_steps) + optimizer = FusedAdam(optimizer_grouped_parameters, + lr=args.learning_rate, + betas=(0.9, 0.999), + bias_correction=True, + eps=1e-6) + if args.fp16: if args.loss_scale == 0: @@ -347,6 +358,14 @@ def prepare_model_and_optimizer(args, device): flat_dist_call([param.data for param in model.parameters()], torch.distributed.broadcast, (0,) ) elif args.n_gpu > 1: model = torch.nn.DataParallel(model) + else: + from optimization import DistributedOptimizer, DistributedAdasumOptimizer + compression = hvd.Compression.none #if args.fp16_allreduce else hvd.Compression.none + optimizer = DistributedAdasumOptimizer(optimizer, + backward_passes_per_step=args.gradient_accumulation_steps, + compression=compression) + hvd.broadcast_parameters(model.state_dict(), root_rank=0) + #hvd.broadcast_optimizer_state(optimizer, root_rank=0) return model, optimizer, checkpoint, global_step @@ -367,9 +386,10 @@ def take_optimizer_step(args, optimizer, model, overflow_buf, global_step): amp_C.multi_tensor_scale(65536, overflow_buf, [master_grads, allreduced_views], - scaler.loss_scale() / (torch.distributed.get_world_size() * args.gradient_accumulation_steps)) + scaler.loss_scale() / (hvd.size() * args.gradient_accumulation_steps)) # 3. sum gradient across ranks. Because of the predivision, this averages the gradient - torch.distributed.all_reduce(flat_raw) + #torch.distributed.all_reduce(flat_raw) + hvd.allreduce_(flat_raw, op=hvd.Sum) # 4. combine unscaling and unflattening of allreduced gradient overflow_buf.zero_() amp_C.multi_tensor_scale(65536, @@ -391,7 +411,7 @@ def take_optimizer_step(args, optimizer, model, overflow_buf, global_step): if is_main_process(): print(("Rank {} :: Gradient overflow. Skipping step, " + "reducing loss scale to {}").format( - torch.distributed.get_rank(), + hvd.rank(), scaler.loss_scale())) if _amp_state.opt_properties.master_weights: for param in optimizer._amp_stash.all_fp32_from_fp16_params: @@ -399,6 +419,14 @@ def take_optimizer_step(args, optimizer, model, overflow_buf, global_step): for param in model.parameters(): param.grad = None else: + # toddm: only for adam + from optimization import warmup_linear as warmup + tmp = warmup(global_step / args.max_steps, args.warmup_proportion) + for group in optimizer.optimizer.param_groups: + group['lr'] = args.learning_rate * tmp + torch.nn.utils.clip_grad_norm_([p for group in optimizer.optimizer.param_groups for p in group['params'] if p is not None], 1.0) + # toddm end + optimizer.step() #optimizer.zero_grad() for param in model.parameters(): @@ -457,11 +485,11 @@ def main(): shared_file_list = {} - if torch.distributed.is_initialized() and torch.distributed.get_world_size() > num_files: - remainder = torch.distributed.get_world_size() % num_files - data_file = files[(f_start_id*torch.distributed.get_world_size()+torch.distributed.get_rank() + remainder*f_start_id)%num_files] + if hvd.size() > num_files: + remainder = hvd.size() % num_files + data_file = files[(f_start_id*hvd.size()+hvd.rank() + remainder*f_start_id)%num_files] else: - data_file = files[(f_start_id*torch.distributed.get_world_size()+torch.distributed.get_rank())%num_files] + data_file = files[(f_start_id*hvd.size()+hvd.rank())%num_files] previous_file = data_file @@ -479,10 +507,10 @@ def main(): for f_id in range(f_start_id + 1 , len(files)): - if torch.distributed.get_world_size() > num_files: - data_file = files[(f_id*torch.distributed.get_world_size()+torch.distributed.get_rank() + remainder*f_id)%num_files] + if hvd.size() > num_files: + data_file = files[(f_id*hvd.size()+hvd.rank() + remainder*f_id)%num_files] else: - data_file = files[(f_id*torch.distributed.get_world_size()+torch.distributed.get_rank())%num_files] + data_file = files[(f_id*hvd.size()+hvd.rank())%num_files] logger.info("file no %s file %s" % (f_id, previous_file)) @@ -523,9 +551,8 @@ def main(): last_num_steps = args.log_freq if last_num_steps == 0 else last_num_steps average_loss = torch.tensor(average_loss, dtype=torch.float32).cuda() average_loss = average_loss / (last_num_steps * divisor) - if (torch.distributed.is_initialized()): - average_loss /= torch.distributed.get_world_size() - torch.distributed.all_reduce(average_loss) + average_loss /= hvd.size() + hvd.allreduce_(average_loss) if is_main_process(): logger.info("Total Steps:{} Final Loss = {}".format(training_steps / args.gradient_accumulation_steps, average_loss.item())) elif training_steps % (args.log_freq * args.gradient_accumulation_steps) == 0: @@ -533,8 +560,8 @@ def main(): print("Step:{} Average Loss = {} Step Loss = {} LR {}".format(global_step, average_loss / ( args.log_freq * divisor), loss.item() * args.gradient_accumulation_steps / divisor, - optimizer.param_groups[0][ - 'lr'])) + optimizer.optimizer.param_groups[0][ + 'lr']), flush=True) average_loss = 0 if global_step >= args.max_steps or training_steps % ( diff --git a/PyTorch/LanguageModeling/BERT/schedulers.py b/PyTorch/LanguageModeling/BERT/schedulers.py index 2cf38841a..670352d39 100755 --- a/PyTorch/LanguageModeling/BERT/schedulers.py +++ b/PyTorch/LanguageModeling/BERT/schedulers.py @@ -15,7 +15,7 @@ import math import torch from torch.optim.optimizer import Optimizer -from apex.optimizers import FP16_Optimizer +#from apex.optimizers import FP16_Optimizer from torch.optim.lr_scheduler import _LRScheduler @@ -24,6 +24,7 @@ def __init__(self, optimizer, last_epoch=-1): # Check if using mixed precision training self.mixed_training = False base_optimizer = optimizer + print("type:" , type(optimizer)) if isinstance(optimizer, FP16_Optimizer): self.mixed_training = True self.fp16_optimizer = optimizer diff --git a/PyTorch/LanguageModeling/BERT/utils.py b/PyTorch/LanguageModeling/BERT/utils.py index 4f8e0d865..a1e0234b1 100755 --- a/PyTorch/LanguageModeling/BERT/utils.py +++ b/PyTorch/LanguageModeling/BERT/utils.py @@ -12,14 +12,11 @@ # limitations under the License. import torch -import torch.distributed as dist +#import torch.distributed as dist +import horovod.torch as hvd def get_rank(): - if not dist.is_available(): - return 0 - if not dist.is_initialized(): - return 0 - return dist.get_rank() + return hvd.rank() def is_main_process(): - return get_rank() == 0 \ No newline at end of file + return get_rank() == 0