sbasu2512's picture
move feature layer to parquet
2a8ebf2
Raw
History Blame Contribute Delete
3.17 kB
"""Build and read model features from the Parquet feature layer."""
from __future__ import annotations
import json
import pandas as pd
from app.config import settings
from database.parquet_store import ParquetFeatureStore
from .correlations import add_correlation_features
from .us_market import build_us_wide
class FeatureEngine:
def __init__(self):
with settings.METADATA_PATH.open() as handle:
self.required_features = set(json.load(handle)["feature_cols"])
@staticmethod
def _read_store(store: ParquetFeatureStore, name: str) -> pd.DataFrame:
frame = store.read(name)
if "timestamp" in frame:
frame["timestamp"] = pd.to_datetime(frame["timestamp"])
return frame
def _merge(self, ohlcv_df, nifty_df, vix_df, macro_df, us_features_df) -> pd.DataFrame:
"""Attach only information available on or before each NSE session."""
us_wide = build_us_wide(us_features_df)
base = ohlcv_df.sort_values(["timestamp", "symbol"]).copy()
def attach(left, right, columns):
if not columns:
return left
right = right[["timestamp", *columns]].sort_values("timestamp")
return pd.merge_asof(left.sort_values("timestamp"), right, on="timestamp", direction="backward")
nifty_cols = [c for c in nifty_df if c != "timestamp" and f"nifty_{c}" in self.required_features]
nifty = nifty_df.rename(columns={c: f"nifty_{c}" for c in nifty_cols})
base = attach(base, nifty, [f"nifty_{c}" for c in nifty_cols])
vix_cols = [c for c in vix_df if c != "timestamp" and c in self.required_features]
if "close_1" in self.required_features and "close" in vix_df:
vix_df = vix_df.rename(columns={"close": "close_1"})
vix_cols.append("close_1")
base = attach(base, vix_df, list(dict.fromkeys(vix_cols)))
base = attach(base, macro_df, [c for c in macro_df if c != "timestamp" and c in self.required_features])
base = attach(base, us_wide, [c for c in us_wide if c != "timestamp" and c in self.required_features])
return add_correlation_features(base.sort_values(["symbol", "timestamp"]).reset_index(drop=True))
def build_all_features(self) -> pd.DataFrame:
store = ParquetFeatureStore(settings.DATABASE_URL)
frames = [self._read_store(store, name) for name in (
settings.OHLCV_FEATURES_TABLE, settings.NIFTY_FEATURES_TABLE,
settings.VIX_FEATURES_TABLE, settings.MACRO_FEATURES_TABLE,
settings.US_FEATURES_TABLE,
)]
return self._merge(*frames)
def build_features_for_date(self, prediction_date: pd.Timestamp) -> pd.DataFrame:
prediction_date = pd.Timestamp(prediction_date).normalize()
df = self._read_store(ParquetFeatureStore(settings.DATABASE_URL), "merged_features")
df = df[df["timestamp"].dt.normalize().eq(prediction_date)].copy()
missing = sorted(self.required_features - set(df.columns))
if missing:
raise RuntimeError(f"Feature-store merge is missing model columns: {missing}")
return df