File size: 3,168 Bytes
2a8ebf2
d60ae27
57384dd
 
 
 
 
 
2a8ebf2
57384dd
 
 
 
 
 
 
 
 
 
2a8ebf2
 
 
 
 
57384dd
d60ae27
 
 
 
57384dd
d60ae27
 
 
 
 
 
 
 
 
57384dd
d60ae27
 
 
 
 
 
 
 
 
2a8ebf2
 
 
d60ae27
 
 
2a8ebf2
d60ae27
57384dd
d60ae27
 
2a8ebf2
 
57384dd
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
"""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