# Import necessary libraries

In [1]:
# General
import os
import cv2
import numpy as np
from tqdm import tqdm
import matplotlib.pyplot as plt
import pickle
import time
import copy
import pandas as pd


# Pytorch
import torch
import torch.nn as nn
import torch.nn.functional as F
import torch.optim as optim
from torch.utils.data import TensorDataset, DataLoader
from torchvision import datasets, transforms


# PySyft
import syft as sy
from syft.frameworks.torch.fl import utils
from syft.workers.websocket_client import WebsocketClientWorker

# Pre-processing the Data

In [2]:
# Set the image size Y where Y represents YxY 
IMG_SIZE = 50
BATCH_SIZE = 100
LR = 0.001

In [3]:
train = datasets.MNIST(r"/media/wilfredo/Willie931GB/EURECOM_SLU_Linux/II_SEMESTER/SLU/PAPER_KDD2022/EXPERIMENTS/PySyft/Datasets/MNIST", 
                      train = True, download = True, 
                      transform = transforms.Compose([transforms.Resize(IMG_SIZE),
                                                      transforms.ToTensor()]))

test = datasets.MNIST(r"/media/wilfredo/Willie931GB/EURECOM_SLU_Linux/II_SEMESTER/SLU/PAPER_KDD2022/EXPERIMENTS/PySyft/Datasets/MNIST", 
                      train = False, download = True, 
                      transform = transforms.Compose([transforms.Resize(IMG_SIZE),
                                                      transforms.ToTensor()]))

In [4]:
# Load the data from the file it was saved in. Take the ENTIRE dataset!
training_data = torch.utils.data.DataLoader(train, batch_size = int(len(train)/2), shuffle = True)
test_data = torch.utils.data.DataLoader(test, batch_size = int(len(test)/2), shuffle = True)

# Create the CNN (based on VGG11)
Source: Page 3/14, Table 1, Configuration A, https://arxiv.org/pdf/1409.1556.pdf

## Individual Client Models

In [5]:
class Net_client1(nn.Module):
    def __init__(self):
        super().__init__()
        # Define your first convolutional layer: input = 1, output = 32 convolutional features, kernel size = 5
        # Remember that kernel = 5 means that the "window" used to scan for features will be 5x5
        self.conv1 = nn.Conv2d(1, 16, 5)
        self.conv2 = nn.Conv2d(16, 32, 5)

    # Function defining only one part of the forward pass (the convolution layers only). This will also write
    # the output dimensions of the conv layers to self._to_linear ONCE, and this information will then be used 
    # as the input data flattened dimensions of the next fully connected layers 
    def convs(self, x):
        # Convolutional layer 1 + activation + max_pooling
        x = self.conv1(x)
        x = F.relu(x)
        x = F.max_pool2d(x, (2, 2))
        x = self.conv2(x)
        x = F.relu(x)
        x = F.max_pool2d(x, (2, 2))
        return x
    
    # Function defining the rest of the forward pass
    def forward(self, x):
        # Run the convs layers first
        x = self.convs(x)
        return x

net_client1 = Net_client1()

In [6]:
class Net_client2(nn.Module):
    def __init__(self):
        super().__init__()
        
        # Start from the third convolutional layer
        self.conv3 = nn.Conv2d(32, 64, 5)

    # Function defining only one part of the forward pass (the convolution layers only). This will also write
    # the output dimensions of the conv layers to self._to_linear ONCE, and this information will then be used 
    # as the input data flattened dimensions of the next fully connected layers 
    def convs(self, x):
        # Convolutional layer 1 + activation + max_pooling
        x = self.conv3(x)
        x = F.relu(x)
        x = F.max_pool2d(x, (2, 2))
        return x
    
    # Function defining the rest of the forward pass
    def forward(self, x):
        # Run the convs layers first
        x = self.convs(x)
        return x

net_client2 = Net_client2()

In [7]:
class Net_client3(nn.Module):
    def __init__(self):
        super().__init__()
        
        # Run the fully connected layers. We know the input of this fc1 layer is 512, because of our previous
        # results with FL, where self.__to__linear told us this result when you run the cell that contains the 
        # NN
        self._to_linear = 256
        self.fc1 = nn.Linear(256, 32)
        self.fc2 = nn.Linear(32, 2)

    
    # Function defining the rest of the forward pass
    def forward(self, x):
        # Reshape the output data from the convs to be flattened
        x = x.view(-1, self._to_linear)
        # Pass the data through the fully connected layers now
        x = F.relu(self.fc1(x))
        # Pass it through the final layer
        x = self.fc2(x)
        # One final softmax function to make the output vector look nicer
        x = F.softmax(x, dim = 1)
        return x

net_client3 = Net_client3()

In [8]:
# Take a look at our models
model_client1 = net_client1
model_client2 = net_client2
model_client3 = net_client3

In [9]:
model_client1

Net_client1(
  (conv1): Conv2d(1, 16, kernel_size=(5, 5), stride=(1, 1))
  (conv2): Conv2d(16, 32, kernel_size=(5, 5), stride=(1, 1))
)

In [10]:
model_client2

Net_client2(
  (conv3): Conv2d(32, 64, kernel_size=(5, 5), stride=(1, 1))
)

In [11]:
model_client3

Net_client3(
  (fc1): Linear(in_features=256, out_features=32, bias=True)
  (fc2): Linear(in_features=32, out_features=2, bias=True)
)

# Establish your loss function

In [12]:
# Set your loss function (MSE for images!)
loss_function = nn.MSELoss()

# Separate your data into data, labels, training, testing, and scale it

In [13]:
# Take the data loaded onto training_data. You NEED to iterate over it to take it, even if you
# want to take the entire thing. Make sure to convert the values to floats
X = next(iter(training_data))[0]
y_unformatted = next(iter(training_data))[1].type(torch.FloatTensor)
X_test = next(iter(test_data))[0]
y_test_unformatted = next(iter(test_data))[1].type(torch.FloatTensor)

# The two other cases in this paper use 2 dimensional labels (0, 1), not only (0)
# MNIST by default comes with labels in the format (9) instead of (9, 0). To change this:
# Create tensors with all zeros of the same size
y_unformatted_addition = torch.zeros(y_unformatted.size())
y_test_unformatted_addition = torch.zeros(y_test_unformatted.size())
# Then stack them together (0 for vertically, -1 for horizontally)
y = torch.stack((y_unformatted, y_unformatted_addition), -1)
y_test = torch.stack((y_test_unformatted, y_test_unformatted_addition), -1)


In [14]:
# Define your training data
# train_X = X[:-val_size]
# train_y = y[:-val_size]
train_X = X
train_y = y

# Define your testing (validation) data
# test_X = X[-val_size:]
# test_y = y[-val_size:]
test_X = X_test
test_y = y_test

# Pipeline Learning

## Establish the virtual workers, their data, their NNs, and their optimizers

In [15]:
# Start the hook
hook = sy.TorchHook(torch)

# Create your virtual workers and our server
worker1 = sy.VirtualWorker(hook, id="worker1")
worker2 = sy.VirtualWorker(hook, id="worker2")
worker3 = sy.VirtualWorker(hook, id="worker3")

# Put the WORKERS into a list for easier access later on
compute_nodes = [worker1, worker2, worker3]

In [16]:
# Clear the workers of any objects, just in case you forgot some were still there from a previous run
worker1.clear_objects()
worker2.clear_objects()
worker3.clear_objects()

<VirtualWorker id:worker3 #objects:0>

In [17]:
# # Establish the NN model for each worker. This is model-centric FL, so it is the same model for all workers
worker1_model = model_client1.copy()
worker2_model = model_client2.copy()
worker3_model = model_client3.copy()
# worker1_model = model_client1
# worker2_model = model_client2
# worker3_model = model_client3

# Establish the optimizer for each worker
worker1_optimizer = optim.SGD(worker1_model.parameters(), lr=LR)
worker2_optimizer = optim.SGD(worker2_model.parameters(), lr=LR)
worker3_optimizer = optim.SGD(worker3_model.parameters(), lr=LR)

In [18]:
# Organize the WORKER models and optimizers into lists. The server stuff must not be mixed with these
models = [worker1_model, worker2_model, worker3_model]
optimizers = [worker1_optimizer, worker2_optimizer, worker3_optimizer]

worker_collection = [(worker1, worker1_model, worker1_optimizer), (worker2, worker2_model, worker2_optimizer), 
                    (worker3, worker3_model, worker3_optimizer)]

## Sequential Training

In [19]:
def train():
    
    # This is a completely sequential algorithm, so:
    batch_count = 0
    batch_times = []
    total_epoch_time = 0
    for i in tqdm(range(0, int(len(train_X)), BATCH_SIZE)):
        # Send the models to their respective workers
        worker1_model.send(worker1)
        worker2_model.send(worker2)
        worker3_model.send(worker3)
        
        # Send the data to first worker in the chain. The labels are only needed on the last worker in the
        # chain
        batch_X = train_X[i : i + BATCH_SIZE]
        batch_y = train_y[i : i + BATCH_SIZE]
        batch_X = batch_X.send(worker1)
        batch_y = batch_y.send(worker3)
        
        # Zero the sequence for all models on both workers and server!
        worker1_optimizer.zero_grad()
        worker2_optimizer.zero_grad()
        worker3_optimizer.zero_grad()
#         print("Zeroed the grads for all workers")

        # Start FP on worker1
        FP1_start_time = time.time()
        intermediate_1 = worker1_model(batch_X)
        FP1_end_time = time.time() - FP1_start_time 
#         print("Finished FP on ", worker1.id)

        # Send the intermediate result to worker2 AND SPLIT THE COMPUTATIONAL GRAPH WITH DETACH()!
        remote_intermediate_1 = intermediate_1.detach().move(worker2).requires_grad_()
#         data_for_edge1 = intermediate_1.detach().move(edge1).requires_grad_()
#         print("Sent FP status to ", worker2.id)

        # Start FP on worker2
        FP2_start_time = time.time()
        intermediate_2 = worker2_model(remote_intermediate_1)
        FP2_end_time = time.time() - FP2_start_time
#         print("Finished FP on ", worker2.id)
        
        # Send the intermediate result to worker3 AND SPLIT THE COMPUTATIONAL GRAPH WITH DETACH()!
        remote_intermediate_2 = intermediate_2.detach().move(worker3).requires_grad_()
#         print("Sent FP status to ", worker3.id)

        # Finish the FP on worker3
        FP3_start_time = time.time()
        pred = worker3_model(remote_intermediate_2)
        FP3_end_time = time.time() - FP3_start_time
        
        total_FP_time = FP1_end_time + FP2_end_time + FP3_end_time 
#         print("Finished FP on ", worker3.id)

        # Calculate loss on worker3
        BP1_start_time = time.time()
        loss = loss_function(pred, batch_y)
        # Do BP on worker3
        loss.backward()
        worker3_optimizer.step()
        BP1_end_time = time.time() - BP1_start_time
#         print("Finished the BP on ", worker3.id)

        # Move gradients to worker2
        intermediate_2.move(worker2)
        grad_intermediate_2 = remote_intermediate_2.grad.copy().move(worker2)
        
        # Do BP on worker2
        BP2_start_time = time.time()
        intermediate_2.backward(grad_intermediate_2)
        worker2_optimizer.step()
        BP2_end_time = time.time() - BP2_start_time
        
        # Send gradients to worker1
        intermediate_1.move(worker1)
        grad_intermediate_1 = remote_intermediate_1.grad.copy().move(worker1)
        # Do BP on worker1
        BP3_start_time = time.time()
        intermediate_1.backward(grad_intermediate_1)
        worker3_optimizer.step()
        BP3_end_time = time.time() - BP3_start_time

        total_BP_time = BP1_end_time + BP2_end_time + BP3_end_time
        
        # Total batch time
        total_batch_time = total_FP_time + total_BP_time
        # Now, for this to be equivalent to the time measured in the other architectures, our
        # constant has to be THE AMOUNT OF DATA. ONE batch in the other architectures means
        # ONE batch PER CLIENT. So 100 batches of data in the other archs are equivalent to:
        # len(compute_nodes) * batch_time in this arch. As such:
        equivalent_batch_time = total_batch_time * len(compute_nodes)
        batch_times.append(equivalent_batch_time)
        # We are multiplying it here by len(compute_nodes) so that each individual cell, when
        # later comparing the CSV's is "equal" in amounts of data. HOWEVER, note we will still
        # have 3x the amount of batch recordings in this code because it needs to do each batch
        # individually. We can't have 3x times the amount of instances AND each instance be 3x
        # times the amount of other archs!! You can either remove the "*len(compute_nodes)"
        # section above OR do the following:
        total_epoch_time += total_batch_time # NOT THE EQUIVALENT BATCH TIME!!
        #because now you will have 3x the amount of instances but each instance counting as
        # NOT the equivalent of 3x the other archs, but each instance counting as its own unit. 
        # So now, you will have 300 batches in PL with normal batch time recorded for each epoch,
        # which is EQUIVALENT to 100 recorded batches in FL with normal batch time recorded
        # TLDR: CAREFUL YOU MULTIPLY BY LEN(COMPUTE_NODES) TWICE!!!
#         print("Total batch time = ", round(total_batch_time, 4), " s")
        # Get back the models and delete the batches to free up memory space for next batch
        worker1_model.get()
        worker2_model.get()
        worker3_model.get()
        
        worker1.clear_objects()
        worker2.clear_objects()
        worker3.clear_objects()
        
#         batch_count += 1
#         if batch_count >= 25:
#             break
        
    # OUTSIDE THE FOR LOOP
    print("TOTAL TIME FOR THIS EPOCH = ", round(total_epoch_time, 4), " s")
#     Return the new model on the server
    return worker1_model, worker2_model, worker3_model, batch_times, total_epoch_time

## Function used for testing

In [20]:
def test(new_worker1_model, new_worker2_model, new_worker3_model):
    
    # Calculate the accuracy
    correct = 0
    total = 0

    # Do not update your gradients while testing
    with torch.no_grad():
        print("Initiated model testing:")
        for i in tqdm(range(len(test_X))):
            
            # Put the model into evaluation mode so it does not update its gradients during this test
            new_worker1_model.eval()
            new_worker2_model.eval()
            new_worker3_model.eval()
            
            # Obtain the real class for the sample
            real_class = torch.argmax(test_y[i])

            # Obtain our prediction for said sample (not arg_maxed yet)
#             output = new_model_server(test_X[i].view(-1, 1, IMG_SIZE, IMG_SIZE))[0]
            output = new_worker1_model(test_X[i].view(-1, 1, IMG_SIZE, IMG_SIZE))
            output = new_worker2_model(output)
            output = new_worker3_model(output)[0]
                                       
#             output = new_worker3_model(new_worker2_model((new_worker1_model(test_X[i].view(-1, 1, IMG_SIZE, IMG_SIZE)))))[0]

            # Obtain our arg_maxed prediction for said sample
            predicted_class = torch.argmax(output)

            # Update counters
            if predicted_class == real_class:
                correct += 1
            total += 1

    print("Accuracy of the new model = ", round(correct/total, 3), "\n \n")

In [21]:
def reset(new_worker1_model, new_worker2_model, new_worker3_model):
    # Clear the workers of any objects, just in case you forgot some were still there from a previous run
    worker1.clear_objects()
    worker2.clear_objects()
    worker3.clear_objects()
    
    # Establish the NN model for each worker. This is model-centric FL, so it is the same model for all workers
    global worker1_model
    worker1_model = new_worker1_model.copy()
    global worker2_model
    worker2_model = new_worker2_model.copy()
    global worker3_model
    worker3_model = new_worker3_model.copy()
#     worker1_model = new_worker1_model
#     worker2_model = new_worker2_model
#     worker3_model = new_worker3_model
    
    # Establish the optimizer for each worker
    global worker1_optimizer
    worker1_optimizer = optim.SGD(worker1_model.parameters(), lr=LR)
    global worker2_optimizer
    worker2_optimizer = optim.SGD(worker2_model.parameters(), lr=LR)
    global worker3_optimizer
    worker3_optimizer = optim.SGD(worker3_model.parameters(), lr=LR)
    
    # Organize the WORKER models and optimizers into lists. The server stuff must not be mixed with these
    global models
    models = [worker1_model, worker2_model, worker3_model]
    global optimizers
    optimizers = [worker1_optimizer, worker2_optimizer, worker3_optimizer]
    global worker_collection
    worker_collection = [(worker1, worker1_model, worker1_optimizer), (worker2, worker2_model, worker2_optimizer), 
                        (worker3, worker3_model, worker3_optimizer)]

In [22]:
# # Get all objects as a dictionary, as keys, or remove a specific object
# worker1.object_store._objects.keys()
# worker1.object_store.rm_obj( obj_id = )

# RUN THE MODEL

In [23]:
# Define your number of epochs
epochs = 5
epoch_times = []

# Train all workers for the set number of epochs
for epoch in range(epochs):
    
    # Start counting the time for this epoch
#     start_time = time.time()
    print(f"Epoch Number {epoch + 1}")
    
    # Train the individual models, and then obtain the federated averaged model
#     train_start_time = time.time()
    new_worker1_model, new_worker2_model, new_worker3_model, batch_times, epoch_time = train()
#     train_total_time = time.time() - train_start_time
#     print("Total TRAIN time for epoch ", epoch, " = ", 
#           round(train_total_time/60, 2), " min")
    
    # Save the epoch time
    epoch_times.append(epoch_time)
    
    # Test your new model to keep a log of how good we're doing per epoch 
    test(new_worker1_model, new_worker2_model, new_worker3_model)
    
    # Stop counting the time
#     epoch_total_time = time.time() - start_time
#     print('Time for this epoch', round(epoch_total_time/60, 2), 'min')
    
    # Re-organize everything before starting next epoch
    reset(new_worker1_model, new_worker2_model, new_worker3_model)
    
    # Save the batch times
    df_batch = pd.DataFrame(batch_times)
    df_batch.to_csv("./Batch_times/MINI_MNIST_PL_epoch_" + str(epoch) + ".csv")

# Save the epoch time
df_epoch = pd.DataFrame(epoch_times)
df_epoch.to_csv("./Epoch_times/MINI_MNIST_PL.csv")

# Clean the global namespace after run is done
%reset -f

Epoch Number 1


100%|██████████| 300/300 [10:21<00:00,  2.07s/it]


TOTAL TIME FOR THIS EPOCH =  138.5664  s
Initiated model testing:


100%|██████████| 5000/5000 [00:12<00:00, 407.97it/s]


Accuracy of the new model =  0.902 
 

Epoch Number 2


100%|██████████| 300/300 [10:19<00:00,  2.07s/it]


TOTAL TIME FOR THIS EPOCH =  135.2971  s
Initiated model testing:


100%|██████████| 5000/5000 [00:12<00:00, 413.71it/s]


Accuracy of the new model =  0.902 
 

Epoch Number 3


100%|██████████| 300/300 [10:12<00:00,  2.04s/it]


TOTAL TIME FOR THIS EPOCH =  132.6342  s
Initiated model testing:


100%|██████████| 5000/5000 [00:10<00:00, 454.57it/s]


Accuracy of the new model =  0.902 
 

Epoch Number 4


100%|██████████| 300/300 [10:40<00:00,  2.14s/it]


TOTAL TIME FOR THIS EPOCH =  145.9953  s
Initiated model testing:


100%|██████████| 5000/5000 [00:17<00:00, 277.82it/s]


Accuracy of the new model =  0.902 
 

Epoch Number 5


100%|██████████| 300/300 [10:07<00:00,  2.02s/it]


TOTAL TIME FOR THIS EPOCH =  131.3853  s
Initiated model testing:


100%|██████████| 5000/5000 [00:12<00:00, 413.91it/s]


Accuracy of the new model =  0.902 
 

