## Creating a Classification Model

In this exercise, you will implement a classification model that uses features of a flight to predict whether or not it will be late.

### Import Spark SQL and Spark ML Libraries

First, import the libraries you will need:

In [1]:
from pyspark.sql import SparkSession

spark = SparkSession.\
        builder.\
        appName("pyspark-notebook").\
        master("spark://spark-master:7077").\
        config("spark.executor.memory", "4098m").\
        getOrCreate()

In [3]:
!pip install numpy

Collecting numpy
  Downloading numpy-1.20.2-cp37-cp37m-manylinux2010_x86_64.whl (15.3 MB)
[K     |████████████████████████████████| 15.3 MB 4.6 MB/s eta 0:00:01    |███████████▌                    | 5.5 MB 5.8 MB/s eta 0:00:02     |█████████████████████▏          | 10.1 MB 4.3 MB/s eta 0:00:02
[?25hInstalling collected packages: numpy
Successfully installed numpy-1.20.2


In [4]:
from pyspark.sql.types import *
from pyspark.sql.functions import *

from pyspark.ml.classification import LogisticRegression
from pyspark.ml.feature import StringIndexer, VectorAssembler

### Load Source Data
The data for this exercise is provided as a CSV file containing details of flights that has already been cleaned up for modeling. The data includes specific characteristics (or *features*) for each flight, as well as a *label* column indicating whether or not the flight was late (a flight with an arrival delay of more than 25 minutes is considered *late*).

You will load this data into a dataframe and display it.

In [5]:
flightSchema = StructType([
  StructField("DayofMonth", IntegerType(), False),
  StructField("DayOfWeek", IntegerType(), False),
  StructField("Carrier", StringType(), False),
  StructField("OriginAirportID", IntegerType(), False),
  StructField("DestAirportID", IntegerType(), False),
  StructField("DepDelay", IntegerType(), False),
  StructField("ArrDelay", IntegerType(), False),
  StructField("Late", IntegerType(), False),
])

data = spark.read.csv('../data/flights.csv', schema=flightSchema, header=True)
data.show()

+----------+---------+-------+---------------+-------------+--------+--------+----+
|DayofMonth|DayOfWeek|Carrier|OriginAirportID|DestAirportID|DepDelay|ArrDelay|Late|
+----------+---------+-------+---------------+-------------+--------+--------+----+
|        21|        2|     WN|          10721|        13342|      26|      57|   1|
|        13|        1|     AA|          15016|        12892|      51|      27|   1|
|         5|        5|     FL|          10397|        11433|       9|       4|   0|
|        22|        1|     US|          11278|        14100|      35|      71|   1|
|        23|        4|     WN|          12451|        10693|       9|       5|   0|
|         5|        7|     AA|          11298|        15016|      39|      42|   1|
|         4|        6|     UA|          13930|        14307|      71|      58|   1|
|        10|        3|     9E|          14307|        11433|      68|     140|   1|
|        29|        7|     UA|          14057|        14771|     130|     12

### Split the Data
It is common practice when building supervised machine learning models to split the source data, using some of it to train the model and reserving some to test the trained model. In this exercise, you will use 70% of the data for training, and reserve 30% for testing.

In [6]:
splits = data.randomSplit([0.7, 0.3])
train = splits[0]
test = splits[1]
train_rows = train.count()
test_rows = test.count()
print ("Training Rows:", train_rows, " Testing Rows:", test_rows)

Training Rows: 446182  Testing Rows: 190901


### Prepare the Training Data
To train the classification model, you need a training data set that includes a vector of numeric features, and a label column. In this exercise, you will use the **StringIndexer** class to generate a numeric category for each discrete **Carrier** string value, and then use the **VectorAssembler** class to transform the numeric features that would be available for a flight that hasn't yet arrived into a vector, and then rename the **Late** column to **label** as this is what we're going to try to predict.

*Note: This is a deliberately simple example. In reality you'd likely perform mulitple data preparation steps, and later in this course we'll examine how to encapsulate these steps in to a pipeline. For now, we'll just use the numeric features as they are to define the training dataset.*

In [7]:
# Carrier is a string, and we need our features to be numeric - so we'll generate a numeric index for each distinct carrier string, and transform the dataframe to add that as a column
carrierIndexer = StringIndexer(inputCol="Carrier", outputCol="CarrierIdx")
numTrain = carrierIndexer.fit(train).transform(train)

# Now we'll assemble a vector of all the numeric feature columns (other than ArrDelay, which we wouldn't have for enroute flights)
assembler = VectorAssembler(inputCols = ["DayofMonth", "DayOfWeek", "CarrierIdx", "OriginAirportID", "DestAirportID", "DepDelay"], outputCol="features")
training = assembler.transform(numTrain).select(col("features"), col("Late").alias("label"))
training.show()

+--------------------+-----+
|            features|label|
+--------------------+-----+
|[1.0,1.0,10.0,103...|    1|
|[1.0,1.0,10.0,103...|    0|
|[1.0,1.0,10.0,105...|    1|
|[1.0,1.0,10.0,107...|    1|
|[1.0,1.0,10.0,107...|    1|
|[1.0,1.0,10.0,107...|    1|
|[1.0,1.0,10.0,107...|    1|
|[1.0,1.0,10.0,107...|    0|
|[1.0,1.0,10.0,108...|    1|
|[1.0,1.0,10.0,108...|    1|
|[1.0,1.0,10.0,108...|    0|
|[1.0,1.0,10.0,110...|    0|
|[1.0,1.0,10.0,110...|    0|
|[1.0,1.0,10.0,110...|    1|
|[1.0,1.0,10.0,110...|    0|
|[1.0,1.0,10.0,111...|    1|
|[1.0,1.0,10.0,111...|    1|
|[1.0,1.0,10.0,111...|    0|
|[1.0,1.0,10.0,111...|    1|
|[1.0,1.0,10.0,111...|    1|
+--------------------+-----+
only showing top 20 rows



### Train a Classification Model
Next, you need to train a classification model using the training data. To do this, create an instance of the classification algorithm you want to use and use its **fit** method to train a model based on the training dataframe. In this exercise, you will use a *Logistic Regression* classification algorithm - but you can use the same technique for any of the classification algorithms supported in the spark.ml API.

In [8]:
lr = LogisticRegression(labelCol="label",featuresCol="features",maxIter=10,regParam=0.3)
model = lr.fit(training)
print ("Model trained!")

Model trained!


### Prepare the Testing Data
Now that you have a trained model, you can test it using the testing data you reserved previously. First, you need to prepare the testing data in the same way as you did the training data by transforming the feature columns into a vector. This time you'll rename the **Late** column to **trueLabel**.

In [9]:
# Transform the test data to add the numeric carrier index
numTest = carrierIndexer.fit(test).transform(test)

# Generate the features vector and label
testing = assembler.transform(numTest).select(col("features"), col("Late").alias("trueLabel"))
testing.show()

+--------------------+---------+
|            features|trueLabel|
+--------------------+---------+
|[1.0,1.0,10.0,107...|        0|
|[1.0,1.0,10.0,110...|        1|
|[1.0,1.0,10.0,110...|        1|
|[1.0,1.0,10.0,111...|        1|
|[1.0,1.0,10.0,111...|        1|
|[1.0,1.0,10.0,111...|        0|
|[1.0,1.0,10.0,111...|        0|
|[1.0,1.0,10.0,111...|        0|
|[1.0,1.0,10.0,112...|        1|
|[1.0,1.0,10.0,114...|        0|
|[1.0,1.0,10.0,114...|        1|
|[1.0,1.0,10.0,114...|        0|
|[1.0,1.0,10.0,114...|        1|
|[1.0,1.0,10.0,114...|        1|
|[1.0,1.0,10.0,114...|        1|
|[1.0,1.0,10.0,123...|        1|
|[1.0,1.0,10.0,124...|        0|
|[1.0,1.0,10.0,124...|        1|
|[1.0,1.0,10.0,124...|        1|
|[1.0,1.0,10.0,124...|        1|
+--------------------+---------+
only showing top 20 rows



### Test the Model
Now you're ready to use the **transform** method of the model to generate some predictions. You can use this approach to predict delay status for flights where the label is unknown; but in this case you are using the test data which includes a known true label value, so you can compare the predicted status to the actual status.

In [10]:
prediction = model.transform(testing)
predicted = prediction.select("features", "probability", col("prediction").cast("Int"), "trueLabel")
predicted.show(100, truncate=False)

+------------------------------------+----------------------------------------+----------+---------+
|features                            |probability                             |prediction|trueLabel|
+------------------------------------+----------------------------------------+----------+---------+
|[1.0,1.0,10.0,10721.0,11193.0,-5.0] |[0.5906470971793083,0.40935290282069176]|0         |0        |
|[1.0,1.0,10.0,11057.0,12478.0,31.0] |[0.4310504674371049,0.5689495325628952] |1         |1        |
|[1.0,1.0,10.0,11057.0,12478.0,104.0]|[0.166136817949245,0.8338631820507549]  |1         |1        |
|[1.0,1.0,10.0,11193.0,10529.0,16.0] |[0.4974853311138068,0.5025146688861932] |1         |1        |
|[1.0,1.0,10.0,11193.0,10821.0,62.0] |[0.2993884409863616,0.7006115590136384] |1         |1        |
|[1.0,1.0,10.0,11193.0,11618.0,20.0] |[0.48061079844947374,0.5193892015505263]|1         |0        |
|[1.0,1.0,10.0,11193.0,13244.0,3.0]  |[0.560180371400626,0.4398196285993739]  |0         |0

Looking at the result, the **prediction** column contains the predicted value for the label, and the **trueLabel** column contains the actual known value from the testing data. The **probability** column shows the probability score for each class (0 or 1). It looks like there are a mix of correct and incorrect predictions, and the ones that are incorrect tend to have fairly close probabilities for each class. Later in this course you'll learn how to measure the accuracy of a model.