In [None]:
#Imports
import pandas as pd
import os
import tensorflow as tf

from keras import backend as K

import sys  
sys.path.append("../../../")  
from utils.models import *
from utils.datahandling import *
from utils.modelrunner import *
from utils.federatedrunner import *

import wandb
import logging
logging.getLogger("wandb").setLevel(logging.ERROR)
logging.getLogger('tensorflow').setLevel(logging.ERROR)

os.environ['WANDB_SILENT'] = 'true'
os.environ['WANDB_CONSOLE'] = 'off'
os.environ['TF_CPP_MIN_LOG_LEVEL'] = '2'

### Data processing

In [None]:
#Get data 
num_users = 30

cwd = os.path.normpath(os.path.dirname(os.path.dirname(os.path.dirname(os.path.dirname(os.getcwd())))))
df = pd.read_csv(cwd+'/data/3final_data/Final_PV_dataset.csv', index_col='Date')
df.index = pd.to_datetime(df.index)
df.fillna(0, inplace=True)

# Get the first date from the index
start_date = df.index.min()
# Calculate the end date as one year from the start date
end_date = start_date + pd.DateOffset(years=1)
# Filter the dataframe to only include the first year of data
df = df[(df.index >= start_date) & (df.index < end_date)]

df_array = []
for idx in range(num_users):
    df_array.append(df[[f'User{idx+1}', 'temp', 'rhum', 'wspd', 'PC1', 'hour sin', 'hour cos', f'User{idx+1}_lag_24hrs']])

df_array[0].head(1)

In [None]:
#Hyperparameters
sequence_length = 25
batch_size = 16
num_features = df_array[0].shape[1]
horizon = 1
max_epochs = 100

loss = tf.keras.losses.MeanSquaredError()
metrics=[
    tf.keras.metrics.RootMeanSquaredError(), 
    tf.keras.metrics.MeanAbsolutePercentageError(),
    tf.keras.metrics.MeanAbsoluteError(),
]

early_stopping = tf.keras.callbacks.EarlyStopping(monitor='val_loss',patience=10,mode='min')
timing_callback = TimingCallback()
custom_callback = CustomCallback()

callbacks=[early_stopping, timing_callback, custom_callback]

#BiLSTM
lstm_architecture = "bilstm"
lstm_layers = 2
lstm_units = 8
lstm_all_results = pd.DataFrame(columns=["user", "architecture", "train_time", "avg_time_epoch", "mse", "rmse", "mape", "mae"])
lstm_results = pd.DataFrame(columns=['architecture', 'train_time', 'avg_time_epoch', 'mse','mse_std', 'rmse','rmse_std','mape','mape_std','mae','mae_std'])

#CNN
cnn_architecture = "CNN_Model"
cnn_layers = 4
cnn_kernel_size = 1
cnn_filter_size = 8
cnn_dense_units = 16
cnn_all_results = pd.DataFrame(columns=["user", "architecture", "train_time", "avg_time_epoch", "mse", "rmse", "mape", "mae"])
cnn_results = pd.DataFrame(columns=['architecture', 'train_time', 'avg_time_epoch', 'mse','mse_std', 'rmse','rmse_std','mape','mape_std','mae','mae_std'])

#Transfotransformer_architecture = "Transformer_ED2_h4_d32"
transformer_layers = 2
transformer_heads = 4
transformer_dense_units = 16
transformer_all_results = pd.DataFrame(columns=["user", "architecture", "train_time", "avg_time_epoch", "mse", "rmse", "mape", "mae"])
transformer_results = pd.DataFrame(columns=['architecture', 'train_time', 'avg_time_epoch', 'mse','mse_std', 'rmse','rmse_std','mape','mape_std','mae','mae_std'])

In [None]:
#Train, Validation and Test datasets
X_train, y_train, X_val, y_val, X_test, y_test = {}, {}, {}, {}, {}, {}

#Create Train, Validation and Test datasets
for idx, df in enumerate(df_array):
    n = len(df)
    train_df = df[0:int(n*0.7)]
    val_df = df[int(n*0.7):int(n*0.9)]
    test_df = df[int(n*0.9):]

    # Min max sclaing
    train_df = min_max_scaling(train_df)
    val_df = min_max_scaling(val_df)
    test_df = min_max_scaling(test_df)

    # Sequencing
    train_sequences = create_sequences(train_df, sequence_length)
    val_sequences = create_sequences(val_df, sequence_length)
    test_sequences = create_sequences(test_df, sequence_length)

    #Split into feature and label
    X_train[f'user{idx+1}'], y_train[f'user{idx+1}'] = prepare_data(train_sequences, batch_size)
    X_val[f'user{idx+1}'], y_val[f'user{idx+1}'] = prepare_data(val_sequences, batch_size)
    X_test[f'user{idx+1}'], y_test[f'user{idx+1}'] = prepare_data(test_sequences, batch_size)

### Federated Learning for Benchmark models

In [None]:
federated_rounds = 3
num_rounds = 5

In [None]:
y = np.loadtxt(f'../../../../data/3final_data/Clusters_KMeans10_dtw.csv', delimiter=',').astype(int)
num_clusters = 10
cluster_users = {i: [] for i in range(num_clusters)}

# Iterate through each cluster
for cluster_number in range(num_clusters):
    users_in_cluster = np.where(y == cluster_number)[0] +1
    cluster_users[cluster_number] = users_in_cluster
cluster_users

In [None]:
#Bilstm
run_federated_benchmark_model(
    num_clusters = num_clusters,
    federated_rounds = federated_rounds,
    cluster_users = cluster_users,
    save_path = os.getcwd(),
    wb_project_name = "TS_FL_PV_Forecasting_Benchmark",
    wb_project = "TS_FL_PV",
    wb_model_name = "global_bilstm",
    df_array = df_array,
    max_epochs = max_epochs,
    batch_size = batch_size,
    X_train = X_train,
    horizon = horizon,
    loss = loss,
    metrics = metrics,
    y_train = y_train,
    X_val = X_val,
    y_val = y_val,
    X_test = X_test,
    y_test = y_test,
    callbacks = callbacks,
    layers = lstm_layers,
    lstm_units = lstm_units
)
 
evaluate_federated_benchmark_model(
    federated_rounds = federated_rounds,
    save_path = os.getcwd(),
    wb_project = "TS_FL_PV",
    wb_model_name = "global_bilstm",
    cluster_users = cluster_users,
    num_rounds = num_rounds,
    horizon = horizon,
    batch_size = batch_size,
    metrics = metrics,
    loss = loss,
    X_train = X_train,
    y_train = y_train,
    X_val = X_val,
    y_val = y_val,
    X_test = X_test,
    y_test = y_test,
    callbacks = callbacks,
    df_array = df_array,
    results = lstm_results,
    all_results = lstm_all_results,
    layers = lstm_layers,
    lstm_units = lstm_units,
    wb_project_name="TS_FL_PV_Forecasting_Benchmark"
)  

In [None]:
#cnn
run_federated_benchmark_model(
    num_clusters = num_clusters,
    federated_rounds = federated_rounds,
    cluster_users = cluster_users,
    save_path = os.getcwd(),
    wb_project_name = "TS_FL_PV_Forecasting_Benchmark",
    wb_project = "TS_FL_PV",
    wb_model_name = "global_cnn",
    df_array = df_array,
    max_epochs = max_epochs,
    batch_size = batch_size,
    X_train = X_train,
    horizon = horizon,
    loss = loss,
    metrics = metrics,
    y_train = y_train,
    X_val = X_val,
    y_val = y_val,
    X_test = X_test,
    y_test = y_test,
    callbacks = callbacks,
    layers = cnn_layers,
    cnn_filter = cnn_filter_size,
    cnn_kernel_size = cnn_kernel_size,
    cnn_dense_units=cnn_dense_units
)

evaluate_federated_benchmark_model(
    federated_rounds = federated_rounds,
    save_path = os.getcwd(),
    wb_project = "TS_FL_PV",
    wb_model_name = "global_cnn",
    cluster_users = cluster_users,
    num_rounds = num_rounds,
    horizon = horizon,
    batch_size = batch_size,
    metrics = metrics,
    loss = loss,
    X_train = X_train,
    y_train = y_train,
    X_val = X_val,
    y_val = y_val,
    X_test = X_test,
    y_test = y_test,
    callbacks = callbacks,
    df_array = df_array,
    results = cnn_results,
    all_results = cnn_all_results,
    layers = cnn_layers,
    cnn_filter = cnn_filter_size,
    cnn_kernel_size = cnn_kernel_size,
    cnn_dense_units=cnn_dense_units,
    wb_project_name="TS_FL_PV_Forecasting_Benchmark"
)

In [None]:
#transformer
run_federated_benchmark_model(
    num_clusters = num_clusters,
    federated_rounds = federated_rounds,
    cluster_users = cluster_users,
    save_path = os.getcwd(),
    wb_project_name = "TS_FL_PV_Forecasting_Benchmark",
    wb_project = "TS_FL_PV",
    wb_model_name = "global_transformer",
    df_array = df_array,
    max_epochs = max_epochs,
    batch_size = batch_size,
    X_train = X_train,
    horizon = horizon,
    loss = loss,
    metrics = metrics,
    y_train = y_train,
    X_val = X_val,
    y_val = y_val,
    X_test = X_test,
    y_test = y_test,
    callbacks = callbacks,
    layers=transformer_layers,
    transformer_dense_units= transformer_dense_units,
    sequence_length = sequence_length,
    transformer_num_features = num_features,
    transformer_num_heads = transformer_heads
)

evaluate_federated_benchmark_model(
    federated_rounds = federated_rounds,
    save_path = os.getcwd(),
    wb_project = "TS_FL_PV",
    wb_model_name = "global_transformer",
    cluster_users = cluster_users,
    num_rounds = num_rounds,
    horizon = horizon,
    batch_size = batch_size,
    metrics = metrics,
    loss = loss,
    X_train = X_train,
    y_train = y_train,
    X_val = X_val,
    y_val = y_val,
    X_test = X_test,
    y_test = y_test,
    callbacks = callbacks,
    df_array = df_array,
    results = transformer_results,###
    all_results = transformer_all_results, ###
    layers=transformer_layers,
    transformer_dense_units= transformer_dense_units,
    sequence_length = sequence_length,
    transformer_num_features = num_features,
    transformer_num_heads = transformer_heads,
    wb_project_name="TS_FL_PV_Forecasting_Benchmark"
)