Spaces:
Runtime error
Runtime error
| """Live V5 feature builder aligned with the cleaned notebook Part A.""" | |
| from pathlib import Path | |
| import numpy as np | |
| import pandas as pd | |
| import polars as pl | |
| EXTERNAL_MARKET_PATHS = [ | |
| Path("data/external_market_features.csv"), | |
| Path("data/external_market_features.parquet"), | |
| Path("data/whale_exchange_flow_features.csv"), | |
| ] | |
| def _parse_timestamp_column(frame: pl.DataFrame, timestamp_col: str = "timestamp") -> pl.DataFrame: | |
| dtype = frame.schema.get(timestamp_col) | |
| if dtype == pl.Datetime: | |
| return frame.with_columns(pl.col(timestamp_col).dt.cast_time_unit("us").alias(timestamp_col)) | |
| return frame.with_columns( | |
| pl.col(timestamp_col).cast(pl.Utf8).str.to_datetime(strict=False).alias(timestamp_col) | |
| ).drop_nulls(timestamp_col) | |
| def load_external_market_features(paths=EXTERNAL_MARKET_PATHS, timestamp_col: str = "timestamp") -> pl.DataFrame | None: | |
| frames = [] | |
| for raw_path in paths: | |
| path = Path(raw_path) | |
| if not path.exists(): | |
| continue | |
| if path.suffix.lower() == ".parquet": | |
| frame = pl.read_parquet(path) | |
| else: | |
| frame = pl.read_csv(path, try_parse_dates=False, infer_schema_length=10000) | |
| if timestamp_col not in frame.columns: | |
| continue | |
| frame = _parse_timestamp_column(frame, timestamp_col).sort(timestamp_col) | |
| numeric_cols = [ | |
| col for col, dtype in zip(frame.columns, frame.dtypes) | |
| if col != timestamp_col and dtype.is_numeric() | |
| ] | |
| if not numeric_cols: | |
| continue | |
| renamed = {col: f"ext_{col}" if not col.startswith("ext_") else col for col in numeric_cols} | |
| frames.append(frame.select([timestamp_col, *numeric_cols]).rename(renamed)) | |
| if not frames: | |
| return None | |
| out = frames[0] | |
| for frame in frames[1:]: | |
| out = out.join_asof(frame, on=timestamp_col, strategy="backward") | |
| return out.sort(timestamp_col) | |
| def add_external_feature_transforms(frame: pl.DataFrame, external_cols: list[str]) -> tuple[pl.DataFrame, list[str]]: | |
| engineered = [] | |
| for col in external_cols: | |
| for lag in [1, 4, 8, 24]: | |
| frame = frame.with_columns(pl.col(col).shift(lag).alias(f"{col}_lag_{lag}h")) | |
| engineered.append(f"{col}_lag_{lag}h") | |
| frame = frame.with_columns([ | |
| (pl.col(col) - pl.col(col).rolling_mean(24)).alias(f"{col}_dev_24h"), | |
| (pl.col(col) / (pl.col(col).rolling_mean(168) + 1e-10) - 1).alias(f"{col}_vs_168h"), | |
| pl.col(col).diff().rolling_mean(8).alias(f"{col}_flow_8h"), | |
| ]) | |
| engineered.extend([f"{col}_dev_24h", f"{col}_vs_168h", f"{col}_flow_8h"]) | |
| return frame, engineered | |
| def _to_hourly_polars(df_raw: pd.DataFrame) -> pl.DataFrame: | |
| df = df_raw.copy() | |
| if "timestamp" not in df.columns: | |
| if isinstance(df.index, pd.DatetimeIndex): | |
| df = df.reset_index().rename(columns={df.index.name or "index": "timestamp"}) | |
| else: | |
| raise ValueError("df_raw must have a timestamp column or DatetimeIndex") | |
| df["timestamp"] = pd.to_datetime(df["timestamp"]) | |
| for col, scale in { | |
| "tx_count": 1.5, | |
| "active_senders": 0.8, | |
| "active_receivers": 0.7, | |
| "total_eth_transferred": 1.2, | |
| "total_gas_used": 0.9, | |
| }.items(): | |
| if col not in df.columns: | |
| df[col] = df["volume"] * scale | |
| frame = pl.from_pandas(df).with_columns( | |
| pl.col("timestamp").dt.cast_time_unit("us").alias("timestamp") | |
| ).sort("timestamp") | |
| diffs = frame["timestamp"].diff().drop_nulls() | |
| is_hourly = len(diffs) > 0 and diffs.dt.total_minutes().median() >= 55 | |
| if is_hourly: | |
| return frame.select([ | |
| "timestamp", "open", "high", "low", "close", "volume", | |
| "tx_count", "active_senders", "active_receivers", | |
| "total_eth_transferred", "total_gas_used", | |
| ]) | |
| return frame.group_by_dynamic("timestamp", every="1h").agg([ | |
| pl.col("open").first().alias("open"), | |
| pl.col("high").max().alias("high"), | |
| pl.col("low").min().alias("low"), | |
| pl.col("close").last().alias("close"), | |
| pl.col("volume").sum().alias("volume"), | |
| pl.col("tx_count").first().alias("tx_count"), | |
| pl.col("active_senders").first().alias("active_senders"), | |
| pl.col("active_receivers").first().alias("active_receivers"), | |
| pl.col("total_eth_transferred").first().alias("total_eth_transferred"), | |
| pl.col("total_gas_used").first().alias("total_gas_used"), | |
| pl.col("close").count().alias("tick_count_1h"), | |
| ]).filter(pl.col("tick_count_1h") >= 30).drop("tick_count_1h").sort("timestamp") | |
| def build_live_v5_features(df_raw: pd.DataFrame) -> tuple[pd.DataFrame, list[str]]: | |
| df_h5 = _to_hourly_polars(df_raw) | |
| external_features = load_external_market_features() | |
| external_feature_cols: list[str] = [] | |
| if external_features is not None: | |
| df_h5 = df_h5.join_asof(external_features, on="timestamp", strategy="backward") | |
| external_feature_cols = [col for col in external_features.columns if col != "timestamp"] | |
| df_h5 = df_h5.with_columns([ | |
| pl.col(col).fill_null(strategy="forward").fill_null(0.0).alias(col) | |
| for col in external_feature_cols | |
| ]) | |
| df_h5, external_engineered_cols = add_external_feature_transforms(df_h5, external_feature_cols) | |
| df_h5 = df_h5.with_columns([ | |
| pl.col(col).fill_null(strategy="forward").fill_null(0.0).alias(col) | |
| for col in external_engineered_cols | |
| ]) | |
| external_feature_cols = external_feature_cols + external_engineered_cols | |
| c, h, l, v, o = pl.col("close"), pl.col("high"), pl.col("low"), pl.col("volume"), pl.col("open") | |
| for lag in [1, 2, 4, 6, 12, 24, 48, 168]: | |
| df_h5 = df_h5.with_columns((c / c.shift(lag) - 1).alias(f"return_{lag}h")) | |
| df_h5 = df_h5.with_columns([ | |
| ((c - o) / (h - l + 1e-10)).alias("candle_body_ratio"), | |
| ((h - l) / (c + 1e-10)).alias("range_pct"), | |
| ((c - l) / (h - l + 1e-10)).alias("close_location_value"), | |
| ((h - c.shift(1)) / (c.shift(1) + 1e-10)).alias("gap_high"), | |
| ((l - c.shift(1)) / (c.shift(1) + 1e-10)).alias("gap_low"), | |
| ]) | |
| for w in [4, 12, 24, 48, 168]: | |
| df_h5 = df_h5.with_columns([ | |
| c.rolling_mean(w).alias(f"sma_{w}h"), | |
| c.rolling_std(w).alias(f"vol_{w}h"), | |
| v.rolling_mean(w).alias(f"vol_avg_{w}h"), | |
| (h - l).rolling_mean(w).alias(f"range_avg_{w}h"), | |
| (c / c.rolling_mean(w) - 1).alias(f"price_vs_sma_{w}h"), | |
| ]) | |
| for period in [6, 14, 24]: | |
| delta = c.diff() | |
| gain = delta.clip(lower_bound=0).rolling_mean(period) | |
| loss = (-delta.clip(upper_bound=0)).rolling_mean(period) | |
| df_h5 = df_h5.with_columns((100 - 100 / (1 + gain / (loss + 1e-10))).alias(f"rsi_{period}h")) | |
| ema12 = c.ewm_mean(span=12) | |
| ema26 = c.ewm_mean(span=26) | |
| macd = ema12 - ema26 | |
| macd_signal = macd.ewm_mean(span=9) | |
| df_h5 = df_h5.with_columns([ | |
| macd.alias("macd_h"), | |
| macd_signal.alias("macd_signal_h"), | |
| (macd - macd_signal).alias("macd_hist_h"), | |
| ]) | |
| df_h5 = df_h5.with_columns([ | |
| (v / (v.rolling_mean(24) + 1e-10)).alias("vol_ratio_24h"), | |
| (v / (v.rolling_mean(168) + 1e-10)).alias("vol_ratio_168h"), | |
| v.rolling_std(24).alias("vol_volatility_24h"), | |
| (v * (c - c.shift(1)).sign()).rolling_sum(24).alias("obv_24h"), | |
| ]) | |
| bb_mid = c.rolling_mean(24) | |
| bb_std = c.rolling_std(24) | |
| df_h5 = df_h5.with_columns([ | |
| ((c - bb_mid) / (bb_std + 1e-10)).alias("bb_zscore_24h"), | |
| (bb_std / (bb_mid + 1e-10)).alias("bb_width_24h"), | |
| ]) | |
| for col in ["tx_count", "active_senders", "active_receivers", "total_eth_transferred", "total_gas_used"]: | |
| df_h5 = df_h5.with_columns([ | |
| (pl.col(col) / (pl.col(col).shift(24) + 1e-10) - 1).alias(f"{col}_change_24h"), | |
| pl.col(col).rolling_mean(24).alias(f"{col}_ma24h"), | |
| pl.col(col).rolling_mean(168).alias(f"{col}_ma168h"), | |
| ]) | |
| for col in ["tx_count", "active_senders", "total_eth_transferred"]: | |
| df_h5 = df_h5.with_columns( | |
| (pl.col(f"{col}_ma24h") / (pl.col(f"{col}_ma168h") + 1e-10) - 1).alias(f"{col}_momentum") | |
| ) | |
| df_h5 = df_h5.with_columns([ | |
| (2 * np.pi * pl.col("timestamp").dt.hour() / 24).sin().alias("hour_sin"), | |
| (2 * np.pi * pl.col("timestamp").dt.hour() / 24).cos().alias("hour_cos"), | |
| (2 * np.pi * pl.col("timestamp").dt.weekday() / 7).sin().alias("dow_sin"), | |
| (2 * np.pi * pl.col("timestamp").dt.weekday() / 7).cos().alias("dow_cos"), | |
| ]) | |
| for lag in [1, 2, 3, 4, 6, 12]: | |
| df_h5 = df_h5.with_columns([ | |
| (c / c.shift(1) - 1).shift(lag).alias(f"ret_lag_{lag}"), | |
| pl.col("range_pct").shift(lag).alias(f"range_lag_{lag}"), | |
| ]) | |
| hr_ret = c / c.shift(1) - 1 | |
| vol_short = hr_ret.rolling_std(window_size=24) | |
| vol_long = hr_ret.rolling_std(window_size=720) | |
| df_h5 = df_h5.with_columns([ | |
| (vol_short / (vol_long + 1e-10)).alias("vol_regime_ratio"), | |
| vol_long.alias("vol_30d"), | |
| (vol_short - vol_long).alias("vol_shift"), | |
| ]) | |
| obv = (v * (c - c.shift(1)).sign()).cum_sum() | |
| obv_slope_24 = obv - obv.shift(24) | |
| price_slope_24 = c - c.shift(24) | |
| df_h5 = df_h5.with_columns([ | |
| obv_slope_24.alias("obv_slope_24h"), | |
| (obv_slope_24.sign() - price_slope_24.sign()).alias("vol_price_divergence"), | |
| ]) | |
| avg_xfer = pl.col("total_eth_transferred") / (pl.col("active_senders") + 1e-10) | |
| df_h5 = df_h5.with_columns([ | |
| avg_xfer.alias("avg_transfer_size"), | |
| (avg_xfer / (avg_xfer.shift(24) + 1e-10) - 1).alias("transfer_size_change_24h"), | |
| (avg_xfer.rolling_mean(24) / (avg_xfer.rolling_mean(168) + 1e-10) - 1).alias("transfer_size_momentum"), | |
| ]) | |
| gpv = pl.col("total_gas_used") / (v + 1e-10) | |
| df_h5 = df_h5.with_columns([ | |
| gpv.alias("gas_per_volume"), | |
| gpv.rolling_mean(24).alias("gas_per_volume_ma24h"), | |
| (gpv / (gpv.shift(24) + 1e-10) - 1).alias("gas_per_volume_change_24h"), | |
| ]) | |
| ret_4h = c / c.shift(4) - 1 | |
| ret_48h = c / c.shift(48) - 1 | |
| ret_168h = c / c.shift(168) - 1 | |
| df_h5 = df_h5.with_columns([ | |
| (ret_4h - ret_48h).alias("momentum_divergence_4_48"), | |
| (ret_48h - ret_168h).alias("momentum_divergence_48_168"), | |
| (ret_4h.sign() - ret_168h.sign()).alias("trend_alignment"), | |
| vol_short.rolling_std(window_size=48).alias("vol_of_vol_48h"), | |
| ]) | |
| for period in [6, 14, 24]: | |
| df_h5 = df_h5.with_columns( | |
| (pl.col(f"rsi_{period}h") * pl.col("vol_ratio_24h")).alias(f"rsi{period}_x_vol") | |
| ) | |
| df_h5 = df_h5.with_columns([ | |
| (pl.col("active_senders") / (pl.col("active_senders").shift(24) + 1e-10) - 1).alias("senders_mom_24h"), | |
| (pl.col("active_senders").rolling_mean(24) / (pl.col("active_senders").rolling_mean(72) + 1e-10) - 1) | |
| .alias("senders_mom_24_vs_72"), | |
| pl.col("active_senders").diff().rolling_mean(24).alias("senders_accel_24h"), | |
| ]) | |
| df_h5 = df_h5.with_columns([ | |
| vol_short.alias("vol_24h_filter"), | |
| (c / c.shift(720) - 1).alias("ret_30d_filter"), | |
| ]) | |
| vol_arr = df_h5["vol_24h_filter"].to_numpy() | |
| ratio_arr = (df_h5["vol_24h_filter"] / (df_h5["vol_30d"] + 1e-10)).to_numpy() | |
| finite_vol = vol_arr[np.isfinite(vol_arr)] | |
| vol_p25 = np.nanpercentile(finite_vol, 25) if len(finite_vol) else np.nan | |
| vol_p75 = np.nanpercentile(finite_vol, 75) if len(finite_vol) else np.nan | |
| df_h5 = df_h5.with_columns([ | |
| pl.Series("vol_regime_low", (vol_arr < vol_p25).astype(np.float64)), | |
| pl.Series("vol_regime_high", (vol_arr > vol_p75).astype(np.float64)), | |
| pl.Series("vol_regime_expanding", (ratio_arr > 1.2).astype(np.float64)), | |
| pl.Series("vol_regime_compressing", (ratio_arr < 0.8).astype(np.float64)), | |
| ]) | |
| exclude = { | |
| "timestamp", "open", "high", "low", "close", "volume", | |
| "tx_count", "active_senders", "active_receivers", | |
| "total_eth_transferred", "total_gas_used", | |
| "vol_24h_filter", "ret_30d_filter", | |
| } | |
| feature_cols = [col for col in df_h5.columns if col not in exclude and not col.startswith("target_ret_")] | |
| df_feat = df_h5.drop_nulls().to_pandas() | |
| return df_feat, feature_cols | |