In [1]:
from pyspark.sql import SparkSession

In [2]:
spark = SparkSession.builder.appName('lrex').getOrCreate()

In [3]:
from pyspark.ml.regression import LinearRegression

## check how the data looks like

In [4]:
training = spark.read.format('libsvm').load('sample_linear_regression_data.txt')

In [5]:
training.show(5)

+-------------------+--------------------+
|              label|            features|
+-------------------+--------------------+
| -9.490009878824548|(10,[0,1,2,3,4,5,...|
| 0.2577820163584905|(10,[0,1,2,3,4,5,...|
| -4.438869807456516|(10,[0,1,2,3,4,5,...|
|-19.782762789614537|(10,[0,1,2,3,4,5,...|
| -7.966593841555266|(10,[0,1,2,3,4,5,...|
+-------------------+--------------------+
only showing top 5 rows



In [6]:
training.head(1)

[Row(label=-9.490009878824548, features=SparseVector(10, {0: 0.4551, 1: 0.3664, 2: -0.3826, 3: -0.4458, 4: 0.3311, 5: 0.8067, 6: -0.2624, 7: -0.4485, 8: -0.0727, 9: 0.5658}))]

In [10]:
training.head(1)[0][1]

SparseVector(10, {0: 0.4551, 1: 0.3664, 2: -0.3826, 3: -0.4458, 4: 0.3311, 5: 0.8067, 6: -0.2624, 7: -0.4485, 8: -0.0727, 9: 0.5658})

In [11]:
# SparseVector(4, [1, 3], [3.0, 4.0] --> a sparse vector that has 4 elements with element 1 = 3, element 3 = 4
# index starts from 0

## Building a model

In [8]:
lr = LinearRegression(featuresCol='features',labelCol='label',predictionCol='prediction')

In [9]:
type(lr)

pyspark.ml.regression.LinearRegression

In [10]:
lrModel = lr.fit(training)

In [11]:
type(lrModel)

pyspark.ml.regression.LinearRegressionModel

In [12]:
lrModel.coefficients

DenseVector([0.0073, 0.8314, -0.8095, 2.4412, 0.5192, 1.1535, -0.2989, -0.5129, -0.6197, 0.6956])

In [13]:
# try checking different model's attributes
lrModel.intercept

0.14228558260358093

In [14]:
training_summary = lrModel.summary

In [15]:
training_summary.predictions.show(5)

+-------------------+--------------------+-------------------+
|              label|            features|         prediction|
+-------------------+--------------------+-------------------+
| -9.490009878824548|(10,[0,1,2,3,4,5,...| 1.5211201432720063|
| 0.2577820163584905|(10,[0,1,2,3,4,5,...|-0.6658770747591632|
| -4.438869807456516|(10,[0,1,2,3,4,5,...| 0.1568703823211514|
|-19.782762789614537|(10,[0,1,2,3,4,5,...| 0.6374146679690593|
| -7.966593841555266|(10,[0,1,2,3,4,5,...|  2.372566473232916|
+-------------------+--------------------+-------------------+
only showing top 5 rows



In [16]:
# try checking different summary attributes
training_summary.r2

0.027839179518600154

## This time let's split the data so we can compare accuracies of the training and test sets

In [17]:
all_data = spark.read.format('libsvm').load('sample_linear_regression_data.txt')

In [18]:
type(all_data)

pyspark.sql.dataframe.DataFrame

In [19]:
# split all_data to 0.7 and 0.3 portions
split_object = all_data.randomSplit([0.7,0.3])

In [20]:
split_object

[DataFrame[label: double, features: vector],
 DataFrame[label: double, features: vector]]

In [21]:
# let's do it properly
train_data, test_data = all_data.randomSplit([0.7,0.3])

In [22]:
train_data.show(5)

+-------------------+--------------------+
|              label|            features|
+-------------------+--------------------+
|-28.046018037776633|(10,[0,1,2,3,4,5,...|
|-26.805483428483072|(10,[0,1,2,3,4,5,...|
|-26.736207182601724|(10,[0,1,2,3,4,5,...|
| -23.51088409032297|(10,[0,1,2,3,4,5,...|
|-23.487440120936512|(10,[0,1,2,3,4,5,...|
+-------------------+--------------------+
only showing top 5 rows



In [23]:
train_data.describe().show()
test_data.describe().show()
all_data.describe().show()

+-------+-------------------+
|summary|              label|
+-------+-------------------+
|  count|                343|
|   mean|-0.2780057485727789|
| stddev| 10.332774128919448|
|    min|-28.046018037776633|
|    max|  27.78383192005107|
+-------+-------------------+

+-------+-------------------+
|summary|              label|
+-------+-------------------+
|  count|                158|
|   mean| 1.4180839979756525|
| stddev| 10.221788321598288|
|    min|-28.571478869743427|
|    max| 27.111027963108548|
+-------+-------------------+

+-------+-------------------+
|summary|              label|
+-------+-------------------+
|  count|                501|
|   mean|0.25688882219498976|
| stddev| 10.317884030544564|
|    min|-28.571478869743427|
|    max|  27.78383192005107|
+-------+-------------------+



In [24]:
correct_model = lr.fit(train_data)

## Evaluate the model on the test data and compare root mean squared errors of the train and test sets

In [25]:
test_results = correct_model.evaluate(test_data)

In [26]:
print('Residual elements = {}'.format(test_results.residuals.count()))
test_results.residuals.show(5)

Residual elements = 158
+-------------------+
|          residuals|
+-------------------+
| -27.68536554356405|
|  -21.6050382526352|
|-22.294637935852606|
| -19.72928855465974|
|-21.393455301067515|
+-------------------+
only showing top 5 rows



In [27]:
print('Root RMS of train data: {}'.format(correct_model.summary.rootMeanSquaredError))
print('Root RMS of test data: {}'.format(test_results.rootMeanSquaredError))

Root RMS of train data: 10.09535702596368
Root RMS of test data: 10.488198014490845


In [28]:
print('Prediction on train data')
correct_model.summary.predictions.show(5)

print('Prediction on test data')
test_results.predictions.show(5)

Prediction on train data
+-------------------+--------------------+-------------------+
|              label|            features|         prediction|
+-------------------+--------------------+-------------------+
|-28.046018037776633|(10,[0,1,2,3,4,5,...|-2.2281026197055427|
|-26.805483428483072|(10,[0,1,2,3,4,5,...|-0.7847575421570169|
|-26.736207182601724|(10,[0,1,2,3,4,5,...|-3.7364820138848494|
| -23.51088409032297|(10,[0,1,2,3,4,5,...|-2.8147375536758923|
|-23.487440120936512|(10,[0,1,2,3,4,5,...| -3.246096788680554|
+-------------------+--------------------+-------------------+
only showing top 5 rows

Prediction on test data
+-------------------+--------------------+-------------------+
|              label|            features|         prediction|
+-------------------+--------------------+-------------------+
|-28.571478869743427|(10,[0,1,2,3,4,5,...|-0.8861133261793763|
|-21.432387764165806|(10,[0,1,2,3,4,5,...|0.17265048846939315|
|-20.212077258958672|(10,[0,1,2,3,4,5,...| 2

## Deploy the model on unlabeled data

In [29]:
# create unlabeled data using test data
unlabeled_data = test_data.select('features')
unlabeled_data.show(5)

+--------------------+
|            features|
+--------------------+
|(10,[0,1,2,3,4,5,...|
|(10,[0,1,2,3,4,5,...|
|(10,[0,1,2,3,4,5,...|
|(10,[0,1,2,3,4,5,...|
|(10,[0,1,2,3,4,5,...|
+--------------------+
only showing top 5 rows



In [30]:
# get predictions on unlabeled data
predictions = correct_model.transform(unlabeled_data)

In [31]:
predictions.show(10)

+--------------------+--------------------+
|            features|          prediction|
+--------------------+--------------------+
|(10,[0,1,2,3,4,5,...| -0.8861133261793763|
|(10,[0,1,2,3,4,5,...| 0.17265048846939315|
|(10,[0,1,2,3,4,5,...|  2.0825606768939338|
|(10,[0,1,2,3,4,5,...|  0.3269525244451873|
|(10,[0,1,2,3,4,5,...|   2.547532828168932|
|(10,[0,1,2,3,4,5,...| -0.5926929560308424|
|(10,[0,1,2,3,4,5,...| -2.5747636207600593|
|(10,[0,1,2,3,4,5,...|-0.27941146762969327|
|(10,[0,1,2,3,4,5,...|   -1.45785877776201|
|(10,[0,1,2,3,4,5,...|   1.481816936891973|
+--------------------+--------------------+
only showing top 10 rows

