In [1]:
import pandas as pd
import pyspark.sql.functions as F
from datetime import datetime
from pyspark.sql.types import *
from pyspark import StorageLevel

import numpy as np
pd.set_option("display.max_rows", 1000)
pd.set_option("display.max_columns", 1000)
pd.set_option("mode.chained_assignment", None)

In [2]:
from pyspark.ml import Pipeline
from pyspark.ml.classification import RandomForestClassifier
from pyspark.ml.feature import IndexToString, StringIndexer, VectorIndexer
# from pyspark.ml.evaluation import MulticlassClassificationEvaluator
from pyspark.ml.feature import OneHotEncoderEstimator, StringIndexer, VectorAssembler

from pyspark.ml.evaluation import BinaryClassificationEvaluator
from pyspark.ml.tuning import CrossValidator, ParamGridBuilder

from pyspark.sql import Row
from pyspark.ml.linalg import Vectors

In [3]:
# !pip install scikit-plot

In [4]:
import sklearn
import scikitplot as skplt
from sklearn.metrics import classification_report, confusion_matrix, precision_score

<hr />
<hr />
<hr />

In [5]:
result_schema = StructType([
                    StructField('experiment_filter', StringType(), True),
                    StructField('undersampling_method', StringType(), True),
                    StructField('undersampling_column', StringType(), True),
                    StructField('filename', StringType(), True),
                    StructField('experiment_id', StringType(), True),
                    StructField('n_covid', IntegerType(), True),
                    StructField('n_not_covid', IntegerType(), True),
                    StructField('model_name', StringType(), True),
                    StructField('model_seed', StringType(), True),
                    StructField('model_maxIter', IntegerType(), True),
                    StructField('model_maxDepth', IntegerType(), True),
                    StructField('model_maxBins', IntegerType(), True),
                    StructField('model_minInstancesPerNode', IntegerType(), True),
                    StructField('model_minInfoGain', FloatType(), True),
                    StructField('model_featureSubsetStrategy', StringType(), True),
                    StructField('model_n_estimators', IntegerType(), True),
                    StructField('model_learning_rate', FloatType(), True),
                    StructField('model_impurity', StringType(), True),
                    StructField('model_AUC_ROC', StringType(), True),
                    StructField('model_AUC_PR', StringType(), True),
                    StructField('model_covid_precision', StringType(), True),
                    StructField('model_covid_recall', StringType(), True),
                    StructField('model_covid_f1', StringType(), True),
                    StructField('model_not_covid_precision', StringType(), True),
                    StructField('model_not_covid_recall', StringType(), True),
                    StructField('model_not_covid_f1', StringType(), True),
                    StructField('model_avg_precision', StringType(), True),
                    StructField('model_avg_recall', StringType(), True),
                    StructField('model_avg_f1', StringType(), True),
                    StructField('model_avg_acc', StringType(), True),
                    StructField('model_TP', StringType(), True),
                    StructField('model_TN', StringType(), True),
                    StructField('model_FN', StringType(), True),
                    StructField('model_FP', StringType(), True),
                    StructField('model_time_exec', StringType(), True),
                    StructField('model_col_set', StringType(), True)
                          ])

<hr />
<hr />
<hr />

In [6]:
cols_sets = ['cols_set_1', 'cols_set_2', 'cols_set_3']
undersamp_col = ['03-STRSAMP-AG', '04-STRSAMP-EW']
dfs = ['ds-1'] #, 'ds-2', 'ds-3']

In [7]:
# lists of params
model_numTrees = [20, 50] 
model_maxDepth = [3, 5, 7] 
model_maxBins = [32, 64]

list_of_param_dicts = []

for numTrees in model_numTrees:
    for maxDepth in model_maxDepth:
        for maxBins in model_maxBins: 
            params_dict = {}
            params_dict['numTrees'] = numTrees
            params_dict['maxDepth'] = maxDepth
            params_dict['maxBins'] = maxBins
            list_of_param_dicts.append(params_dict)

print("There is {} set of params.".format(len(list_of_param_dicts)))
# list_of_param_dicts

There is 12 set of params.


In [8]:
prefix = 'gs://ai-covid19-datalake/trusted/experiment_map/'

<hr />
<hr />
<hr />

In [9]:
# filename = 'gs://ai-covid19-datalake/trusted/experiment_map/03-STRSAMP-AG/ds-1/cols_set_1/experiment0.parquet'
# df = spark.read.parquet(filename)
# df.limit(2).toPandas()

In [10]:
# params_dict = {'numTrees': 50,
#                'maxDepth': 3,
#                'maxBins': 32}
# cols = 'cols_set_1'
# experiment_filter = 'ds-1'
# undersampling_method = '03-STRSAMP-AG', 
# experiment_id = 0

In [11]:
# run_rf(df, params_dict, cols, filename, experiment_filter, undersampling_method, experiment_id)

<hr />
<hr />
<hr />

In [12]:
def run_rf(exp_df, params_dict, cols, filename, experiment_filter, 
            undersampling_method, experiment_id):
    import time
    start_time = time.time()
    
    n_covid = exp_df.filter(F.col('CLASSI_FIN') == 1.0).count()
    n_not_covid = exp_df.filter(F.col('CLASSI_FIN') == 0.0).count()
    
    
    id_cols = ['NU_NOTIFIC', 'CLASSI_FIN']

    labelIndexer = StringIndexer(inputCol="CLASSI_FIN", outputCol="indexedLabel").fit(exp_df)    
    
    input_cols = [x for x in exp_df.columns if x not in id_cols]
    assembler = VectorAssembler(inputCols = input_cols, outputCol= 'features')
    exp_df = assembler.transform(exp_df)
    
    # Automatically identify categorical features, and index them.
    # Set maxCategories so features with > 4 distinct values are treated as continuous.
    featureIndexer = VectorIndexer(inputCol="features", outputCol="indexedFeatures", maxCategories=30).fit(exp_df)
    
    # Split the data into training and test sets (30% held out for testing)
    (trainingData, testData) = exp_df.randomSplit([0.7, 0.3])
    trainingData = trainingData.persist(StorageLevel.MEMORY_ONLY)
    testData = testData.persist(StorageLevel.MEMORY_ONLY)
    
    # Train a RandomForest model.
    rf = RandomForestClassifier(labelCol="indexedLabel", featuresCol="indexedFeatures",
                               numTrees = params_dict['numTrees'], 
                               maxDepth = params_dict['maxDepth'], 
                               maxBins = params_dict['maxBins'])
    
    # Convert indexed labels back to original labels.
    labelConverter = IndexToString(inputCol="prediction", outputCol="predictedLabel",
                               labels=labelIndexer.labels)

    # Chain indexers and forest in a Pipeline
    pipeline = Pipeline(stages=[labelIndexer, featureIndexer, rf, labelConverter])

    # Train model.  This also runs the indexers.
    model = pipeline.fit(trainingData)
    
    # Make predictions.
    predictions = model.transform(testData)    
    
    
    pred = predictions.select(['CLASSI_FIN', 'predictedLabel'])\
                  .withColumn('predictedLabel', F.col('predictedLabel').cast('double'))\
                  .withColumn('predictedLabel', F.when(F.col('predictedLabel') == 1.0, 'covid').otherwise('n-covid'))\
                  .withColumn('CLASSI_FIN', F.when(F.col('CLASSI_FIN') == 1.0, 'covid').otherwise('n-covid'))\
                  .toPandas()

    y_true = pred['CLASSI_FIN'].tolist()
    y_pred = pred['predictedLabel'].tolist()
    
    report = classification_report(y_true, y_pred, output_dict=True)
    
    
    evaluator_ROC = BinaryClassificationEvaluator(labelCol="indexedLabel", rawPredictionCol="prediction", metricName="areaUnderROC")
    accuracy_ROC = evaluator_ROC.evaluate(predictions)


    
    evaluator_PR = BinaryClassificationEvaluator(labelCol="indexedLabel", rawPredictionCol="prediction", metricName="areaUnderPR")
    accuracy_PR = evaluator_PR.evaluate(predictions)
    
    conf_matrix = confusion_matrix(y_true, y_pred)

    result_dict = {}
    
    result_dict['experiment_filter'] = experiment_filter
    result_dict['undersampling_method'] = undersampling_method
    result_dict['filename'] = filename
    result_dict['experiment_id'] = experiment_id
    result_dict['n_covid'] = n_covid
    result_dict['n_not_covid'] = n_not_covid
    result_dict['model_name'] = 'RF'
    result_dict['params'] = params_dict
    result_dict['model_AUC_ROC'] = accuracy_ROC
    result_dict['model_AUC_PR'] = accuracy_PR
    result_dict['model_covid_precision'] = report['covid']['precision']
    result_dict['model_covid_recall'] = report['covid']['recall']
    result_dict['model_covid_f1'] = report['covid']['f1-score']
    result_dict['model_not_covid_precision'] = report['n-covid']['precision']
    result_dict['model_not_covid_recall'] = report['n-covid']['recall']
    result_dict['model_not_covid_f1'] = report['n-covid']['f1-score']
    result_dict['model_avg_precision'] = report['macro avg']['precision']
    result_dict['model_avg_recall'] = report['macro avg']['recall']
    result_dict['model_avg_f1'] = report['macro avg']['f1-score']
    result_dict['model_avg_acc'] = report['accuracy']
    result_dict['model_TP'] = conf_matrix[0][0]
    result_dict['model_TN'] = conf_matrix[1][1]
    result_dict['model_FN'] = conf_matrix[0][1]
    result_dict['model_FP'] = conf_matrix[1][0]
    result_dict['model_time_exec'] = time.time() - start_time
    result_dict['model_col_set'] = cols
    
    return result_dict

<hr />
<hr />
<hr />

# Running GBT on 10 samples for each experiment
### 3x col sets -> ['cols_set_1', 'cols_set_2', 'cols_set_3']
### 3x model_maxIter -> [100, 200, 300]
### 3x model_maxDepth -> [5, 10, 15]
### 3x model_maxBins -> [16, 32, 64]
Total: 10 * 3 * 3 * 3 * 3 = 810

In [13]:
experiments = []

### Datasets: strat_samp_lab_agegrp

In [None]:
for uc in undersamp_col: 
    for ds in dfs:
        for col_set in cols_sets:
            for params_dict in list_of_param_dicts: 
                for id_exp in range(50):
                    filename = prefix + uc + '/' + ds + '/' + col_set + '/' + 'experiment' + str(id_exp) + '.parquet'
                    exp_dataframe = spark.read.parquet(filename)
                    print('read {}'.format(filename))
                    
                    undersampling_method = uc
                    experiment_filter = ds
                    experiment_id = id_exp

                    try:                     
                        model = run_rf(exp_dataframe, params_dict, col_set, filename, experiment_filter, undersampling_method, experiment_id)
                        experiments.append(model)

                        print("Parameters ==> {}\n Results: \n AUC_PR: {} \n Precision: {} \n Time: {}".format(str(params_dict), str(model['model_AUC_PR']), str(model['model_avg_precision']), str(model['model_time_exec'])))
                        print('=========================== \n')
                    except:
                        print('=========== W A R N I N G =========== \n')
                        print('Something wrong with the exp: {}, {}, {}'.format(filename, params_dict, col_set))

read gs://ai-covid19-datalake/trusted/experiment_map/03-STRSAMP-AG/ds-1/cols_set_1/experiment0.parquet
Parameters ==> {'numTrees': 20, 'maxDepth': 3, 'maxBins': 32}
 Results: 
 AUC_PR: 0.8890270259383006 
 Precision: 0.9256869778335173 
 Time: 20.402421951293945

read gs://ai-covid19-datalake/trusted/experiment_map/03-STRSAMP-AG/ds-1/cols_set_1/experiment1.parquet
Parameters ==> {'numTrees': 20, 'maxDepth': 3, 'maxBins': 32}
 Results: 
 AUC_PR: 0.8865246453136948 
 Precision: 0.9246920836637833 
 Time: 10.293002843856812

read gs://ai-covid19-datalake/trusted/experiment_map/03-STRSAMP-AG/ds-1/cols_set_1/experiment2.parquet
Parameters ==> {'numTrees': 20, 'maxDepth': 3, 'maxBins': 32}
 Results: 
 AUC_PR: 0.8867577390328641 
 Precision: 0.9252727295630192 
 Time: 14.857552766799927

read gs://ai-covid19-datalake/trusted/experiment_map/03-STRSAMP-AG/ds-1/cols_set_1/experiment3.parquet
Parameters ==> {'numTrees': 20, 'maxDepth': 3, 'maxBins': 32}
 Results: 
 AUC_PR: 0.8861001227810739 
 Pr

<hr />
<hr />
<hr />

In [None]:
for i in range(len(experiments)):
    for d in list(experiments[i].keys()):
        experiments[i][d] = str(experiments[i][d])

In [None]:
# experiments

In [None]:
cols = ['experiment_filter', 'undersampling_method', 'filename', 'experiment_id', 'n_covid', 'n_not_covid', 'model_name', 'params', 'model_AUC_ROC', 'model_AUC_PR', 'model_covid_precision', 'model_covid_recall', 'model_covid_f1', 'model_not_covid_precision', 'model_not_covid_recall', 'model_not_covid_f1', 'model_avg_precision', 'model_avg_recall', 'model_avg_f1', 'model_avg_acc', 'model_TP', 'model_TN', 'model_FN', 'model_FP', 'model_time_exec', 'model_col_set']

In [None]:
intermed_results = spark.createDataFrame(data=experiments).select(cols)
intermed_results.toPandas()



Unnamed: 0,experiment_filter,undersampling_method,filename,experiment_id,n_covid,n_not_covid,model_name,params,model_AUC_ROC,model_AUC_PR,model_covid_precision,model_covid_recall,model_covid_f1,model_not_covid_precision,model_not_covid_recall,model_not_covid_f1,model_avg_precision,model_avg_recall,model_avg_f1,model_avg_acc,model_TP,model_TN,model_FN,model_FP,model_time_exec,model_col_set
0,ds-1,03-STRSAMP-AG,gs://ai-covid19-datalake/trusted/experiment_ma...,0,86324,44538,RF,"{'numTrees': 20, 'maxDepth': 3, 'maxBins': 32}",0.870944647125077,0.8890270259383006,0.8862876254180602,0.9858172518019065,0.9334067143643368,0.9650863302489745,0.7560720424482475,0.8478880321823667,0.9256869778335173,0.870944647125077,0.8906473732733518,0.9073672391354276,25440,10117,366,3264,20.402421951293945,cols_set_1
1,ds-1,03-STRSAMP-AG,gs://ai-covid19-datalake/trusted/experiment_ma...,1,86460,44538,RF,"{'numTrees': 20, 'maxDepth': 3, 'maxBins': 32}",0.8683937780124868,0.8865246453136948,0.8851053490228747,0.9856971664927133,0.9326968799151396,0.964278818304692,0.7510903895322605,0.844436929320257,0.9246920836637833,0.8683937780124868,0.8885669046176983,0.9060433528225291,25499,9988,370,3310,10.293002843856812,cols_set_1
2,ds-1,03-STRSAMP-AG,gs://ai-covid19-datalake/trusted/experiment_ma...,2,87179,44538,RF,"{'numTrees': 20, 'maxDepth': 3, 'maxBins': 32}",0.8655394494971446,0.8867577390328641,0.8831478266789203,0.9872155848108972,0.9322865201846895,0.9673976324471182,0.7438633141833918,0.8410308321734363,0.9252727295630192,0.8655394494971445,0.8866586761790629,0.9050270883205241,25946,9970,336,3433,14.857552766799927,cols_set_1
3,ds-1,03-STRSAMP-AG,gs://ai-covid19-datalake/trusted/experiment_ma...,3,86086,44538,RF,"{'numTrees': 20, 'maxDepth': 3, 'maxBins': 32}",0.8688023424339175,0.8861001227810739,0.8838445664719667,0.9846477556109726,0.9315270481983229,0.9625332826169646,0.7529569292568623,0.8449434450519637,0.9231889245444657,0.8688023424339175,0.8882352466251433,0.9050042191934947,25270,10122,394,3321,9.830419063568115,cols_set_1
4,ds-1,03-STRSAMP-AG,gs://ai-covid19-datalake/trusted/experiment_ma...,4,84842,44538,RF,"{'numTrees': 20, 'maxDepth': 3, 'maxBins': 32}",0.8680173388579784,0.8869464317562252,0.8838934195725534,0.9858078174618732,0.9320730238161431,0.9647962656812215,0.7502268602540835,0.8440889947675162,0.9243448426268874,0.8680173388579784,0.8880810092918296,0.905372957062818,25145,9921,362,3303,9.404830932617188,cols_set_1
...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...
3595,ds-1,04-STRSAMP-EW,gs://ai-covid19-datalake/trusted/experiment_ma...,45,85323,44538,RF,"{'numTrees': 50, 'maxDepth': 7, 'maxBins': 64}",0.9059409872321467,0.8958252952689226,0.9202291758483914,0.974676158244836,0.9466704448852367,0.9449285170459352,0.8372058162194573,0.8878114692206811,0.9325788464471633,0.9059409872321467,0.9172409570529589,0.9277062152679966,25056,11170,651,2172,23.117098569869995,cols_set_3
3596,ds-1,04-STRSAMP-EW,gs://ai-covid19-datalake/trusted/experiment_ma...,46,85681,44538,RF,"{'numTrees': 50, 'maxDepth': 7, 'maxBins': 64}",0.9027541925904995,0.8941765881577614,0.9167366921193965,0.9747494367182037,0.944853425714985,0.945087437695362,0.8307589484627952,0.8842429751412876,0.9309120649073792,0.9027541925904994,0.9145482004281362,0.9252958579881657,25092,11187,650,2279,22.554786682128906,cols_set_3
3597,ds-1,04-STRSAMP-EW,gs://ai-covid19-datalake/trusted/experiment_ma...,47,84689,44538,RF,"{'numTrees': 50, 'maxDepth': 7, 'maxBins': 64}",0.9088350128897971,0.9058817051614398,0.9200881380829967,0.9793222061525232,0.9487815499971598,0.9552226172337904,0.8383478196270708,0.89297725024728,0.9376553776583936,0.9088350128897971,0.9208794001222199,0.9307191886077246,25054,11285,529,2176,22.263594150543213,cols_set_3
3598,ds-1,04-STRSAMP-EW,gs://ai-covid19-datalake/trusted/experiment_ma...,48,85958,44538,RF,"{'numTrees': 50, 'maxDepth': 7, 'maxBins': 64}",0.9061993777204436,0.9010914506558221,0.918833321057671,0.9777028880442024,0.9473544320619671,0.9513592067020004,0.8346958673966849,0.8892173704606289,0.9350962638798357,0.9061993777204436,0.918285901261298,0.9286265829300937,24950,11129,569,2204,23.29175353050232,cols_set_3


In [None]:
intermed_results.write.parquet('gs://ai-covid19-datalake/trusted/intermed_results/STRSAMP/RF_experiments-ds1.parquet', mode='overwrite')

In [None]:
print('finished')

finished


In [None]:
intermed_results.show()

+-----------------+--------------------+--------------------+-------------+-------+-----------+----------+--------------------+------------------+------------------+---------------------+------------------+------------------+-------------------------+----------------------+------------------+-------------------+------------------+------------------+------------------+--------+--------+--------+--------+------------------+-------------+
|experiment_filter|undersampling_method|            filename|experiment_id|n_covid|n_not_covid|model_name|              params|     model_AUC_ROC|      model_AUC_PR|model_covid_precision|model_covid_recall|    model_covid_f1|model_not_covid_precision|model_not_covid_recall|model_not_covid_f1|model_avg_precision|  model_avg_recall|      model_avg_f1|     model_avg_acc|model_TP|model_TN|model_FN|model_FP|   model_time_exec|model_col_set|
+-----------------+--------------------+--------------------+-------------+-------+-----------+----------+--------------