From 25c12258e4366eb7e825c6458ad4ab396dc4857e Mon Sep 17 00:00:00 2001 From: William Falcon Date: Tue, 3 Mar 2020 12:17:49 -0500 Subject: [PATCH] added community examples (#1030) --- docs/source/multi_gpu.rst | 151 ++++++++++++++++++++++++++++++++++++-- 1 file changed, 144 insertions(+), 7 deletions(-) diff --git a/docs/source/multi_gpu.rst b/docs/source/multi_gpu.rst index f44cb3f9..2b1776f3 100644 --- a/docs/source/multi_gpu.rst +++ b/docs/source/multi_gpu.rst @@ -1,10 +1,83 @@ Multi-GPU training -===================== - +=================== Lightning supports multiple ways of doing distributed training. -Data Parallel (dp) +Preparing your code ------------------- +To train on CPU/GPU/TPU without changing your code, we need to build a few good habits :) + +Delete .cuda() or .to() calls +^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ + +Delete any calls to .cuda() or .to(device). + +.. code-block:: python + + # before lightning + def forward(self, x): + x = x.cuda(0) + layer_1.cuda(0) + x_hat = layer_1(x) + + # after lightning + def forward(self, x): + x_hat = layer_1(x) + +Init using type_as +^^^^^^^^^^^^^^^^^^ +When you need to create a new tensor, use `type_as`. +This will make your code scale to any arbitrary number of GPUs or TPUs with Lightning + +.. code-block:: python + + # before lightning + def forward(self, x): + z = torch.Tensor(2, 3) + z = z.cuda(0) + + # with lightning + def forward(self, x): + z = torch.Tensor(2, 3) + z = z.type_as(x.type()) + +Remove samplers +^^^^^^^^^^^^^^^ +For multi-node or TPU training, in PyTorch we must use `torch.nn.DistributedSampler`. The +sampler makes sure each GPU sees the appropriate part of your data. + +.. code-block:: python + + # without lightning + def train_dataloader(self): + dataset = MNIST(...) + sampler = None + + if self.on_tpu: + sampler = DistributedSampler(dataset) + + return DataLoader(dataset, sampler=sampler) + +With Lightning, you don't need to do this because it takes care of adding the correct samplers +when needed. + +.. code-block:: python + + # with lightning + def train_dataloader(self): + dataset = MNIST(...) + return DataLoader(dataset) + +Distributed modes +----------------- +Lightning allows multiple ways of training + +- Data Parallel (`distributed_backend='dp'`) (multiple-gpus, 1 machine) +- DistributedDataParallel (`distributed_backend='ddp'`) (multiple-gpus across many machines). +- DistributedDataParallel2 (`distributed_backend='ddp2'`) (dp in a machine, ddp across machines). +- TPUs (`num_tpu_cores=8|x`) (tpu or TPU pod) + +Data Parallel (dp) +^^^^^^^^^^^^^^^^^^ `DataParallel `_ splits a batch across k GPUs. That is, if you have a batch of 32 and use dp with 2 gpus, each GPU will process 16 samples, after which the root node will aggregate the results. @@ -14,7 +87,7 @@ each GPU will process 16 samples, after which the root node will aggregate the r trainer = pl.Trainer(gpus=2, distributed_backend='dp') Distributed Data Parallel ---------------------------- +^^^^^^^^^^^^^^^^^^^^^^^^^ `DistributedDataParallel `_ works as follows. 1. Each GPU across every node gets its own process. @@ -40,7 +113,7 @@ Distributed Data Parallel trainer = pl.Trainer(gpus=8, distributed_backend='ddp', num_nodes=4) Distributed Data Parallel 2 ------------------------------ +^^^^^^^^^^^^^^^^^^^^^^^^^^^ In certain cases, it's advantageous to use all batches on the same machine instead of a subset. For instance you might want to compute a NCE loss where it pays to have more negative samples. @@ -61,9 +134,73 @@ In this case, we can use ddp2 which behaves like dp in a machine and ddp across # train on 32 GPUs (4 nodes) trainer = pl.Trainer(gpus=8, distributed_backend='ddp2', num_nodes=4) +DP/DDP2 caveats +^^^^^^^^^^^^^^^ +In DP and DDP2 each GPU within a machine sees a portion of a batch. +DP and ddp2 roughly do the following: + +.. code-block:: python + + def distributed_forward(batch, model): + batch = torch.Tensor(32, 8) + gpu_0_batch = batch[:8] + gpu_1_batch = batch[8:16] + gpu_2_batch = batch[16:24] + gpu_3_batch = batch[24:] + + y_0 = model_copy_gpu_0(gpu_0_batch) + y_1 = model_copy_gpu_0(gpu_1_batch) + y_2 = model_copy_gpu_0(gpu_2_batch) + y_3 = model_copy_gpu_0(gpu_3_batch) + + return [y_0, y_1, y_2, y_3] + +So, when Lightning calls any of the `training_step`, `validation_step`, `test_step` +you will only be operating on one of those pieces. + +.. code-block:: python + + # the batch here is a portion of the FULL batch + def training_step(self, batch, batch_idx): + y_0 = batch + +For most metrics, this doesn't really matter. However, if you want +full batch statistics or want to use the outputs of the training_step +to do something like a softmax, you can use the `training_end` step. + +.. code-block:: python + + def training_end(self, outputs): + outputs = torch.cat(outputs, dim=1) + softmax = softmax(outputs, dim=1) + out = softmax.mean() + return out + +In pseudocode, the full sequence is: + +.. code-block:: python + + # get data + batch = next(dataloader) + + # copy model and data to each gpu + batch_splits = split_batch(batch, num_gpus) + models = copy_model_to_gpus(model) + + # in parallel, operate on each batch chunk + all_results = [] + for gpu_num in gpus: + batch_split = batch_splits[gpu_num] + gpu_model = models[gpu_num] + out = gpu_model(batch_split) + all_results.append(out) + + # calculate statistics for all parts of the batch + full out = model.training_end(all_results) + Implement Your Own Distributed (DDP) training ----------------------------------------------- -If you need your own way to init PyTorch DDP you can override :meth:`pytorch_lightning.core.LightningModule.init_ddp_connection`. +^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ +If you need your own way to init PyTorch DDP you can override :meth:`pytorch_lightning.core.LightningModule.`. If you also need to use your own DDP implementation, override: :meth:`pytorch_lightning.core.LightningModule.configure_ddp`.