{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"pygments_lexer":"ipython3","nbconvert_exporter":"python","version":"3.6.4","file_extension":".py","codemirror_mode":{"name":"ipython","version":3},"name":"python","mimetype":"text/x-python"},"kaggle":{"accelerator":"none","dataSources":[{"sourceId":96164,"databundleVersionId":12993472,"sourceType":"competition"}],"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"# This Python 3 environment comes with many helpful analytics libraries installed\n# It is defined by the kaggle/python Docker image: https://github.com/kaggle/docker-python\n# For example, here's several helpful packages to load\n\nimport numpy as np # linear algebra\nimport pandas as pd # data processing, CSV file I/O (e.g. pd.read_csv)\n\n# Input data files are available in the read-only \"../input/\" directory\n# For example, running this (by clicking run or pressing Shift+Enter) will list all files under the input directory\n\nimport os\nfor dirname, _, filenames in os.walk('/kaggle/input'):\n    for filename in filenames:\n        print(os.path.join(dirname, filename))\n\n# You can write up to 20GB to the current directory (/kaggle/working/) that gets preserved as output when you create a version using \"Save & Run All\" \n# You can also write temporary files to /kaggle/temp/, but they won't be saved outside of the current session","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ============================================\n# INSTALLATIONS\n# ============================================\n!pip install -q prophet\n!pip install -q koolbox\n!pip install -q scikit-learn==1.5.2\n!pip install -q autogluon\n!pip install -q flaml[automl]\n!pip install -q mljar-supervised\n!pip install -q h2o\n!pip install -q optuna\n!pip install -q lightgbm\n!pip install -q xgboost\n!pip install -q catboost\n!pip install -q ngboost\n!pip install -q shap\n!pip install -q tsfresh\n\n# ============================================\n# IMPORTS\n# ============================================\nimport numpy as np\nimport pandas as pd\nimport matplotlib.pyplot as plt\nimport seaborn as sns\nfrom prophet import Prophet\nimport warnings\nwarnings.filterwarnings('ignore')\n\nfrom sklearn.ensemble import HistGradientBoostingRegressor, RandomForestRegressor, ExtraTreesRegressor\nfrom sklearn.model_selection import KFold, cross_val_score, TimeSeriesSplit\nfrom sklearn.linear_model import Ridge, Lasso, ElasticNet, RidgeCV, LassoCV\nfrom sklearn.preprocessing import StandardScaler, RobustScaler, QuantileTransformer\nfrom sklearn.decomposition import PCA, FastICA, TruncatedSVD\nfrom sklearn.feature_selection import SelectFromModel, mutual_info_regression\nfrom sklearn.base import clone, BaseEstimator, RegressorMixin, TransformerMixin\nfrom sklearn.pipeline import Pipeline\nfrom sklearn.cluster import KMeans\nfrom lightgbm import LGBMRegressor, LGBMClassifier\nfrom xgboost import XGBRegressor\nfrom catboost import CatBoostRegressor\nfrom scipy.stats import pearsonr, spearmanr\nfrom scipy.signal import savgol_filter\nfrom ngboost import NGBRegressor\nfrom ngboost.distns import Normal, LogNormal\nimport joblib\nimport gc\nimport os\nfrom typing import Dict, List, Tuple, Optional\nimport shap\nfrom tsfresh import extract_features\nfrom tsfresh.utilities.dataframe_functions import impute\n\n# ============================================\n# CONFIGURATION\n# ============================================\nclass CFG:\n    train_path = \"/kaggle/input/drw-crypto-market-prediction/train.parquet\"\n    test_path = \"/kaggle/input/drw-crypto-market-prediction/test.parquet\"\n    sample_sub_path = \"/kaggle/input/drw-crypto-market-prediction/sample_submission.csv\"\n    \n    target = \"label\"\n    n_folds = 5\n    seed = 42\n    \n    # Feature list\n    X_FEATURES = ['X363', 'X321', 'X405', 'X730', 'X523', 'X756', 'X589', 'X462', 'X779',\n                  'X25', 'X532', 'X520', 'X329', 'X383', 'X751', 'X535', 'X639', 'X596', 'X761',\n                  \"X752\", \"X287\", \"X298\", \"X759\", \"X302\", \"X55\", \"X56\", \"X52\", \"X303\", \"X51\",\n                  \"X598\", \"X385\", \"X603\", \"X674\", \"X415\", \"X345\", \"X174\", \"X178\", \"X168\", \"X612\",\n                  \"bid_qty\", \"ask_qty\", \"buy_qty\", \"sell_qty\"]\n    \n    # Advanced configuration\n    use_recent_months = 6  # Only use recent months of data\n    noise_percentile = 80  # Percentile for noise detection\n    min_feature_importance = 0.001  # Minimum feature importance threshold\n\n# ============================================\n# UTILITY FUNCTIONS\n# ============================================\ndef _pearsonr(y_true, y_pred):\n    return pearsonr(y_true, y_pred)[0]\n\ndef reduce_mem_usage(dataframe, dataset):    \n    print(f'Reducing memory usage for: {dataset}')\n    initial_mem_usage = dataframe.memory_usage().sum() / 1024**2\n    \n    for col in dataframe.columns:\n        col_type = dataframe[col].dtype\n        c_min = dataframe[col].min()\n        c_max = dataframe[col].max()\n        \n        if str(col_type)[:3] == 'int':\n            if c_min > np.iinfo(np.int8).min and c_max < np.iinfo(np.int8).max:\n                dataframe[col] = dataframe[col].astype(np.int8)\n            elif c_min > np.iinfo(np.int16).min and c_max < np.iinfo(np.int16).max:\n                dataframe[col] = dataframe[col].astype(np.int16)\n            elif c_min > np.iinfo(np.int32).min and c_max < np.iinfo(np.int32).max:\n                dataframe[col] = dataframe[col].astype(np.int32)\n            elif c_min > np.iinfo(np.int64).min and c_max < np.iinfo(np.int64).max:\n                dataframe[col] = dataframe[col].astype(np.int64)\n        else:\n            if c_min > np.finfo(np.float16).min and c_max < np.finfo(np.float16).max:\n                dataframe[col] = dataframe[col].astype(np.float16)\n            elif c_min > np.finfo(np.float32).min and c_max < np.finfo(np.float32).max:\n                dataframe[col] = dataframe[col].astype(np.float32)\n            else:\n                dataframe[col] = dataframe[col].astype(np.float64)\n\n    final_mem_usage = dataframe.memory_usage().sum() / 1024**2\n    print(f'--- Memory usage before: {initial_mem_usage:.2f} MB')\n    print(f'--- Memory usage after: {final_mem_usage:.2f} MB')\n    print(f'--- Decreased memory usage by {100 * (initial_mem_usage - final_mem_usage) / initial_mem_usage:.1f}%\\n')\n\n    return dataframe\n\n# ============================================\n# ADVANCED FEATURE ENGINEERING\n# ============================================\ndef feature_engineering(df):\n    # Original features\n    df['bid_ask_interaction'] = df['bid_qty'] * df['ask_qty']\n    df['bid_buy_interaction'] = df['bid_qty'] * df['buy_qty']\n    df['bid_sell_interaction'] = df['bid_qty'] * df['sell_qty']\n    df['ask_buy_interaction'] = df['ask_qty'] * df['buy_qty']\n    df['ask_sell_interaction'] = df['ask_qty'] * df['sell_qty']\n    \n    df['volume_weighted_sell'] = df['sell_qty'] * df['volume']\n    df['buy_sell_ratio'] = df['buy_qty'] / (df['sell_qty'] + 1e-10)\n    df['selling_pressure'] = df['sell_qty'] / (df['volume'] + 1e-10)\n    df['log_volume'] = np.log1p(df['volume'])\n    \n    df['effective_spread_proxy'] = np.abs(df['buy_qty'] - df['sell_qty']) / (df['volume'] + 1e-10)\n    df['bid_ask_imbalance'] = (df['bid_qty'] - df['ask_qty']) / (df['bid_qty'] + df['ask_qty'] + 1e-10)\n    df['order_flow_imbalance'] = (df['buy_qty'] - df['sell_qty']) / (df['buy_qty'] + df['sell_qty'] + 1e-10)\n    df['liquidity_ratio'] = (df['bid_qty'] + df['ask_qty']) / (df['volume'] + 1e-10)\n    \n    # New microstructure features\n    df['net_order_flow'] = df['buy_qty'] - df['sell_qty']\n    df['normalized_net_flow'] = df['net_order_flow'] / (df['volume'] + 1e-10)\n    df['buying_pressure'] = df['buy_qty'] / (df['volume'] + 1e-10)\n    df['volume_weighted_buy'] = df['buy_qty'] * df['volume']\n    \n    # Liquidity Depth Measures\n    df['total_depth'] = df['bid_qty'] + df['ask_qty']\n    df['depth_imbalance'] = (df['bid_qty'] - df['ask_qty']) / (df['total_depth'] + 1e-10)\n    df['relative_spread'] = np.abs(df['bid_qty'] - df['ask_qty']) / (df['total_depth'] + 1e-10)\n    df['log_depth'] = np.log1p(df['total_depth'])\n    \n    # Order Flow Toxicity Proxies\n    df['kyle_lambda'] = np.abs(df['net_order_flow']) / (df['volume'] + 1e-10)\n    df['flow_toxicity'] = np.abs(df['order_flow_imbalance']) * df['volume']\n    df['aggressive_flow_ratio'] = (df['buy_qty'] + df['sell_qty']) / (df['total_depth'] + 1e-10)\n    \n    # Market Activity Indicators\n    df['volume_depth_ratio'] = df['volume'] / (df['total_depth'] + 1e-10)\n    df['activity_intensity'] = (df['buy_qty'] + df['sell_qty']) / (df['volume'] + 1e-10)\n    df['log_buy_qty'] = np.log1p(df['buy_qty'])\n    df['log_sell_qty'] = np.log1p(df['sell_qty'])\n    df['log_bid_qty'] = np.log1p(df['bid_qty'])\n    df['log_ask_qty'] = np.log1p(df['ask_qty'])\n    \n    # Handle infinities and NaN\n    df = df.replace([np.inf, -np.inf], np.nan)\n    \n    # Fill NaN with median\n    for col in df.columns:\n        if df[col].isna().any():\n            median_val = df[col].median()\n            df[col] = df[col].fillna(median_val if not pd.isna(median_val) else 0)\n    \n    return df\n\ndef add_lag_features(df, features, lags=[1, 5, 15, 30, 60]):\n    \"\"\"Add lag features for time series\"\"\"\n    new_features = []\n    for feature in features:\n        if feature in df.columns:\n            for lag in lags:\n                lag_col = f'{feature}_lag_{lag}'\n                df[lag_col] = df[feature].shift(lag)\n                new_features.append(lag_col)\n    return df, new_features\n\ndef add_rolling_features(df, features, windows=[5, 15, 30, 60]):\n    \"\"\"Add rolling statistics\"\"\"\n    new_features = []\n    for feature in features:\n        if feature in df.columns:\n            for window in windows:\n                # Rolling mean\n                mean_col = f'{feature}_rolling_mean_{window}'\n                df[mean_col] = df[feature].rolling(window, min_periods=1).mean()\n                new_features.append(mean_col)\n                \n                # Rolling std\n                std_col = f'{feature}_rolling_std_{window}'\n                df[std_col] = df[feature].rolling(window, min_periods=1).std()\n                new_features.append(std_col)\n                \n                # Rolling min/max spread\n                spread_col = f'{feature}_rolling_spread_{window}'\n                rolling_max = df[feature].rolling(window, min_periods=1).max()\n                rolling_min = df[feature].rolling(window, min_periods=1).min()\n                df[spread_col] = rolling_max - rolling_min\n                new_features.append(spread_col)\n    \n    return df, new_features\n\n# ============================================\n# ADVANCED NOISE REDUCTION & FEATURE COMPRESSION\n# ============================================\nclass HierarchicalNoiseReducer(BaseEstimator, TransformerMixin):\n    \"\"\"Hierarchical noise reduction with multiple stages\"\"\"\n    \n    def __init__(self, n_components=50, noise_threshold=0.1, n_clusters=5):\n        self.n_components = n_components\n        self.noise_threshold = noise_threshold\n        self.n_clusters = n_clusters\n        self.scalers = {}\n        self.reducers = {}\n        self.noise_filters = {}\n        self.feature_clusters = None\n        \n    def fit(self, X, y=None):\n        X_array = X.values if hasattr(X, 'values') else X\n        \n        # Stage 1: Identify feature clusters\n        self.kmeans = KMeans(n_clusters=self.n_clusters, random_state=42)\n        feature_corr = np.corrcoef(X_array.T)\n        self.feature_clusters = self.kmeans.fit_predict(feature_corr)\n        \n        # Stage 2: Apply cluster-specific noise reduction\n        for cluster_id in range(self.n_clusters):\n            cluster_mask = self.feature_clusters == cluster_id\n            X_cluster = X_array[:, cluster_mask]\n            \n            # Robust scaling per cluster\n            scaler = RobustScaler()\n            X_scaled = scaler.fit_transform(X_cluster)\n            self.scalers[cluster_id] = scaler\n            \n            # Noise estimation using rolling variance\n            noise_levels = []\n            for col in range(X_scaled.shape[1]):\n                # Use multiple methods to estimate noise\n                diff_noise = np.std(np.diff(X_scaled[:, col])) / np.sqrt(2)\n                mad_noise = np.median(np.abs(X_scaled[:, col] - np.median(X_scaled[:, col]))) * 1.4826\n                combined_noise = (diff_noise + mad_noise) / 2\n                noise_levels.append(combined_noise)\n            \n            noise_levels = np.array(noise_levels)\n            noise_threshold = np.percentile(noise_levels, CFG.noise_percentile)\n            self.noise_filters[cluster_id] = noise_levels < noise_threshold\n            \n            # Apply dimensionality reduction to clean features\n            clean_features = X_scaled[:, self.noise_filters[cluster_id]]\n            if clean_features.shape[1] > 0:\n                n_comp = min(self.n_components // self.n_clusters, clean_features.shape[1])\n                reducer = PCA(n_components=n_comp, random_state=42)\n                reducer.fit(clean_features)\n                self.reducers[cluster_id] = reducer\n            \n        return self\n    \n    def transform(self, X):\n        X_array = X.values if hasattr(X, 'values') else X\n        transformed_features = []\n        \n        for cluster_id in range(self.n_clusters):\n            cluster_mask = self.feature_clusters == cluster_id\n            X_cluster = X_array[:, cluster_mask]\n            \n            # Scale\n            X_scaled = self.scalers[cluster_id].transform(X_cluster)\n            \n            # Apply noise filter and reduce\n            if cluster_id in self.reducers:\n                clean_features = X_scaled[:, self.noise_filters[cluster_id]]\n                reduced_features = self.reducers[cluster_id].transform(clean_features)\n                transformed_features.append(reduced_features)\n            \n            # Keep some noisy features with smoothing\n            noisy_features = X_scaled[:, ~self.noise_filters[cluster_id]]\n            if noisy_features.shape[1] > 0:\n                # Apply Savitzky-Golay filter for smoothing\n                smoothed = np.zeros_like(noisy_features)\n                for col in range(noisy_features.shape[1]):\n                    try:\n                        smoothed[:, col] = savgol_filter(noisy_features[:, col], \n                                                        window_length=min(11, noisy_features.shape[0]//2*2-1), \n                                                        polyorder=3)\n                    except:\n                        smoothed[:, col] = noisy_features[:, col]\n                \n                # Take top principal components of smoothed noisy features\n                if smoothed.shape[1] > 5:\n                    reducer = TruncatedSVD(n_components=5, random_state=42)\n                    smoothed = reducer.fit_transform(smoothed)\n                transformed_features.append(smoothed)\n        \n        return np.hstack(transformed_features)\n\nclass AdaptiveFeatureSelector(BaseEstimator, TransformerMixin):\n    \"\"\"Adaptive feature selection based on stability and importance\"\"\"\n    \n    def __init__(self, n_features=100, stability_threshold=0.7):\n        self.n_features = n_features\n        self.stability_threshold = stability_threshold\n        self.selected_features = None\n        self.feature_scores = None\n        \n    def fit(self, X, y):\n        X_array = X.values if hasattr(X, 'values') else X\n        n_samples, n_features = X_array.shape\n        \n        # Calculate multiple importance metrics\n        # 1. Mutual information\n        mi_scores = mutual_info_regression(X_array, y, random_state=42)\n        \n        # 2. Correlation with target\n        corr_scores = np.abs([pearsonr(X_array[:, i], y)[0] for i in range(n_features)])\n        \n        # 3. Feature stability (using bootstrap)\n        stability_scores = np.zeros(n_features)\n        n_bootstrap = 10\n        \n        for _ in range(n_bootstrap):\n            idx = np.random.choice(n_samples, size=n_samples, replace=True)\n            X_boot = X_array[idx]\n            y_boot = y[idx]\n            \n            # Use simple model for speed\n            model = LGBMRegressor(n_estimators=50, num_leaves=31, random_state=42, verbose=-1)\n            model.fit(X_boot, y_boot)\n            \n            importances = model.feature_importances_\n            top_features = np.argsort(importances)[-self.n_features:]\n            stability_scores[top_features] += 1\n        \n        stability_scores /= n_bootstrap\n        \n        # Combine scores\n        mi_scores_norm = mi_scores / (mi_scores.max() + 1e-10)\n        corr_scores_norm = corr_scores / (corr_scores.max() + 1e-10)\n        \n        self.feature_scores = (mi_scores_norm + corr_scores_norm + stability_scores) / 3\n        \n        # Select features above stability threshold and top n_features\n        stable_mask = stability_scores >= self.stability_threshold\n        stable_indices = np.where(stable_mask)[0]\n        \n        if len(stable_indices) < self.n_features:\n            # Add more features based on combined score\n            remaining_indices = np.where(~stable_mask)[0]\n            remaining_scores = self.feature_scores[remaining_indices]\n            additional_indices = remaining_indices[np.argsort(remaining_scores)[-( self.n_features - len(stable_indices)):]]\n            self.selected_features = np.concatenate([stable_indices, additional_indices])\n        else:\n            # Select top n_features from stable features\n            stable_scores = self.feature_scores[stable_indices]\n            top_stable_indices = stable_indices[np.argsort(stable_scores)[-self.n_features:]]\n            self.selected_features = top_stable_indices\n        \n        return self\n    \n    def transform(self, X):\n        X_array = X.values if hasattr(X, 'values') else X\n        return X_array[:, self.selected_features]\n\n# ============================================\n# NGBOOST WITH UNCERTAINTY\n# ============================================\nclass UncertaintyAwareNGBoost(BaseEstimator, RegressorMixin):\n    \"\"\"NGBoost with uncertainty-based sample weighting\"\"\"\n    \n    def __init__(self, n_estimators=500, learning_rate=0.01, minibatch_frac=0.5, \n                 uncertainty_weight=True):\n        self.n_estimators = n_estimators\n        self.learning_rate = learning_rate\n        self.minibatch_frac = minibatch_frac\n        self.uncertainty_weight = uncertainty_weight\n        self.model = None\n        self.uncertainty_threshold = None\n        \n    def fit(self, X, y, X_val=None, y_val=None):\n        # Initial fit to get uncertainty estimates\n        self.model = NGBRegressor(\n            Dist=Normal,\n            n_estimators=100,  # Quick initial fit\n            learning_rate=self.learning_rate * 2,\n            minibatch_frac=self.minibatch_frac,\n            verbose=False,\n            random_state=42\n        )\n        \n        self.model.fit(X, y)\n        \n        if self.uncertainty_weight:\n            # Get uncertainty estimates\n            y_dist = self.model.pred_dist(X)\n            uncertainties = y_dist.std()\n            \n            # Create sample weights (lower weight for high uncertainty)\n            self.uncertainty_threshold = np.percentile(uncertainties, 75)\n            sample_weights = 1.0 / (1.0 + uncertainties / self.uncertainty_threshold)\n            sample_weights = sample_weights / sample_weights.mean()  # Normalize\n        else:\n            sample_weights = None\n        \n        # Final fit with all estimators\n        self.model = NGBRegressor(\n            Dist=Normal,\n            n_estimators=self.n_estimators,\n            learning_rate=self.learning_rate,\n            minibatch_frac=self.minibatch_frac,\n            verbose=False,\n            random_state=42\n        )\n        \n        if X_val is not None and y_val is not None:\n            self.model.fit(X, y, X_val=X_val, Y_val=y_val, sample_weight=sample_weights)\n        else:\n            self.model.fit(X, y, sample_weight=sample_weights)\n        \n        return self\n    \n    def predict(self, X):\n        return self.model.predict(X)\n    \n    def predict_with_uncertainty(self, X):\n        y_dist = self.model.pred_dist(X)\n        return y_dist.mean(), y_dist.std()\n\n# ============================================\n# HIERARCHICAL ENSEMBLE\n# ============================================\nclass HierarchicalEnsemble(BaseEstimator, RegressorMixin):\n    \"\"\"Multi-level ensemble with uncertainty weighting\"\"\"\n    \n    def __init__(self, base_models, meta_model=None, use_uncertainty=True):\n        self.base_models = base_models\n        self.meta_model = meta_model or RidgeCV(alphas=[0.001, 0.01, 0.1, 1.0, 10.0])\n        self.use_uncertainty = use_uncertainty\n        self.level1_models = {}\n        self.level2_model = None\n        \n    def fit(self, X, y, cv=None):\n        if cv is None:\n            cv = TimeSeriesSplit(n_splits=5)\n        \n        # Level 1: Train base models with cross-validation\n        level1_predictions = {}\n        level1_uncertainties = {}\n        \n        for name, model in self.base_models.items():\n            print(f\"Training {name}...\")\n            oof_preds = np.zeros(len(X))\n            oof_uncertainty = np.ones(len(X))\n            \n            for fold, (train_idx, val_idx) in enumerate(cv.split(X, y)):\n                X_train, X_val = X.iloc[train_idx], X.iloc[val_idx]\n                y_train, y_val = y.iloc[train_idx], y.iloc[val_idx]\n                \n                # Clone and train model\n                model_clone = clone(model)\n                \n                if hasattr(model_clone, 'predict_with_uncertainty'):\n                    model_clone.fit(X_train, y_train, X_val, y_val)\n                    preds, uncertainty = model_clone.predict_with_uncertainty(X_val)\n                    oof_preds[val_idx] = preds\n                    oof_uncertainty[val_idx] = uncertainty\n                else:\n                    model_clone.fit(X_train, y_train)\n                    oof_preds[val_idx] = model_clone.predict(X_val)\n                \n                self.level1_models[f\"{name}_fold{fold}\"] = model_clone\n            \n            level1_predictions[name] = oof_preds\n            level1_uncertainties[name] = oof_uncertainty\n            \n            score = _pearsonr(y, oof_preds)\n            print(f\"{name} OOF Score: {score:.6f}\")\n        \n        # Level 2: Train meta-model\n        X_meta = pd.DataFrame(level1_predictions)\n        \n        if self.use_uncertainty:\n            # Add uncertainty-weighted features\n            for name in level1_predictions:\n                uncertainty = level1_uncertainties[name]\n                weight = 1.0 / (1.0 + uncertainty / uncertainty.mean())\n                X_meta[f\"{name}_weighted\"] = level1_predictions[name] * weight\n        \n        self.level2_model = clone(self.meta_model)\n        self.level2_model.fit(X_meta, y)\n        \n        # Print meta-model weights\n        if hasattr(self.level2_model, 'coef_'):\n            print(\"\\nMeta-model weights:\")\n            for feature, weight in zip(X_meta.columns, self.level2_model.coef_):\n                print(f\"{feature}: {weight:.4f}\")\n        \n        return self\n    \n    def predict(self, X):\n        level1_predictions = {}\n        level1_uncertainties = {}\n        \n        for name, model in self.base_models.items():\n            preds = []\n            uncertainties = []\n            \n            # Average predictions from all folds\n            for fold in range(5):  # Assuming 5 folds\n                model_key = f\"{name}_fold{fold}\"\n                if model_key in self.level1_models:\n                    fold_model = self.level1_models[model_key]\n                    \n                    if hasattr(fold_model, 'predict_with_uncertainty'):\n                        pred, unc = fold_model.predict_with_uncertainty(X)\n                        preds.append(pred)\n                        uncertainties.append(unc)\n                    else:\n                        preds.append(fold_model.predict(X))\n                        uncertainties.append(np.ones(len(X)))\n            \n            level1_predictions[name] = np.mean(preds, axis=0)\n            level1_uncertainties[name] = np.mean(uncertainties, axis=0)\n        \n        X_meta = pd.DataFrame(level1_predictions)\n        \n        if self.use_uncertainty:\n            for name in level1_predictions:\n                uncertainty = level1_uncertainties[name]\n                weight = 1.0 / (1.0 + uncertainty / uncertainty.mean())\n                X_meta[f\"{name}_weighted\"] = level1_predictions[name] * weight\n        \n        return self.level2_model.predict(X_meta)\n\n# ============================================\n# MAIN PIPELINE\n# ============================================\ndef main():\n    print(\"Starting Advanced DRW Crypto Market Prediction Pipeline...\")\n    print(\"Focus: Minimizing train-CV gap through noise reduction and hierarchical learning\")\n    \n    # Load data\n    print(\"\\n1. Loading data...\")\n    train = pd.read_parquet(CFG.train_path).reset_index(drop=True)\n    test = pd.read_parquet(CFG.test_path).reset_index(drop=True)\n    \n    # Use only recent data as suggested\n    if CFG.use_recent_months > 0:\n        total_minutes = CFG.use_recent_months * 30 * 24 * 60  # Approximate\n        train = train.tail(total_minutes).reset_index(drop=True)\n        print(f\"Using only last {CFG.use_recent_months} months of data: {len(train)} samples\")\n    \n    # Select features\n    selected_columns = CFG.X_FEATURES + [\"volume\"]\n    train = train[selected_columns + [CFG.target]]\n    test = test[selected_columns]\n    \n    # Add timestamp if missing\n    if '__index_level_0__' not in train.columns:\n        train['__index_level_0__'] = pd.date_range('2023-03-01', periods=len(train), freq='T')\n    if '__index_level_0__' not in test.columns:\n        test['__index_level_0__'] = pd.date_range('2024-03-01', periods=len(test), freq='T')\n    \n    # Apply feature engineering\n    print(\"\\n2. Advanced Feature Engineering...\")\n    train = feature_engineering(train)\n    test = feature_engineering(test)\n    \n    # Add time series features\n    important_features = ['volume', 'order_flow_imbalance', 'bid_ask_imbalance', 'liquidity_ratio']\n    train, lag_features = add_lag_features(train, important_features)\n    test, _ = add_lag_features(test, important_features)\n    \n    train, rolling_features = add_rolling_features(train, important_features)\n    test, _ = add_rolling_features(test, important_features)\n    \n    # Fill NaN from lag/rolling features\n    train = train.fillna(method='ffill').fillna(0)\n    test = test.fillna(method='ffill').fillna(0)\n    \n    # Remove base features and timestamp\n    to_remove = [\"bid_qty\", \"ask_qty\", \"buy_qty\", \"sell_qty\", \"volume\", \"__index_level_0__\"]\n    train = train.drop(columns=[col for col in to_remove if col in train.columns])\n    test = test.drop(columns=[col for col in to_remove if col in test.columns])\n    \n    # Reduce memory\n    train = reduce_mem_usage(train, \"train\")\n    test = reduce_mem_usage(test, \"test\")\n    \n    # Prepare data\n    X = train.drop(CFG.target, axis=1)\n    y = train[CFG.target]\n    X_test = test\n    \n    print(f\"Shape after feature engineering: {X.shape}\")\n    \n    # ============================================\n    # HIERARCHICAL NOISE REDUCTION\n    # ============================================\n    print(\"\\n3. Applying Hierarchical Noise Reduction...\")\n    noise_reducer = HierarchicalNoiseReducer(n_components=100, noise_threshold=0.15, n_clusters=8)\n    X_reduced = pd.DataFrame(noise_reducer.fit_transform(X))\n    X_test_reduced = pd.DataFrame(noise_reducer.transform(X_test))\n    \n    print(f\"Shape after noise reduction: {X_reduced.shape}\")\n    \n    # ============================================\n    # ADAPTIVE FEATURE SELECTION\n    # ============================================\n    print(\"\\n4. Adaptive Feature Selection...\")\n    feature_selector = AdaptiveFeatureSelector(n_features=150, stability_threshold=0.6)\n    X_selected = pd.DataFrame(feature_selector.fit_transform(X_reduced, y))\n    X_test_selected = pd.DataFrame(feature_selector.transform(X_test_reduced))\n    \n    print(f\"Shape after feature selection: {X_selected.shape}\")\n    \n    # ============================================\n    # PREPARE MODELS\n    # ============================================\n    print(\"\\n5. Preparing Advanced Models...\")\n    \n    # Base models with focus on regularization\n    base_models = {\n        'NGBoost': UncertaintyAwareNGBoost(\n            n_estimators=500,\n            learning_rate=0.01,\n            minibatch_frac=0.5,\n            uncertainty_weight=True\n        ),\n        \n        'LightGBM_Regularized': LGBMRegressor(\n            n_estimators=300,\n            learning_rate=0.02,\n            num_leaves=31,\n            subsample=0.6,\n            colsample_bytree=0.6,\n            reg_alpha=50,\n            reg_lambda=50,\n            min_child_samples=100,\n            min_split_gain=0.1,\n            random_state=42,\n            verbose=-1\n        ),\n        \n        'XGBoost_Regularized': XGBRegressor(\n            n_estimators=300,\n            learning_rate=0.02,\n            max_depth=6,\n            subsample=0.6,\n            colsample_bytree=0.6,\n            reg_alpha=50,\n            reg_lambda=50,\n            min_child_weight=100,\n            gamma=5,\n            random_state=42,\n            verbosity=0\n        ),\n        \n        'CatBoost_Regularized': CatBoostRegressor(\n            iterations=300,\n            learning_rate=0.02,\n            depth=6,\n            l2_leaf_reg=50,\n            subsample=0.6,\n            colsample_bylevel=0.6,\n            min_data_in_leaf=100,\n            random_seed=42,\n            verbose=False\n        ),\n        \n        'Ridge_Strong': Ridge(alpha=100.0, solver='saga', max_iter=10000),\n        \n        'ElasticNet_Strong': ElasticNet(alpha=1.0, l1_ratio=0.5, max_iter=10000)\n    }\n    \n    # ============================================\n    # HIERARCHICAL ENSEMBLE TRAINING\n    # ============================================\n    print(\"\\n6. Training Hierarchical Ensemble...\")\n    \n    # Use TimeSeriesSplit for better time series validation\n    cv = TimeSeriesSplit(n_splits=5, test_size=len(X_selected)//10)\n    \n    ensemble = HierarchicalEnsemble(\n        base_models=base_models,\n        meta_model=RidgeCV(alphas=[0.01, 0.1, 1.0, 10.0, 100.0], cv=5),\n        use_uncertainty=True\n    )\n    \n    ensemble.fit(X_selected, y, cv=cv)\n    \n    # Get predictions\n    final_predictions = ensemble.predict(X_test_selected)\n    \n    # ============================================\n    # ANALYZE TRAIN-CV GAP\n    # ============================================\n    print(\"\\n7. Analyzing Train-CV Gap...\")\n    \n    # Calculate in-sample predictions for gap analysis\n    train_predictions = ensemble.predict(X_selected)\n    train_score = _pearsonr(y, train_predictions)\n    \n    # Cross-validation scores\n    cv_scores = []\n    for train_idx, val_idx in cv.split(X_selected, y):\n        X_val = X_selected.iloc[val_idx]\n        y_val = y.iloc[val_idx]\n        val_pred = ensemble.predict(X_val)\n        cv_scores.append(_pearsonr(y_val, val_pred))\n    \n    cv_score = np.mean(cv_scores)\n    train_cv_gap = train_score - cv_score\n    \n    print(f\"\\nTrain Score: {train_score:.6f}\")\n    print(f\"CV Score: {cv_score:.6f}\")\n    print(f\"Train-CV Gap: {train_cv_gap:.6f}\")\n    \n    # ============================================\n    # POST-PROCESSING FOR STABILITY\n    # ============================================\n    print(\"\\n8. Post-processing for Stability...\")\n    \n    # Apply slight smoothing to predictions\n    window_size = 5\n    final_predictions_smoothed = pd.Series(final_predictions).rolling(\n        window=window_size, center=True, min_periods=1\n    ).mean().values\n    \n    # Clip extreme predictions based on training distribution\n    train_percentiles = np.percentile(y, [1, 99])\n    final_predictions_clipped = np.clip(\n        final_predictions_smoothed,\n        train_percentiles[0],\n        train_percentiles[1]\n    )\n    \n    # ============================================\n    # SAVE RESULTS\n    # ============================================\n    print(\"\\n9. Saving Results...\")\n    \n    # Save submission\n    sub = pd.read_csv(CFG.sample_sub_path)\n    sub[\"prediction\"] = final_predictions_clipped\n    sub.to_csv(\"submission_advanced.csv\", index=False)\n    print(\"Submission saved to submission_advanced.csv\")\n    \n    # Save models and preprocessors\n    joblib.dump(noise_reducer, \"noise_reducer.pkl\")\n    joblib.dump(feature_selector, \"feature_selector.pkl\")\n    joblib.dump(ensemble, \"ensemble_model.pkl\")\n    \n    # Diagnostic plots\n    plt.figure(figsize=(15, 10))\n    \n    # Plot 1: Feature importance from adaptive selection\n    plt.subplot(2, 2, 1)\n    feature_scores = feature_selector.feature_scores[feature_selector.selected_features]\n    plt.hist(feature_scores, bins=30, alpha=0.7, color='blue')\n    plt.title('Selected Feature Scores Distribution')\n    plt.xlabel('Feature Score')\n    plt.ylabel('Count')\n    \n    # Plot 2: Prediction distribution\n    plt.subplot(2, 2, 2)\n    plt.hist(y, bins=50, alpha=0.5, label='Training', color='blue')\n    plt.hist(final_predictions_clipped, bins=50, alpha=0.5, label='Test Predictions', color='red')\n    plt.title('Prediction Distribution')\n    plt.xlabel('Value')\n    plt.ylabel('Count')\n    plt.legend()\n    \n    # Plot 3: Cross-validation scores\n    plt.subplot(2, 2, 3)\n    plt.plot(cv_scores, marker='o', linewidth=2, markersize=8)\n    plt.axhline(y=cv_score, color='r', linestyle='--', label=f'Mean: {cv_score:.6f}')\n    plt.title('Cross-Validation Scores by Fold')\n    plt.xlabel('Fold')\n    plt.ylabel('Pearson Correlation')\n    plt.legend()\n    plt.grid(True, alpha=0.3)\n    \n    # Plot 4: Train vs CV predictions scatter\n    plt.subplot(2, 2, 4)\n    sample_size = min(5000, len(y))\n    sample_idx = np.random.choice(len(y), sample_size, replace=False)\n    plt.scatter(y.iloc[sample_idx], train_predictions[sample_idx], alpha=0.5, s=1)\n    plt.plot([y.min(), y.max()], [y.min(), y.max()], 'r--', lw=2)\n    plt.title(f'Train Predictions (Sample, r={train_score:.6f})')\n    plt.xlabel('True Values')\n    plt.ylabel('Predicted Values')\n    \n    plt.tight_layout()\n    plt.savefig('advanced_model_diagnostics.png', dpi=300)\n    plt.show()\n    \n    print(\"\\nPipeline completed successfully!\")\n    print(f\"Final Train-CV Gap: {train_cv_gap:.6f}\")\n    \n    return final_predictions_clipped, train_score, cv_score, train_cv_gap\n\n# ============================================\n# RUN MAIN PIPELINE\n# ============================================\nif __name__ == \"__main__\":\n    final_predictions, train_score, cv_score, train_cv_gap = main()","metadata":{"trusted":true},"outputs":[],"execution_count":null}]}