In [1]:
epochs = 20
n_train_items = 1280

In [2]:
import torch
import torch.nn as nn
import torch.nn.functional as F
import torch.optim as optim
from torchvision import datasets, transforms

In [3]:
import syft as sy  # <-- NEW: import the Pysyft library
hook = sy.TorchHook(torch)  # <-- NEW: hook PyTorch ie add extra functionalities to support Federated Learning
# simulation functions
def connect_to_workers(n_workers):
    return [
        sy.VirtualWorker(hook, id=f"worker{i+1}")
        for i in range(n_workers)
    ]


workers = connect_to_workers(n_workers=20)


  _np_qint8 = np.dtype([("qint8", np.int8, 1)])
  _np_quint8 = np.dtype([("quint8", np.uint8, 1)])
  _np_qint16 = np.dtype([("qint16", np.int16, 1)])
  _np_quint16 = np.dtype([("quint16", np.uint16, 1)])
  _np_qint32 = np.dtype([("qint32", np.int32, 1)])
  np_resource = np.dtype([("resource", np.ubyte, 1)])


In [31]:
class Arguments():
    def __init__(self):
        self.batch_size = 64
        self.test_batch_size = 64
        self.epochs = 20
        self.lr = 0.01
        self.momentum = 0.5
        self.no_cuda = False
        self.seed = 1
        self.log_interval = 3
        self.save_model = False

args = Arguments()

use_cuda = not args.no_cuda and torch.cuda.is_available()

torch.manual_seed(args.seed)

device = torch.device("cuda" if use_cuda else "cpu")

kwargs = {'num_workers': 1, 'pin_memory': True} if use_cuda else {}

In [5]:
# federated_train_loader = sy.FederatedDataLoader( # <-- this is now a FederatedDataLoader 
#     datasets.MNIST('../data', train=True, download=True,
#                    transform=transforms.Compose([
#                        transforms.ToTensor(),
#                        transforms.Normalize((0.1307,), (0.3081,))
#                    ]))
#     .federate(workers), # <-- NEW: we distribute the dataset across all the workers, it's now a FederatedDataset
#     batch_size=args.batch_size, shuffle=True, **kwargs)

# test_loader = torch.utils.data.DataLoader(
#     datasets.MNIST('../data', train=False, transform=transforms.Compose([
#                        transforms.ToTensor(),
#                        transforms.Normalize((0.1307,), (0.3081,))
#                    ])),
#     batch_size=args.test_batch_size, shuffle=True, **kwargs)

train_loader = torch.utils.data.DataLoader(
    datasets.MNIST('../data', train=True, download=True, transform=transforms.Compose([
                       transforms.ToTensor(),
                       transforms.Normalize((0.1307,), (0.3081,))
                   ])),
    batch_size=args.batch_size
)
test_loader = torch.utils.data.DataLoader(
    datasets.MNIST('../data', train=False, download=True, transform=transforms.Compose([
                       transforms.ToTensor(),
                       transforms.Normalize((0.1307,), (0.3081,))
                   ])),
    batch_size=args.test_batch_size
)

    
#---

less_train_dataloader = [
        ((data), (target))
        for i, (data, target) in enumerate(train_loader)
        if i < n_train_items / args.batch_size
    ]
less_test_dataloader = [
        ((data), (target))
        for i, (data, target) in enumerate(test_loader)
        if i < n_train_items / args.batch_size
    ]




In [6]:
# from PIL import Image
# import numpy 
# #mnist_dataset.__getitem__(2)[1]
# a = (mnist_dataset.__getitem__(0)[0]).numpy()
# a.dtype = 'uint8'
# print(a)
# Image.fromarray(a[0], mode= 'P')

In [7]:
class Net(nn.Module):
    def __init__(self):
        super(Net, self).__init__()
        self.conv1 = nn.Conv2d(1, 20, 5, 1)
        self.conv2 = nn.Conv2d(20, 50, 5, 1)
        self.fc1 = nn.Linear(4*4*50, 500)
        self.fc2 = nn.Linear(500, 10)

    def forward(self, x):
        x = F.relu(self.conv1(x))
        x = F.max_pool2d(x, 2, 2)
        x = F.relu(self.conv2(x))
        x = F.max_pool2d(x, 2, 2)
        x = x.view(-1, 4*4*50)
        x = F.relu(self.fc1(x))
        x = self.fc2(x)
        return F.log_softmax(x, dim=1)

In [27]:
def train(args, model, device, less_train_dataloader, optimizer, epoch, workers):
    model.train()
    for batch_idx, (data, target) in enumerate(less_train_dataloader): # <-- now it is a distributed dataset
        model.send(workers[batch_idx%len(workers)]) # <-- NEW: send the model to the right location
        
        data_on_worker = data.send(workers[batch_idx%len(workers)])
        target_on_worker = target.send(workers[batch_idx%len(workers)])
        
        data_on_worker, target_on_worker = data_on_worker.to(device), target_on_worker.to(device)
        
        optimizer.zero_grad()
        output = model(data_on_worker)
        loss = F.nll_loss(output, target_on_worker)
        loss.backward()
        optimizer.step()
        model.get() # <-- NEW: get the model back
        if batch_idx % args.log_interval == 0:
            loss = loss.get() # <-- NEW: get the loss back
            print('Train Epoch: {} [{}/{} ({:.0f}%)]\tLoss: {:.6f}'.format(
                epoch, batch_idx * args.batch_size, len(less_train_dataloader) * args.batch_size,
                100. * batch_idx / len(less_train_dataloader), loss.item()))

In [28]:
args.batch_size

64

In [29]:
def test(args, model, device, test_loader):
    model.eval()
    test_loss = 0
    correct = 0
    with torch.no_grad():
        for data, target in test_loader:
            data, target = data.to(device), target.to(device)
            
            output = model(data)
            test_loss += F.nll_loss(output, target, reduction='sum').item() # sum up batch loss
            pred = output.argmax(1, keepdim=True) # get the index of the max log-probability 
            correct += pred.eq(target.view_as(pred)).sum().item()

    test_loss /= len(test_loader*args.batch_size)

    print('\nTest set: Average loss: {:.4f}, Accuracy: {}/{} ({:.0f}%)\n'.format(
        test_loss, correct, len(test_loader* args.batch_size),
        100. * correct / (len(test_loader)*args.batch_size)))

In [32]:
%%time
model = Net().to(device)
optimizer = optim.SGD(model.parameters(), lr=args.lr) # TODO momentum is not supported at the moment

for epoch in range(1, args.epochs + 1):
    train(args, model, device, less_train_dataloader, optimizer, epoch, workers)
    test(args, model, device, less_test_dataloader)

if (args.save_model):
    torch.save(model.state_dict(), "mnist_cnn.pt")


Test set: Average loss: 2.2409, Accuracy: 372/1280 (29%)


Test set: Average loss: 2.1541, Accuracy: 481/1280 (38%)


Test set: Average loss: 2.0013, Accuracy: 615/1280 (48%)


Test set: Average loss: 1.7253, Accuracy: 707/1280 (55%)


Test set: Average loss: 1.3563, Accuracy: 806/1280 (63%)


Test set: Average loss: 1.0623, Accuracy: 901/1280 (70%)


Test set: Average loss: 0.8874, Accuracy: 940/1280 (73%)


Test set: Average loss: 0.7978, Accuracy: 947/1280 (74%)


Test set: Average loss: 0.7504, Accuracy: 962/1280 (75%)


Test set: Average loss: 0.7129, Accuracy: 972/1280 (76%)


Test set: Average loss: 0.6783, Accuracy: 979/1280 (76%)


Test set: Average loss: 0.6435, Accuracy: 984/1280 (77%)


Test set: Average loss: 0.6124, Accuracy: 999/1280 (78%)


Test set: Average loss: 0.5811, Accuracy: 1005/1280 (79%)


Test set: Average loss: 0.5544, Accuracy: 1021/1280 (80%)


Test set: Average loss: 0.5301, Accuracy: 1042/1280 (81%)


Test set: Average loss: 0.5099, Accuracy: 1050/1280 

In [14]:
print(workers[0])

<VirtualWorker id:worker1 #objects:6>


In [176]:
model = Net()
for batch_idx, (data, target) in enumerate(less_train_dataloader):
    data = data.send(workers[0])
    print(data)
#     if batch_idx<3:
#         model.send(workers[0])
#         print(data.size())
        
#         pre = model(data)


(Wrapper)>[PointerTensor | me:60407973905 -> worker1:78163135291]
(Wrapper)>[PointerTensor | me:74050153873 -> worker1:11814332395]
(Wrapper)>[PointerTensor | me:93651540692 -> worker1:98042234477]
(Wrapper)>[PointerTensor | me:1792853962 -> worker1:10207084528]
(Wrapper)>[PointerTensor | me:84209531306 -> worker1:30075516339]
(Wrapper)>[PointerTensor | me:4457758866 -> worker1:81018670476]
(Wrapper)>[PointerTensor | me:37733711282 -> worker1:94469615650]
(Wrapper)>[PointerTensor | me:41956409820 -> worker1:31235277075]
(Wrapper)>[PointerTensor | me:894364867 -> worker1:14226006381]
(Wrapper)>[PointerTensor | me:38101863086 -> worker1:87171777304]
(Wrapper)>[PointerTensor | me:95239140209 -> worker1:51638303354]
(Wrapper)>[PointerTensor | me:35612061462 -> worker1:68752432269]
(Wrapper)>[PointerTensor | me:64152108992 -> worker1:31182305053]
(Wrapper)>[PointerTensor | me:51875992271 -> worker1:65242662746]
(Wrapper)>[PointerTensor | me:86903732979 -> worker1:47402955970]
(Wrapper)>[Poi

In [183]:
workers[0]._objects[50490937571]

tensor([[[[-0.4242, -0.4242, -0.4242,  ..., -0.4242, -0.4242, -0.4242],
          [-0.4242, -0.4242, -0.4242,  ..., -0.4242, -0.4242, -0.4242],
          [-0.4242, -0.4242, -0.4242,  ..., -0.4242, -0.4242, -0.4242],
          ...,
          [-0.4242, -0.4242, -0.4242,  ..., -0.4242, -0.4242, -0.4242],
          [-0.4242, -0.4242, -0.4242,  ..., -0.4242, -0.4242, -0.4242],
          [-0.4242, -0.4242, -0.4242,  ..., -0.4242, -0.4242, -0.4242]]],


        [[[-0.4242, -0.4242, -0.4242,  ..., -0.4242, -0.4242, -0.4242],
          [-0.4242, -0.4242, -0.4242,  ..., -0.4242, -0.4242, -0.4242],
          [-0.4242, -0.4242, -0.4242,  ..., -0.4242, -0.4242, -0.4242],
          ...,
          [-0.4242, -0.4242, -0.4242,  ..., -0.4242, -0.4242, -0.4242],
          [-0.4242, -0.4242, -0.4242,  ..., -0.4242, -0.4242, -0.4242],
          [-0.4242, -0.4242, -0.4242,  ..., -0.4242, -0.4242, -0.4242]]],


        [[[-0.4242, -0.4242, -0.4242,  ..., -0.4242, -0.4242, -0.4242],
          [-0.4242, -0.424