In [1]:
# Reg fetch new batch of features and compute predictions and save to feature store
# 

In [1]:
%load_ext autoreload
%autoreload 2

In [2]:
import sys
import os

# Add the parent directory to the Python path
sys.path.append(os.path.abspath(os.path.join(os.getcwd(), "..")))
import src.config as config

In [3]:
# from src.inference import get_feature_store
# from datetime import datetime, timedelta
# import pandas as pd  

# # Get the current datetime64[us, Etc/UTC]  
# current_date = pd.Timestamp.now(tz='Etc/UTC')
# feature_store = get_feature_store()

# # read time-series data from the feature store
# fetch_data_to = current_date - timedelta(hours=1)
# fetch_data_from = current_date - timedelta(days=1*29)
# print(f"Fetching data from {fetch_data_from} to {fetch_data_to}")
# feature_view = feature_store.get_feature_view(
#     name=config.FEATURE_VIEW_NAME, version=config.FEATURE_VIEW_VERSION
# )

# ts_data = feature_view.get_batch_data(
#     start_time=(fetch_data_from - timedelta(days=1)),
#     end_time=(fetch_data_to + timedelta(days=1)),
# )
# ts_data = ts_data[ts_data.pickup_hour.between(fetch_data_from, fetch_data_to)]
# ts_data.sort_values(["pickup_location_id", "pickup_hour"]).reset_index(drop=True)
# ts_data["pickup_hour"] = ts_data["pickup_hour"].dt.tz_localize(None)

# from src.data_utils import transform_ts_data_info_features
# features = transform_ts_data_info_features(ts_data, window_size=24*28, step_size=23)


from src.inference import get_feature_store
from src.config import FEATURE_VIEW_NAME, FEATURE_VIEW_VERSION
from datetime import datetime, timedelta
import pandas as pd

# Step 1: Get current UTC time
current_date = pd.Timestamp.now(tz="Etc/UTC")

# Step 2: Connect to feature store
feature_store = get_feature_store()

# Step 3: Define time range
fetch_data_to = current_date - timedelta(hours=1)
fetch_data_from = current_date - timedelta(days=29)
print(f"Fetching data from {fetch_data_from} to {fetch_data_to}")

# Step 4: Load feature view
feature_view = feature_store.get_feature_view(
    name=FEATURE_VIEW_NAME,
    version=FEATURE_VIEW_VERSION
)

# Step 5: Pull raw data with Spark enabled
ts_data = feature_view.get_batch_data(
    start_time=(fetch_data_from - timedelta(days=1)),
    end_time=(fetch_data_to + timedelta(days=1)),
    write_options={"use_spark": True}  # ✅ KEEPING THIS
)

# Step 6: Filter and sort
ts_data = ts_data[ts_data["pickup_hour"].between(fetch_data_from, fetch_data_to)]
ts_data = ts_data.sort_values(["pickup_location_id", "pickup_hour"]).reset_index(drop=True)

# Step 7: Remove timezone for model compatibility
ts_data["pickup_hour"] = ts_data["pickup_hour"].dt.tz_localize(None)

# Step 8: Feature transformation
from src.data_utils import transform_ts_data_into_features
features = transform_ts_data_into_features(ts_data, window_size=24 * 28, step_size=23)


2025-05-07 21:53:08,802 INFO: Initializing external client
2025-05-07 21:53:08,812 INFO: Base URL: https://c.app.hopsworks.ai:443
2025-05-07 21:53:10,161 INFO: Python Engine initialized.

Logged in to project, explore it here https://c.app.hopsworks.ai:443/p/1213633
Fetching data from 2025-04-09 01:53:08.802108+00:00 to 2025-05-08 00:53:08.802108+00:00
Finished: Reading data from Hopsworks, using Hopsworks Feature Query Service (4.23s) 


In [4]:
from src.inference import load_model_from_registry

model = load_model_from_registry()

2025-05-07 21:53:32,527 INFO: Closing external client and cleaning up certificates.
Connection closed.
2025-05-07 21:53:32,539 INFO: Initializing external client
2025-05-07 21:53:32,539 INFO: Base URL: https://c.app.hopsworks.ai:443
2025-05-07 21:53:33,401 INFO: Python Engine initialized.

Logged in to project, explore it here https://c.app.hopsworks.ai:443/p/1213633


Downloading: 0.000%|          | 0/320762 elapsed<00:00 remaining<?

Downloading model artifact (0 dirs, 1 files)... DONE

In [5]:
from src.inference import get_model_predictions
predictions = get_model_predictions(model, features)

In [6]:
predictions["pickup_hour"] = current_date.ceil('h')
predictions

Unnamed: 0,pickup_location_id,pickup_hour,predicted_demand
0,5187.03,2025-05-08 02:00:00+00:00,0.0
1,5282.02,2025-05-08 02:00:00+00:00,0.0
2,5746.14,2025-05-08 02:00:00+00:00,0.0
3,6098.12,2025-05-08 02:00:00+00:00,0.0
4,6322.01,2025-05-08 02:00:00+00:00,0.0
...,...,...,...
90,JC108,2025-05-08 02:00:00+00:00,0.0
91,JC109,2025-05-08 02:00:00+00:00,0.0
92,JC110,2025-05-08 02:00:00+00:00,0.0
93,JC115,2025-05-08 02:00:00+00:00,-0.0


In [7]:
from src.inference import get_feature_store

feature_group = get_feature_store().get_or_create_feature_group(
    name=config.FEATURE_GROUP_MODEL_PREDICTION,
    version=1,
    description="Predictions from LGBM Model",
    primary_key=["pickup_location_id", "pickup_hour"],
    event_time="pickup_hour",
)


2025-05-07 21:57:09,249 INFO: Closing external client and cleaning up certificates.
Connection closed.
2025-05-07 21:57:09,255 INFO: Initializing external client
2025-05-07 21:57:09,256 INFO: Base URL: https://c.app.hopsworks.ai:443
2025-05-07 21:57:10,136 INFO: Python Engine initialized.

Logged in to project, explore it here https://c.app.hopsworks.ai:443/p/1213633


In [8]:
feature_group.insert(predictions, write_options={"wait_for_job": False})

Feature Group created successfully, explore it at 
https://c.app.hopsworks.ai:443/p/1213633/fs/1201258/fg/1440567


Uploading Dataframe: 100.00% |███████████████████████████████| Rows 95/95 | Elapsed Time: 00:00 | Remaining Time: 00:00


Launching job: citi_bike_hourly_model_prediction_1_offline_fg_materialization
Job started successfully, you can follow the progress at 
https://c.app.hopsworks.ai:443/p/1213633/jobs/named/citi_bike_hourly_model_prediction_1_offline_fg_materialization/executions


(Job('citi_bike_hourly_model_prediction_1_offline_fg_materialization', 'SPARK'),
 None)