diff --git a/build/cuda.cu.o b/build/cuda.cu.o index 00a1896..1d8c326 100644 Binary files a/build/cuda.cu.o and b/build/cuda.cu.o differ diff --git a/build/distributed.o b/build/distributed.o index 11a92ef..58e0ef2 100644 Binary files a/build/distributed.o and b/build/distributed.o differ diff --git a/build/tensor.o b/build/tensor.o index 17fca73..d48dad3 100644 Binary files a/build/tensor.o and b/build/tensor.o differ diff --git a/norch/__pycache__/tensor.cpython-38.pyc b/norch/__pycache__/tensor.cpython-38.pyc index 2381a0f..8505635 100644 Binary files a/norch/__pycache__/tensor.cpython-38.pyc and b/norch/__pycache__/tensor.cpython-38.pyc differ diff --git a/norch/csrc/tensor.cpp b/norch/csrc/tensor.cpp index cc6b922..74d8178 100644 --- a/norch/csrc/tensor.cpp +++ b/norch/csrc/tensor.cpp @@ -81,11 +81,6 @@ extern "C" { else if ((strcmp(target_device, "cpu") == 0) && (strcmp(tensor->device, "cuda") == 0)) { cuda_to_cpu(tensor); } - - else { - fprintf(stderr, "Could not send tensor to device %d", device_id); - exit(1); - } } Tensor* add_tensor(Tensor* tensor1, Tensor* tensor2) { diff --git a/norch/libtensor.so b/norch/libtensor.so index f50c17c..bc35ee2 100755 Binary files a/norch/libtensor.so and b/norch/libtensor.so differ diff --git a/norch/nn/__pycache__/parameter.cpython-38.pyc b/norch/nn/__pycache__/parameter.cpython-38.pyc index fc69338..14bb069 100644 Binary files a/norch/nn/__pycache__/parameter.cpython-38.pyc and b/norch/nn/__pycache__/parameter.cpython-38.pyc differ diff --git a/norch/nn/parallel.py b/norch/nn/parallel.py index 16eca59..c021537 100644 --- a/norch/nn/parallel.py +++ b/norch/nn/parallel.py @@ -1,6 +1,7 @@ from .module import * import norch.distributed as dist import os +import norch class DistributedDataParallel(Module): def __init__(self, module): @@ -8,16 +9,40 @@ class DistributedDataParallel(Module): self.module = module + + self.broadcast_parameters() + self.register_grads_hooks() + def forward(self, *inputs, **kwargs): return self.module(*inputs, **kwargs) - def backward(self): - self.module.backward() - for module, name, _ in self.parameters(): - parameter = getattr(module, name) - dist.allreduce_mean_tensor(parameter) - setattr(module, name, parameter) + def broadcast_parameters(self): + """ + Broadcast parameters of device 0 to all devices + """ + for _, _, parameter in self.parameters(): + dist.broadcast_tensor(parameter) + + def allreduce_grads_hook(grad): + """ + Everytime a gradient is assign to some value, it calculates mean of this gradient among all devices + """ + if isinstance(grad, norch.Tensor): + dist.allreduce_sum_tensor(grad) + grad /= dist.get_world_size() + return grad + + def register_grads_hooks(self): + """ + Everytime a gradient is assign it calls this allreduce hook + """ + for _, _, parameter in self.parameters(): + parameter.register_hook(self.allreduce_grads_hook) + + + + diff --git a/norch/optim/optimizers/__pycache__/sgd.cpython-38.pyc b/norch/optim/optimizers/__pycache__/sgd.cpython-38.pyc index c1b0789..7336b8e 100644 Binary files a/norch/optim/optimizers/__pycache__/sgd.cpython-38.pyc and b/norch/optim/optimizers/__pycache__/sgd.cpython-38.pyc differ diff --git a/norch/tensor.py b/norch/tensor.py index 2e80e72..0a1ecab 100644 --- a/norch/tensor.py +++ b/norch/tensor.py @@ -39,6 +39,7 @@ class Tensor: self.numel *= s self.requires_grad = requires_grad + self.hooks = [] self.grad = None self.grad_fn = None @@ -59,6 +60,7 @@ class Tensor: self.ndim = None, self.device = device self.requires_grad = None + self.hooks = [] self.grad = None self.grad_fn = None @@ -79,6 +81,15 @@ class Tensor: flat_data, shape = flatten_recursively(nested_list) return flat_data, shape + def __setattr__(self, name, value): + if name == 'grad': + for hook in self.hooks: + value = hook(value) + super().__setattr__(name, value) + + def register_hook(self, function): + self.hooks.append(function) + def ones_like(self): Tensor._C.ones_like_tensor.argtypes = [ctypes.POINTER(CTensor)] diff --git a/train.py b/train.py index 96cbc8d..4eb30c6 100644 --- a/train.py +++ b/train.py @@ -3,52 +3,21 @@ import norch import norch.distributed as dist import norch.distributed +import norch.nn as nn +import norch.optim as optim +from norch.utils.data.dataloader import DataLoader +from norch.nn.parallel import DistributedDataParallel +from norch.utils.data.distributed import DistributedSampler +from norch.norchvision import transforms as T +import numpy as np +import matplotlib.pyplot as plt +import random +random.seed(1) + def main(): - local_rank = int(os.getenv('OMPI_COMM_WORLD_LOCAL_RANK', -1)) - rank = int(os.getenv('OMPI_COMM_WORLD_RANK', -1)) - world_size = int(os.getenv('OMPI_COMM_WORLD_SIZE', -1)) - - dist.init_process_group(rank, world_size) - - tensor = norch.Tensor([1,1,1]).to(rank) - tensor = (rank + 1) * tensor - print(f"BEFORE on rank {rank}: {tensor} \n\n") - - dist.allreduce_sum_tensor(tensor) - - print(f"AFTER ALLREDUCE on rank {rank}: {tensor} \n\n") - - print("###############\n\n\n") - - tensor = tensor * 10 - print(f"BEFORE BROADCAST on rank {rank}: {tensor} \n\n") - - dist.broadcast_tensor(tensor) - - print(f"AFTER BROADCAST on rank {rank}: {tensor} \n\n") - -def main2(): - import norch - import norch.nn as nn - import norch.optim as optim - from norch.utils.data.dataloader import DataLoader - from norch.nn.parallel import DistributedDataParallel - from norch.utils.data.distributed import DistributedSampler - from norch.norchvision import transforms as T - import numpy as np - import matplotlib.pyplot as plt - import random - random.seed(1) - - local_rank = int(os.getenv('OMPI_COMM_WORLD_LOCAL_RANK', -1)) - rank = int(os.getenv('OMPI_COMM_WORLD_RANK', -1)) - world_size = int(os.getenv('OMPI_COMM_WORLD_SIZE', -1)) - - dist.init_process_group(rank, world_size) - BATCH_SIZE = 32 - device = local_rank + device = "cpu" epochs = 10 transform = T.Compose( @@ -65,8 +34,7 @@ def main2(): ) train_data, test_data = norch.norchvision.datasets.MNIST.splits(transform=transform, target_transform=target_transform) - distributed_sampler = DistributedSampler(dataset=train_data, num_replicas=world_size, rank=local_rank) - train_loader = norch.utils.data.DataLoader(train_data, batch_size=BATCH_SIZE, sampler=distributed_sampler) + train_loader = norch.utils.data.DataLoader(train_data, batch_size=BATCH_SIZE) class MyModel(nn.Module): def __init__(self): @@ -90,9 +58,6 @@ def main2(): optimizer = optim.SGD(model.parameters(), lr=0.01) loss_list = [] - print(f"Local rank: {local_rank}") - print(f"World size: {world_size}") - for epoch in range(epochs): for idx, batch in enumerate(train_loader): @@ -102,17 +67,19 @@ def main2(): target = target.to(device) outputs = model(inputs) - loss = criterion(outputs, target) optimizer.zero_grad() + print("#####################\n\nantes backward") + print(model.module.fc1.bias.grad) loss.backward() - print(f"AFTER rank {local_rank}: {model.module.fc2.bias.grad}") - print("\n\n") + print(model.module.fc1.bias.grad) + print("\n\n\n###############\n\npós backward") optimizer.step() + print("@@") break break diff --git a/train_multigpu.py b/train_multigpu.py new file mode 100644 index 0000000..20e3b6e --- /dev/null +++ b/train_multigpu.py @@ -0,0 +1,120 @@ +import os +import norch +import norch.distributed as dist +import norch.distributed + +def main(): + + local_rank = int(os.getenv('OMPI_COMM_WORLD_LOCAL_RANK', -1)) + rank = int(os.getenv('OMPI_COMM_WORLD_RANK', -1)) + world_size = int(os.getenv('OMPI_COMM_WORLD_SIZE', -1)) + + dist.init_process_group(rank, world_size) + + tensor = norch.Tensor([1,1,1]).to(rank) + tensor = (rank + 1) * tensor + print(f"BEFORE on rank {rank}: {tensor} \n\n") + + dist.allreduce_sum_tensor(tensor) + + print(f"AFTER ALLREDUCE on rank {rank}: {tensor} \n\n") + + print("###############\n\n\n") + + tensor = tensor * 10 + print(f"BEFORE BROADCAST on rank {rank}: {tensor} \n\n") + + dist.broadcast_tensor(tensor) + + print(f"AFTER BROADCAST on rank {rank}: {tensor} \n\n") + +def main2(): + import norch + import norch.nn as nn + import norch.optim as optim + from norch.utils.data.dataloader import DataLoader + from norch.nn.parallel import DistributedDataParallel + from norch.utils.data.distributed import DistributedSampler + from norch.norchvision import transforms as T + import numpy as np + import matplotlib.pyplot as plt + import random + random.seed(1) + + local_rank = int(os.getenv('OMPI_COMM_WORLD_LOCAL_RANK', -1)) + rank = int(os.getenv('OMPI_COMM_WORLD_RANK', -1)) + world_size = int(os.getenv('OMPI_COMM_WORLD_SIZE', -1)) + + dist.init_process_group(rank, world_size) + + BATCH_SIZE = 32 + device = local_rank + epochs = 10 + + transform = T.Compose( + [ + T.ToTensor(), + T.Reshape([-1, 784, 1]) + ] + ) + + target_transform = T.Compose( + [ + T.ToTensor() + ] + ) + + train_data, test_data = norch.norchvision.datasets.MNIST.splits(transform=transform, target_transform=target_transform) + distributed_sampler = DistributedSampler(dataset=train_data, num_replicas=world_size, rank=local_rank) + train_loader = norch.utils.data.DataLoader(train_data, batch_size=BATCH_SIZE, sampler=distributed_sampler) + + class MyModel(nn.Module): + def __init__(self): + super(MyModel, self).__init__() + self.fc1 = nn.Linear(784, 30) + self.sigmoid1 = nn.Sigmoid() + self.fc2 = nn.Linear(30, 10) + self.sigmoid2 = nn.Sigmoid() + + def forward(self, x): + out = self.fc1(x) + out = self.sigmoid1(out) + out = self.fc2(out) + out = self.sigmoid2(out) + + return out + + model = MyModel().to(device) + model = DistributedDataParallel(model) + print(f"parameter bias on Rank {rank}: {model.module.fc1.bias}\n\n") + criterion = nn.CrossEntropyLoss() + optimizer = optim.SGD(model.parameters(), lr=0.01) + loss_list = [] + + for epoch in range(epochs): + for idx, batch in enumerate(train_loader): + + inputs, target = batch + + inputs = inputs.to(device) + target = target.to(device) + + outputs = model(inputs) + + loss = criterion(outputs, target) + + optimizer.zero_grad() + + loss.backward() + print(f"GRADIENT AFTER rank {local_rank}: {model.module.fc2.bias.grad}") + print("\n\n") + + + optimizer.step() + break + + break + +if __name__ == "__main__": + main() +