In [1]:
# autoreload
%load_ext autoreload
%autoreload 2

# set current working directory
import os
os.chdir('..')

import src.config as config

In [2]:
# connect to hopsworks feature store
import hopsworks

# connect to project
project = hopsworks.login(project=config.HOPSWORKS_PROJECT_NAME, api_key_value=config.HOPSWORKS_API_KEY)

# connect to feature store
feature_store = project.get_feature_store()

# connect to the feature group
feature_group = feature_store.get_or_create_feature_group(
    name=config.FEATURE_GROUP_NAME,
    version=config.FEATURE_GROUP_VERSION,
    description='Time series data at hourly frequency',
    primary_key=['pickup_datetime', 'pickup_hour'],
    event_time='pickup_hour',)

Connected. Call `.close()` to terminate connection gracefully.

Logged in to project, explore it here https://c.app.hopsworks.ai:443/p/20648
Connected. Call `.close()` to terminate connection gracefully.


In [3]:
# create feature view (if it doesn't exist)

try:
    # create feature view if it doesn't exist
    feature_store.create_feature_view(
        name=config.FEATURE_VIEW_NAME,
        version=config.FEATURE_VIEW_VERSION,
        query=feature_group.select_all()
    )
except:
    print('Feature view already exists. Skipping creation.')

# get feature view
feature_view = feature_store.get_feature_view(config.FEATURE_VIEW_NAME, config.FEATURE_VIEW_VERSION)

Feature view already exists. Skipping creation.


In [4]:
ts_data, _ = feature_view.training_data(description='Time series hourly taxi rides')

2023-03-04 12:46:58,081 INFO: USE `taxi_demand_1_featurestore`
2023-03-04 12:46:58,422 INFO: SELECT `fg0`.`pickup_hour` `pickup_hour`, `fg0`.`rides` `rides`, `fg0`.`pickup_location_id` `pickup_location_id`
FROM `taxi_demand_1_featurestore`.`time_series_hourly_feature_group_1` `fg0`




In [5]:
ts_data.sort_values(by=['pickup_location_id', 'pickup_hour'], inplace=True)
ts_data

Unnamed: 0,pickup_hour,rides,pickup_location_id
2080335,2022-01-01 00:00:00,0,1
2207139,2022-01-01 01:00:00,0,1
1940894,2022-01-01 02:00:00,0,1
1896226,2022-01-01 03:00:00,0,1
1418782,2022-01-01 04:00:00,1,1
...,...,...,...
210063,2023-03-04 11:00:00,5,265
228628,2023-03-04 12:00:00,4,265
392227,2023-03-04 13:00:00,6,265
197414,2023-03-04 14:00:00,7,265


In [6]:
from src.data import create_ts_dataset

features, targets = create_ts_dataset(
    ts_data,
    n_features=24*28, # 1 month
    step_size=23)

features_and_target = features.copy()
features_and_target['target_rides_next_hour'] = targets

print(f'{features_and_target.shape=}')

100%|██████████| 262/262 [10:19<00:00,  2.36s/it]


features_and_target.shape=(91107, 675)


In [7]:
from datetime import date, timedelta
from pytz import timezone
import pandas as pd
from src.data_split import train_test_split

# training data range: January 2022 to Current Date - 1 month
# test data range: Current Date - 1 month to Current Date
cutoff_date = pd.to_datetime(date.today() - timedelta(days=28))

print(f'{cutoff_date=}')

X_train, y_train, X_test, y_test = train_test_split(
    df=features_and_target,
    cutoff_date=cutoff_date,
    target_column_name='target_rides_next_hour')

print(f'{X_train.shape=}')
print(f'{y_train.shape=}')
print(f'{X_test.shape=}')
print(f'{y_test.shape=}')

cutoff_date=Timestamp('2023-02-04 00:00:00')
X_train.shape=(83658, 674)
y_train.shape=(83658,)
X_test.shape=(7449, 674)
y_test.shape=(7449,)


In [8]:
import numpy as np

from sklearn.model_selection import TimeSeriesSplit
from sklearn.metrics import mean_absolute_error
import optuna

from src import model

# define objective function
def objective(trial: optuna.trial.Trial) -> float:
    '''Takes in hyperparameters as input, and trains a model that computes the average validation error based on TimeSeriesSplit cross validation'''

    # define hyperparameters
    params = {
        "num_leaves": trial.suggest_int("num_leaves", 2, 256),
        "colsample_bytree": trial.suggest_float("colsample_bytree", 0.2, 1.0),
        "subsample": trial.suggest_float("subsample", 0.2, 1.0),
        "min_child_samples": trial.suggest_int("min_child_samples", 3, 100),   
    }

    tss = TimeSeriesSplit(n_splits=4)
    scores = []
    for train_index, val_index in tss.split(X_train):
        # split data
        X_train_, X_val = X_train.iloc[train_index], X_train.iloc[val_index]
        y_train_, y_val = y_train.iloc[train_index], y_train.iloc[val_index]

        # create model
        pipeline = model.get_pipeline(**params)

        # fit model
        pipeline.fit(X_train_, y_train_)

        # compute validation error
        y_pred = pipeline.predict(X_val)
        mae = mean_absolute_error(y_val, y_pred)

        scores.append(mae)
    
    return np.mean(scores)

In [10]:
import warnings
warnings.filterwarnings('ignore')

# optuna study
study = optuna.create_study(direction='minimize', study_name='lightgbm')
study.optimize(objective, n_trials=1)

[32m[I 2023-03-04 13:14:08,576][0m A new study created in memory with name: lightgbm[0m
[32m[I 2023-03-04 13:14:45,166][0m Trial 0 finished with value: 3.1604979000753546 and parameters: {'num_leaves': 241, 'colsample_bytree': 0.7747895635208406, 'subsample': 0.7278038159065294, 'min_child_samples': 75}. Best is trial 0 with value: 3.1604979000753546.[0m


In [11]:
# print best parameters
best_params = study.best_trial.params
print(f'{best_params=}')

best_params={'num_leaves': 241, 'colsample_bytree': 0.7747895635208406, 'subsample': 0.7278038159065294, 'min_child_samples': 75}


In [12]:
# fit best params on full training set
pipeline = model.get_pipeline(**best_params)
pipeline.fit(X_train, y_train)

In [13]:
# compute test error on test set
predictions = pipeline.predict(X_test)
test_mae = mean_absolute_error(y_test, predictions)
print(f'{test_mae=:.4f}')

test_mae=5.5708


In [14]:
# save trained model
import joblib
from src.paths import MODELS_DIR

joblib.dump(pipeline, MODELS_DIR / 'model.pkl')

['/Users/ani/Projects/taxi_demand_forecasting/models/model.pkl']

In [15]:
# define schema for hopsworks model reigistry
from hsml.schema import Schema
from hsml.model_schema import ModelSchema

input_schema = Schema(X_train)
output_schema = Schema(y_train)
model_schema = ModelSchema(input_schema, output_schema)


In [26]:
os.chdir('..')

# upload model to hopsworks model registry
model_registry = project.get_model_registry()

model = model_registry.sklearn.create_model(
    name='taxt_demand_forecaster_next_hour',
    metrics={'test_mae': test_mae},
    description='LightGBM model that predicts the number of taxi rides in the next hour',
    model_schema=model_schema,
    input_example=X_train.sample()
)

model.save(MODELS_DIR / 'model.pkl')

Connected. Call `.close()` to terminate connection gracefully.


  0%|          | 0/6 [00:00<?, ?it/s]

Model created, explore it at https://c.app.hopsworks.ai:443/p/20648/models/taxt_demand_forecaster_next_hour/1


Model(name: 'taxt_demand_forecaster_next_hour', version: 1)