Spaces:
Running
Running
| """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"]) | |
| 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 | |