What is the best PySpark practice to predict recent time-series data & forecast next dates value?

Viewed 182

I'm experimenting with the RandomForestRegressor. The main dataframe created by Aggregations on Windows over Event-Time. Let's say below are train and test sets in the form of the dataframes:

#trainset dataframe
+-----+------------------------------------------+------------------------+
|Name |window_frame_24_Hours                     |Total_count_24_Hours    |
+-----+------------------------------------------+------------------------+
| xxx |{2021-11-28 00:00:00, 2021-11-29 00:00:00}|316666                  |
| xxx |{2021-11-27 00:00:00, 2021-11-28 00:00:00}|324526                  |
| xxx |{2021-11-26 00:00:00, 2021-11-27 00:00:00}|261414                  |
| xxx |{2021-11-25 00:00:00, 2021-11-26 00:00:00}|268632                  |
| xxx |{2021-11-24 00:00:00, 2021-11-25 00:00:00}|284578                  |
| xxx |{2021-11-23 00:00:00, 2021-11-24 00:00:00}|232226                  |
  ...                 ....                          ...
| xxx |{2021-11-02 00:00:00, 2021-11-03 00:00:00}|94100                   |
| xxx |{2021-11-01 00:00:00, 2021-11-02 00:00:00}|106666                  |
| xxx |{2021-10-31 00:00:00, 2021-11-01 00:00:00}|130108                  |
| xxx |{2021-10-30 00:00:00, 2021-10-31 00:00:00}|11235                   |
+-----+------------------------------------------+------------------------+

#testset dataframe
+-----+------------------------------------------+------------------------+
|Name |window_frame_24_Hours                     |Total_count_24_Hours    |
+-----+------------------------------------------+------------------------+
| xxx |{2021-11-28 00:00:00, 2021-11-29 00:00:00}|316666                   | <-- predict
+-----+------------------------------------------+------------------------+

#forecast dataframe
+-----+------------------------------------------+------------------------+
|Name |window_frame_24_Hours                     |Total_count_24_Hours    |
+-----+------------------------------------------+------------------------+
| xxx |{2021-11-29 00:00:00, 2021-11-30 00:00:00}|???????                 | <-- forecast
+-----+------------------------------------------+------------------------+

The problem here is I split data for tainting and test as below:

time_frame = '24'
date_from_val = '2021-10-30'
date_to_val   = '2021-11-28'
X_train = df_filtered.filter(F.col(f"window_frame_{time_frame}_Hours").between(f"{date_from_val}", f"{date_to_val}")).sort(f"window_frame_{time_frame}_Hours")
X_test  = df_filtered.filter(F.col('date') > "2021-11-07").sort('date')

following is the RandomForestRegressor model:

#Dependencies
from pyspark.ml import Pipeline
from pyspark.ml.feature import StringIndexer
from pyspark.ml.feature import VectorAssembler
from pyspark.ml.feature import MinMaxScaler
from pyspark.ml.regression import RandomForestRegressor, GBTRegressor
from pyspark.ml.evaluation import RegressionEvaluator
from pyspark.ml.tuning import CrossValidator, ParamGridBuilder

value = "Total_count_24_Hours"

assembler = VectorAssembler(inputCols=[value], outputCol='features')
scaler = MinMaxScaler(inputCol='features', outputCol='scaled_features')

# Create a RandomForest model.
rf = RandomForestRegressor(featuresCol="scaled_features",  labelCol=value)

# Chain model, assembler and scaler into a Pipeline.
pipeline = Pipeline(stages=[assembler, scaler, rf])

# Train model on training data. 
rf_model = pipeline.fit(X_train)

#save/store the model with Name and window time frame (Start and end dates in X_train dataframe)
#.... no idea


#load/restore the trained model based on Name and time-frame window data recorded and trained from the root path
#.... no idea

# Make predictions.
predictions = rf_model.transform(X_test)

# Select example rows to display.
predictions = predictions.select("window_frame_24_Hours.end", value, "prediction", "features", "scaled_features").sort('date')
##predictions.show(45, truncate = False)

# Select (prediction, true label) and compute test error
evaluator = RegressionEvaluator(labelCol=value, predictionCol="prediction", metricName="rmse")

# Error computation of model
rmse = evaluator.evaluate(predictions)

# summary of model
rfModel = rf_model.stages 

the only approach that crossed my mind is I add further rows for next days with the default value of 0 (here Total_count_24_Hours) and concat to the frame and ask the model to fit_predict those rows supposed to forecast. but I'm wondering if there is another elegant approach to this (independent of data split on date/Timestamp)?

I assume that another way could be: passing the date/Timestamp column to the index, but it would be expensive in sense of the Spark process. That would be much pythonic approach rather than Spark.

0 Answers
Related