{"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":"nvidiaTeslaT4","dataSources":[{"sourceId":96164,"databundleVersionId":12993472,"sourceType":"competition"}],"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":true}},"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":"# Import all required libraries\nimport numpy as np\nimport pandas as pd\nimport os\nimport gc\nimport warnings\nwarnings.filterwarnings('ignore')\n\n# Print input files\nfor dirname, _, filenames in os.walk('/kaggle/input'):\n    for filename in filenames:\n        print(os.path.join(dirname, filename))\n\n# Install required packages\nprint(\"Installing required packages...\")\n!pip install koolbox scikit-learn==1.5.2 --quiet\n\n# Import modeling libraries\nfrom sklearn.model_selection import KFold\nfrom lightgbm import LGBMRegressor\nfrom scipy.stats import pearsonr, kurtosis, skew\nfrom xgboost import XGBRegressor\nfrom sklearn.base import clone, BaseEstimator, RegressorMixin\nfrom koolbox import Trainer\nimport matplotlib.pyplot as plt\nimport seaborn as sns\nfrom sklearn.preprocessing import RobustScaler, QuantileTransformer\nimport pickle\nfrom dataclasses import dataclass\nfrom typing import Dict, Tuple, List\n\n# PyTorch imports for GANDALF\nimport torch\nimport torch.nn as nn\nimport torch.optim as optim\nfrom torch.utils.data import DataLoader, TensorDataset\n\n# TensorFlow imports for AutoEncoder\nimport tensorflow as tf\n\n# Configuration\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@dataclass\nclass CompressionProfile:\n    \"\"\"Stores compression parameters for each feature\"\"\"\n    method: str\n    strength: float\n    median: float\n    mad: float\n    noise_level: float\n    outlier_score: float\n    optimal_compression: float\n\nclass AdaptiveCompressor:\n    \"\"\"Learns optimal compression for each feature based on its characteristics\"\"\"\n    \n    def __init__(self):\n        self.compression_profiles: Dict[str, CompressionProfile] = {}\n        self.label_profile: CompressionProfile = None\n        \n    def analyze_feature(self, data: np.ndarray, feature_name: str) -> Dict[str, float]:\n        \"\"\"Analyze feature characteristics to determine optimal compression\"\"\"\n        # Remove NaN values for analysis\n        clean_data = data[~np.isnan(data)]\n        \n        if len(clean_data) == 0:\n            return {\n                'noise_level': 0,\n                'outlier_score': 0,\n                'distribution_score': 0,\n                'optimal_compression': 0\n            }\n        \n        # Calculate robust statistics\n        median = np.median(clean_data)\n        mad = np.median(np.abs(clean_data - median))\n        \n        # Noise level estimation (using high-frequency components)\n        if len(clean_data) > 100:\n            # Simple noise estimation using difference between consecutive values\n            diffs = np.diff(np.sort(clean_data))\n            noise_level = np.std(diffs[diffs != 0]) / (mad + 1e-8)\n        else:\n            noise_level = 0.1\n            \n        # Outlier score (proportion of extreme values)\n        if mad > 0:\n            z_scores = np.abs((clean_data - median) / mad)\n            outlier_score = np.mean(z_scores > 3)\n        else:\n            outlier_score = 0\n            \n        # Distribution characteristics\n        try:\n            kurt = kurtosis(clean_data)\n            skewness = abs(skew(clean_data))\n            distribution_score = np.tanh(kurt / 10) * 0.5 + np.tanh(skewness / 5) * 0.5\n        except:\n            distribution_score = 0.5\n            \n        # Determine optimal compression based on characteristics\n        # Higher noise, more outliers, or extreme distributions = more compression\n        optimal_compression = np.clip(\n            0.2 * noise_level + 0.4 * outlier_score + 0.4 * distribution_score,\n            0, 0.8\n        )\n        \n        return {\n            'median': median,\n            'mad': mad,\n            'noise_level': noise_level,\n            'outlier_score': outlier_score,\n            'distribution_score': distribution_score,\n            'optimal_compression': optimal_compression\n        }\n    \n    def fit(self, X: pd.DataFrame, y: pd.Series = None, \n            feature_compression_override: Dict[str, float] = None,\n            label_compression: float = 0.0):\n        \"\"\"Learn compression parameters for each feature\"\"\"\n        \n        print(\"Learning adaptive compression parameters...\")\n        \n        # Analyze each feature\n        for col in X.columns:\n            analysis = self.analyze_feature(X[col].values, col)\n            \n            # Use override if provided, otherwise use adaptive compression\n            if feature_compression_override and col in feature_compression_override:\n                compression = feature_compression_override[col]\n            else:\n                compression = analysis['optimal_compression']\n            \n            # Determine best method based on distribution\n            if analysis['outlier_score'] > 0.1:\n                method = 'tanh'  # Best for heavy outliers\n            elif abs(analysis.get('skewness', 0)) > 2:\n                method = 'log'   # Best for skewed distributions\n            else:\n                method = 'sigmoid'  # General purpose\n            \n            self.compression_profiles[col] = CompressionProfile(\n                method=method,\n                strength=compression,\n                median=analysis['median'],\n                mad=analysis['mad'],\n                noise_level=analysis['noise_level'],\n                outlier_score=analysis['outlier_score'],\n                optimal_compression=analysis['optimal_compression']\n            )\n        \n        # Analyze label if provided\n        if y is not None and label_compression > 0:\n            label_analysis = self.analyze_feature(y.values, 'label')\n            self.label_profile = CompressionProfile(\n                method='tanh',  # Conservative for labels\n                strength=label_compression,\n                median=label_analysis['median'],\n                mad=label_analysis['mad'],\n                noise_level=label_analysis['noise_level'],\n                outlier_score=label_analysis['outlier_score'],\n                optimal_compression=label_compression\n            )\n            print(f\"  Label: tanh compression at {label_compression:.2f}\")\n    \n    def transform(self, X: pd.DataFrame, y: pd.Series = None) -> Tuple[pd.DataFrame, pd.Series]:\n        \"\"\"Apply learned compression to features and optionally labels\"\"\"\n        X_compressed = X.copy()\n        \n        # Compress each feature using its profile\n        for col in X.columns:\n            if col in self.compression_profiles:\n                profile = self.compression_profiles[col]\n                X_compressed[col] = self._compress_data(\n                    X[col].values,\n                    profile\n                )\n        \n        # Compress labels if profile exists\n        y_compressed = y\n        if y is not None and self.label_profile is not None:\n            y_compressed = pd.Series(\n                self._compress_data(y.values, self.label_profile),\n                index=y.index\n            )\n        \n        return X_compressed, y_compressed\n    \n    def _compress_data(self, data: np.ndarray, profile: CompressionProfile) -> np.ndarray:\n        \"\"\"Apply compression using the specified profile\"\"\"\n        # Handle NaN values\n        nan_mask = np.isnan(data)\n        data_clean = data[~nan_mask]\n        \n        if len(data_clean) == 0 or profile.strength == 0:\n            return data\n        \n        # Normalize using stored statistics\n        if profile.mad > 0:\n            normalized = (data - profile.median) / (profile.mad * 6)\n        else:\n            return data\n        \n        # Apply compression method\n        if profile.method == 'tanh':\n            compressed = np.tanh(normalized * (1 - profile.strength))\n        elif profile.method == 'sigmoid':\n            compressed = 2 / (1 + np.exp(-normalized * (1 - profile.strength) * 2)) - 1\n        elif profile.method == 'arctan':\n            compressed = (2/np.pi) * np.arctan(normalized * (1 - profile.strength) * np.pi/2)\n        elif profile.method == 'log':\n            sign = np.sign(normalized)\n            abs_norm = np.abs(normalized)\n            compressed = sign * np.log1p(abs_norm * (1 - profile.strength)) / np.log1p(1 - profile.strength)\n        \n        # Scale back\n        result = compressed * (profile.mad * 6) + profile.median\n        \n        # Preserve NaN values\n        if nan_mask.any():\n            full_result = np.full_like(data, np.nan, dtype=np.float64)\n            full_result[~nan_mask] = result[~nan_mask]\n            return full_result\n        \n        return result\n    \n    def save(self, filepath: str):\n        \"\"\"Save compression profiles\"\"\"\n        with open(filepath, 'wb') as f:\n            pickle.dump({\n                'compression_profiles': self.compression_profiles,\n                'label_profile': self.label_profile\n            }, f)\n    \n    def load(self, filepath: str):\n        \"\"\"Load compression profiles\"\"\"\n        with open(filepath, 'rb') as f:\n            data = pickle.load(f)\n            self.compression_profiles = data['compression_profiles']\n            self.label_profile = data['label_profile']\n\ndef create_microstructure_features(df):\n    \"\"\"Create advanced microstructure features from market data\"\"\"\n    print(\"Creating microstructure features...\")\n    \n    # Price-related features (using volume as proxy for price movements)\n    df['volume_log'] = np.log1p(df['volume'])\n    df['volume_sqrt'] = np.sqrt(df['volume'])\n    \n    # Bid-Ask Spread features\n    df['bid_ask_spread'] = df['ask_qty'] - df['bid_qty']\n    df['bid_ask_ratio'] = df['bid_qty'] / (df['ask_qty'] + 1e-8)\n    df['bid_ask_imbalance'] = (df['bid_qty'] - df['ask_qty']) / (df['bid_qty'] + df['ask_qty'] + 1e-8)\n    df['bid_ask_pressure'] = np.log1p(df['bid_qty']) - np.log1p(df['ask_qty'])\n    \n    # Order Flow Imbalance\n    df['order_flow_imbalance'] = (df['buy_qty'] - df['sell_qty']) / (df['buy_qty'] + df['sell_qty'] + 1e-8)\n    df['order_flow_ratio'] = df['buy_qty'] / (df['sell_qty'] + 1e-8)\n    df['net_order_flow'] = df['buy_qty'] - df['sell_qty']\n    df['order_flow_log_ratio'] = np.log1p(df['buy_qty']) - np.log1p(df['sell_qty'])\n    \n    # Trade Intensity\n    df['trade_intensity'] = (df['buy_qty'] + df['sell_qty']) / df['volume'].clip(lower=1)\n    df['trade_size_imbalance'] = (df['buy_qty'] - df['sell_qty']) / df['volume'].clip(lower=1)\n    \n    # Liquidity measures\n    df['total_liquidity'] = df['bid_qty'] + df['ask_qty']\n    df['liquidity_imbalance'] = df['bid_ask_imbalance'] * df['total_liquidity']\n    df['liquidity_ratio'] = df['total_liquidity'] / (df['volume'] + 1e-8)\n    df['liquidity_consumption'] = (df['buy_qty'] + df['sell_qty']) / (df['total_liquidity'] + 1e-8)\n    \n    # Microstructure volatility proxies\n    df['qty_volatility'] = np.sqrt(df[['bid_qty', 'ask_qty', 'buy_qty', 'sell_qty']].var(axis=1))\n    df['flow_volatility'] = np.abs(df['order_flow_imbalance'])\n    \n    # Advanced ratios\n    df['vwap_proxy'] = df['volume'] / (df['buy_qty'] + df['sell_qty'] + 1e-8)\n    df['aggressive_ratio'] = (df['buy_qty'] + df['sell_qty']) / (df['bid_qty'] + df['ask_qty'] + 1e-8)\n    df['passive_aggressive_imbalance'] = (df['bid_qty'] + df['ask_qty'] - df['buy_qty'] - df['sell_qty']) / (df['volume'] + 1e-8)\n    \n    # Market depth features\n    df['depth_imbalance'] = (df['bid_qty'] * df['buy_qty'] - df['ask_qty'] * df['sell_qty']) / (df['volume']**2 + 1e-8)\n    df['weighted_order_flow'] = df['order_flow_imbalance'] * np.log1p(df['volume'])\n    \n    # Interaction features\n    df['bid_volume_interaction'] = df['bid_qty'] * df['volume']\n    df['ask_volume_interaction'] = df['ask_qty'] * df['volume']\n    df['imbalance_volume_interaction'] = df['bid_ask_imbalance'] * df['volume']\n    \n    # Power transformations\n    df['bid_qty_squared'] = df['bid_qty'] ** 2\n    df['ask_qty_squared'] = df['ask_qty'] ** 2\n    df['order_flow_imbalance_squared'] = df['order_flow_imbalance'] ** 2\n    \n    # Replace infinities and NaNs\n    df = df.replace([np.inf, -np.inf], np.nan)\n    \n    return df\n\ndef reduce_mem_usage(dataframe, dataset):    \n    print('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\n        c_min = dataframe[col].min()\n        c_max = dataframe[col].max()\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('--- Memory usage before: {:.2f} MB'.format(initial_mem_usage))\n    print('--- Memory usage after: {:.2f} MB'.format(final_mem_usage))\n    print('--- Decreased memory usage by {:.1f}%\\n'.format(100 * (initial_mem_usage - final_mem_usage) / initial_mem_usage))\n\n    return dataframe\n\ndef _pearsonr(y_true, y_pred):\n    return pearsonr(y_true, y_pred)[0]\n\n# Define model parameters\nlgbm_params = {\n    \"boosting_type\": \"gbdt\",\n    \"colsample_bytree\": 0.5625888953382505,\n    \"learning_rate\": 0.029312951475451557,\n    \"min_child_samples\": 63,\n    \"min_child_weight\": 0.11456572852335424,\n    \"n_estimators\": 126,\n    \"n_jobs\": -1,\n    \"num_leaves\": 37,\n    \"random_state\": 42,\n    \"reg_alpha\": 85.2476527854083,\n    \"reg_lambda\": 99.38305361388907,\n    \"subsample\": 0.450669817684892,\n    \"verbose\": -1\n}\n\nlgbm_goss_params = {\n    \"boosting_type\": \"goss\",\n    \"colsample_bytree\": 0.34695458228489784,\n    \"learning_rate\": 0.031023014900595287,\n    \"min_child_samples\": 30,\n    \"min_child_weight\": 0.4727729225033618,\n    \"n_estimators\": 220,\n    \"n_jobs\": -1,\n    \"num_leaves\": 58,\n    \"random_state\": 42,\n    \"reg_alpha\": 38.665994901468224,\n    \"reg_lambda\": 92.76991677464294,\n    \"subsample\": 0.4810891284493255,\n    \"verbose\": -1\n}\n\nxgb_params = {\n    \"colsample_bylevel\": 0.4778015829774066,\n    \"colsample_bynode\": 0.362764358742407,\n    \"colsample_bytree\": 0.7107423488010493,\n    \"gamma\": 1.7094857725240398,\n    \"learning_rate\": 0.02213323588455387,\n    \"max_depth\": 20,\n    \"max_leaves\": 12,\n    \"min_child_weight\": 16,\n    \"n_estimators\": 1667,\n    \"n_jobs\": -1,\n    \"random_state\": 42,\n    \"reg_alpha\": 39.352415706891264,\n    \"reg_lambda\": 75.44843704068275,\n    \"subsample\": 0.06566669853471274,\n    \"verbosity\": 0\n}\n\n# GANDALF Model Implementation\nclass GANDALF(BaseEstimator, RegressorMixin):\n    def __init__(self, n_estimators=100, learning_rate=0.01, max_depth=5, \n                 feature_fraction=0.8, bagging_fraction=0.8, lambda_reg=1.0,\n                 min_data_in_leaf=20, num_iterations=100, random_state=42):\n        self.n_estimators = n_estimators\n        self.learning_rate = learning_rate\n        self.max_depth = max_depth\n        self.feature_fraction = feature_fraction\n        self.bagging_fraction = bagging_fraction\n        self.lambda_reg = lambda_reg\n        self.min_data_in_leaf = min_data_in_leaf\n        self.num_iterations = num_iterations\n        self.random_state = random_state\n        self.device = torch.device('cuda' if torch.cuda.is_available() else 'cpu')\n        \n    def _build_network(self, input_dim):\n        \"\"\"Build the neural network architecture for GANDALF\"\"\"\n        class GANDALFNet(nn.Module):\n            def __init__(self, input_dim, hidden_dims=[256, 128, 64]):\n                super(GANDALFNet, self).__init__()\n                \n                layers = []\n                prev_dim = input_dim\n                \n                for hidden_dim in hidden_dims:\n                    layers.extend([\n                        nn.Linear(prev_dim, hidden_dim),\n                        nn.BatchNorm1d(hidden_dim),\n                        nn.ReLU(),\n                        nn.Dropout(0.3)\n                    ])\n                    prev_dim = hidden_dim\n                \n                layers.append(nn.Linear(prev_dim, 1))\n                \n                self.network = nn.Sequential(*layers)\n                \n            def forward(self, x):\n                return self.network(x)\n        \n        return GANDALFNet(input_dim).to(self.device)\n    \n    def fit(self, X, y):\n        # Convert to tensors\n        X_tensor = torch.FloatTensor(X.values if hasattr(X, 'values') else X).to(self.device)\n        y_tensor = torch.FloatTensor(y.values if hasattr(y, 'values') else y).reshape(-1, 1).to(self.device)\n        \n        # Build network\n        self.model = self._build_network(X_tensor.shape[1])\n        \n        # Create dataset and dataloader\n        dataset = TensorDataset(X_tensor, y_tensor)\n        dataloader = DataLoader(dataset, batch_size=1024, shuffle=True)\n        \n        # Optimizer\n        optimizer = optim.AdamW(self.model.parameters(), lr=self.learning_rate, weight_decay=self.lambda_reg)\n        criterion = nn.MSELoss()\n        \n        # Training loop\n        self.model.train()\n        for epoch in range(self.num_iterations):\n            total_loss = 0\n            for batch_X, batch_y in dataloader:\n                optimizer.zero_grad()\n                outputs = self.model(batch_X)\n                loss = criterion(outputs, batch_y)\n                loss.backward()\n                optimizer.step()\n                total_loss += loss.item()\n            \n            if (epoch + 1) % 20 == 0:\n                print(f'Epoch [{epoch+1}/{self.num_iterations}], Loss: {total_loss/len(dataloader):.4f}')\n        \n        return self\n    \n    def predict(self, X):\n        self.model.eval()\n        X_tensor = torch.FloatTensor(X.values if hasattr(X, 'values') else X).to(self.device)\n        \n        with torch.no_grad():\n            predictions = self.model(X_tensor).cpu().numpy().flatten()\n        \n        return predictions\n\n# AutoEncoder MLP\nclass AutoEncoderMLP(BaseEstimator, RegressorMixin):\n    def __init__(self, num_columns, hidden_units, dropout_rates, lr=1e-3):\n        self.num_columns = num_columns\n        self.hidden_units = hidden_units\n        self.dropout_rates = dropout_rates\n        self.lr = lr\n        self.model = self._build_model()\n    \n    def _build_model(self):\n        inp = tf.keras.layers.Input(shape=(self.num_columns,))\n        x0 = tf.keras.layers.BatchNormalization()(inp)\n\n        encoder = tf.keras.layers.GaussianNoise(self.dropout_rates[0])(x0)\n        encoder = tf.keras.layers.Dense(self.hidden_units[0])(encoder)\n        encoder = tf.keras.layers.BatchNormalization()(encoder)\n        encoder = tf.keras.layers.Activation('swish')(encoder)\n\n        decoder = tf.keras.layers.Dropout(self.dropout_rates[1])(encoder)\n        decoder = tf.keras.layers.Dense(self.num_columns, name='decoder')(decoder)\n\n        x_reg = tf.keras.layers.Dense(self.hidden_units[1])(encoder)\n        x_reg = tf.keras.layers.BatchNormalization()(x_reg)\n        x_reg = tf.keras.layers.Activation('swish')(x_reg)\n        x_reg = tf.keras.layers.Dropout(self.dropout_rates[2])(x_reg)\n\n        out_reg = tf.keras.layers.Dense(1, activation='linear', name='target')(x_reg)\n\n        model = tf.keras.models.Model(inputs=inp, outputs=[decoder, out_reg])\n        model.compile(\n            optimizer=tf.keras.optimizers.Adam(learning_rate=self.lr),\n            loss={\"decoder\": tf.keras.losses.MeanSquaredError(),\n                  \"target\": tf.keras.losses.MeanSquaredError()},\n            loss_weights={\"decoder\": 0.3, \"target\": 1.0}\n        )\n        return model\n\n    def fit(self, X, y):\n        self.model.fit(\n            X, {\"decoder\": X, \"target\": y},\n            epochs=50,\n            batch_size=8192,\n            validation_split=0.2,\n            callbacks=[\n                tf.keras.callbacks.EarlyStopping(patience=10, restore_best_weights=True),\n                tf.keras.callbacks.ReduceLROnPlateau(patience=5)\n            ],\n            verbose=0\n        )\n        return self\n\n    def predict(self, X):\n        _, y_pred = self.model.predict(X, verbose=0)\n        return y_pred.flatten()\n\ndef run_pipeline_with_adaptive_compression(feature_compression: float, \n                                         label_compression: float,\n                                         use_adaptive: bool,\n                                         strategy_name: str):\n    \"\"\"\n    Run the entire pipeline with adaptive compression strategy\n    \"\"\"\n    print(f\"\\n{'='*80}\")\n    print(f\"Running pipeline with compression strategy: {strategy_name}\")\n    print(f\"Feature compression: {feature_compression:.2f}, Label compression: {label_compression:.2f}\")\n    print(f\"Adaptive: {use_adaptive}\")\n    print(f\"{'='*80}\\n\")\n    \n    # Load 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 for memory efficiency\n    train_size = len(train)\n    if train_size > 200000:\n        train = train.iloc[-200000:].reset_index(drop=True)\n        print(f\"Using last 200,000 rows of training data for memory efficiency\")\n    \n    # Create microstructure features\n    train = create_microstructure_features(train)\n    test = create_microstructure_features(test)\n    \n    # Define enhanced feature set\n    # Original high-performing features from your analysis\n    TOP_FEATURES = ['X759', 'X752', 'X298', 'X29', 'X586', 'X43', 'X30', 'X31', 'X32', 'X33', 'X35']\n    \n    # Extended feature set based on original code\n    ORIGINAL_FEATURES = ['X363', 'X321', 'X405', 'X730', 'X523', 'X756', 'X589', 'X462', 'X779',\n                        'X25', 'X532', 'X520', 'X329', 'X383', 'X751', 'X535', 'X639', 'X596', 'X761',\n                        'X287', 'X302', 'X55', 'X56', 'X52', 'X303', 'X51',\n                        'X598', 'X385', 'X603', 'X674', 'X415', 'X345', 'X174', 'X178', 'X168', 'X612']\n    \n    # Market microstructure features\n    MICROSTRUCTURE_FEATURES = [\n        'bid_qty', 'ask_qty', 'buy_qty', 'sell_qty', 'volume',\n        'volume_log', 'volume_sqrt',\n        'bid_ask_spread', 'bid_ask_ratio', 'bid_ask_imbalance', 'bid_ask_pressure',\n        'order_flow_imbalance', 'order_flow_ratio', 'net_order_flow', 'order_flow_log_ratio',\n        'trade_intensity', 'trade_size_imbalance',\n        'total_liquidity', 'liquidity_imbalance', 'liquidity_ratio', 'liquidity_consumption',\n        'qty_volatility', 'flow_volatility',\n        'vwap_proxy', 'aggressive_ratio', 'passive_aggressive_imbalance',\n        'depth_imbalance', 'weighted_order_flow',\n        'bid_volume_interaction', 'ask_volume_interaction', 'imbalance_volume_interaction',\n        'bid_qty_squared', 'ask_qty_squared', 'order_flow_imbalance_squared'\n    ]\n    \n    # Combine all features (NO automatic addition of first 100 features)\n    ALL_FEATURES = list(set(TOP_FEATURES + ORIGINAL_FEATURES + MICROSTRUCTURE_FEATURES))\n    \n    # Filter to only existing columns\n    selected_columns = [col for col in ALL_FEATURES if col in train.columns]\n    \n    print(f\"Total features selected: {len(selected_columns)}\")\n    print(f\"- Top ranked features: {len([f for f in TOP_FEATURES if f in selected_columns])}\")\n    print(f\"- Original features: {len([f for f in ORIGINAL_FEATURES if f in selected_columns])}\")\n    print(f\"- Microstructure features: {len([f for f in MICROSTRUCTURE_FEATURES if f in selected_columns])}\")\n    \n    # Select columns\n    train = train[selected_columns + [CFG.target]]\n    test = test[selected_columns]\n    \n    # Handle missing values\n    train = train.fillna(0)\n    test = test.fillna(0)\n    \n    # Prepare data\n    X_train = train.drop(CFG.target, axis=1)\n    y_train = train[CFG.target]\n    X_test = test\n    \n    # Initialize and fit compressor\n    compressor = AdaptiveCompressor()\n    \n    if use_adaptive:\n        # Let the compressor determine optimal compression for each feature\n        compressor.fit(X_train, y_train, label_compression=label_compression)\n    else:\n        # Use uniform compression for all features\n        feature_overrides = {col: feature_compression for col in X_train.columns}\n        compressor.fit(X_train, y_train, \n                      feature_compression_override=feature_overrides,\n                      label_compression=label_compression)\n    \n    # Apply compression\n    X_train_compressed, y_train_compressed = compressor.transform(X_train, y_train)\n    X_test_compressed, _ = compressor.transform(X_test, None)\n    \n    # Save compressor for future use\n    compressor.save(f'compressor_{strategy_name}.pkl')\n    \n    # Visualize compression effects for a few features\n    if use_adaptive:\n        fig, axes = plt.subplots(2, 2, figsize=(12, 10))\n        axes = axes.ravel()\n        \n        features_to_plot = []\n        # Try to plot some top features first\n        for feat in TOP_FEATURES[:4]:\n            if feat in X_train.columns:\n                features_to_plot.append(feat)\n        # Fill with other features if needed\n        for feat in selected_columns:\n            if len(features_to_plot) >= 4:\n                break\n            if feat not in features_to_plot:\n                features_to_plot.append(feat)\n        \n        for idx, feat in enumerate(features_to_plot[:4]):\n            ax = axes[idx]\n            \n            original = X_train[feat].values\n            compressed = X_train_compressed[feat].values\n            \n            # Remove NaN for plotting\n            mask = ~(np.isnan(original) | np.isnan(compressed))\n            \n            ax.scatter(original[mask], compressed[mask], alpha=0.1, s=1)\n            ax.plot([original[mask].min(), original[mask].max()], \n                   [original[mask].min(), original[mask].max()], 'r--', alpha=0.5)\n            \n            profile = compressor.compression_profiles[feat]\n            ax.set_title(f'{feat}\\nMethod: {profile.method}, Strength: {profile.strength:.2f}')\n            ax.set_xlabel('Original')\n            ax.set_ylabel('Compressed')\n        \n        plt.tight_layout()\n        plt.savefig(f'adaptive_compression_{strategy_name}.png', dpi=150, bbox_inches='tight')\n        plt.close()\n    \n    # Reduce memory\n    X_train_compressed = reduce_mem_usage(X_train_compressed, \"train\")\n    X_test_compressed = reduce_mem_usage(X_test_compressed, \"test\")\n    \n    # Initialize storage\n    scores = {}\n    oof_preds = {}\n    test_preds = {}\n    \n    # Train models\n    # Train LightGBM (gbdt)\n    print(f\"\\nTraining LightGBM (gbdt) for {strategy_name}...\")\n    lgbm_trainer = Trainer(\n        LGBMRegressor(**lgbm_params),\n        cv=KFold(n_splits=5, shuffle=False),\n        metric=_pearsonr,\n        task=\"regression\",\n        metric_precision=6\n    )\n    lgbm_trainer.fit(X_train_compressed, y_train_compressed)\n    scores[\"LightGBM (gbdt)\"] = lgbm_trainer.fold_scores\n    oof_preds[\"LightGBM (gbdt)\"] = lgbm_trainer.oof_preds\n    test_preds[\"LightGBM (gbdt)\"] = lgbm_trainer.predict(X_test_compressed)\n    \n    # Try to extract feature importance\n    try:\n        if hasattr(lgbm_trainer, 'models') and len(lgbm_trainer.models) > 0:\n            feature_importance = pd.DataFrame({\n                'feature': X_train_compressed.columns,\n                'importance': lgbm_trainer.models[0].feature_importances_\n            }).sort_values('importance', ascending=False)\n            \n            # Save top 50 features\n            feature_importance.head(50).to_csv(f'feature_importance_{strategy_name}.csv', index=False)\n            \n            # Plot feature importance for top features\n            plt.figure(figsize=(10, 8))\n            top_features = feature_importance.head(30)\n            plt.barh(range(len(top_features)), top_features['importance'])\n            plt.yticks(range(len(top_features)), top_features['feature'])\n            plt.xlabel('Importance')\n            plt.title(f'Top 30 Features - {strategy_name}')\n            plt.tight_layout()\n            plt.savefig(f'feature_importance_{strategy_name}.png', dpi=150, bbox_inches='tight')\n            plt.close()\n    except:\n        print(\"Could not extract feature importance\")\n    \n    # Train LightGBM (goss)\n    print(f\"\\nTraining LightGBM (goss) for {strategy_name}...\")\n    lgbm_goss_trainer = Trainer(\n        LGBMRegressor(**lgbm_goss_params),\n        cv=KFold(n_splits=5, shuffle=False),\n        metric=_pearsonr,\n        task=\"regression\",\n        metric_precision=6\n    )\n    lgbm_goss_trainer.fit(X_train_compressed, y_train_compressed)\n    scores[\"LightGBM (goss)\"] = lgbm_goss_trainer.fold_scores\n    oof_preds[\"LightGBM (goss)\"] = lgbm_goss_trainer.oof_preds\n    test_preds[\"LightGBM (goss)\"] = lgbm_goss_trainer.predict(X_test_compressed)\n    \n    # Train XGBoost\n    print(f\"\\nTraining XGBoost for {strategy_name}...\")\n    xgb_trainer = Trainer(\n        XGBRegressor(**xgb_params),\n        cv=KFold(n_splits=5, shuffle=False),\n        metric=_pearsonr,\n        task=\"regression\",\n        metric_precision=6\n    )\n    xgb_trainer.fit(X_train_compressed, y_train_compressed)\n    scores[\"XGBoost\"] = xgb_trainer.fold_scores\n    oof_preds[\"XGBoost\"] = xgb_trainer.oof_preds\n    test_preds[\"XGBoost\"] = xgb_trainer.predict(X_test_compressed)\n    \n    # Train GANDALF (reduce iterations for memory)\n    print(f\"\\nTraining GANDALF for {strategy_name}...\")\n    gandalf_model = GANDALF(\n        n_estimators=100,\n        learning_rate=0.001,\n        max_depth=5,\n        feature_fraction=0.8,\n        bagging_fraction=0.8,\n        lambda_reg=0.1,\n        min_data_in_leaf=20,\n        num_iterations=30,  # Reduced from 50\n        random_state=42\n    )\n    gandalf_trainer = Trainer(\n        gandalf_model,\n        cv=KFold(n_splits=5, shuffle=False),\n        metric=_pearsonr,\n        task=\"regression\",\n        metric_precision=6\n    )\n    gandalf_trainer.fit(X_train_compressed, y_train_compressed)\n    scores[\"GANDALF\"] = gandalf_trainer.fold_scores\n    oof_preds[\"GANDALF\"] = gandalf_trainer.oof_preds\n    test_preds[\"GANDALF\"] = gandalf_trainer.predict(X_test_compressed)\n    \n    # Create ensemble predictions\n    X_ensemble = pd.DataFrame(oof_preds)\n    X_test_ensemble = pd.DataFrame(test_preds)\n    \n    # Train AutoEncoder\n    print(f\"\\nTraining AutoEncoder for {strategy_name}...\")\n    ae_model = AutoEncoderMLP(\n        num_columns=X_ensemble.shape[1],\n        hidden_units=[128, 128],\n        dropout_rates=[0.05, 0.1, 0.2],\n        lr=1e-3\n    )\n    ae_trainer = Trainer(\n        ae_model,\n        cv=KFold(n_splits=5, shuffle=False),\n        metric=_pearsonr,\n        task=\"regression\",\n        metric_precision=6\n    )\n    ae_trainer.fit(X_ensemble, y_train_compressed)\n    scores[\"AutoEncoder\"] = ae_trainer.fold_scores\n    oof_preds[\"AutoEncoder\"] = ae_trainer.oof_preds\n    ae_test_preds = ae_trainer.predict(X_test_ensemble)\n    \n    # Create weighted ensemble\n    ensemble_weights = {\n        \"LightGBM (gbdt)\": 0.285,\n        \"LightGBM (goss)\": 0.285,\n        \"XGBoost\": 0.285,\n        \"GANDALF\": 0.05,\n        \"AutoEncoder\": 0.095\n    }\n    \n    # Calculate weighted predictions\n    weighted_test_pred = np.zeros(len(X_test))\n    for model_name, weight in ensemble_weights.items():\n        if model_name == \"AutoEncoder\":\n            weighted_test_pred += weight * ae_test_preds\n        else:\n            weighted_test_pred += weight * test_preds[model_name]\n    \n    # If labels were compressed, we need to inverse transform predictions\n    if label_compression > 0 and compressor.label_profile is not None:\n        print(\"Note: Predictions are in compressed space. Consider inverse transformation for submission.\")\n    \n    # Print results\n    print(f\"\\nResults for {strategy_name}:\")\n    scores_df = pd.DataFrame(scores)\n    mean_scores = scores_df.mean()\n    print(mean_scores.sort_values(ascending=False))\n    \n    # Export compression profile summary\n    if use_adaptive:\n        compression_summary = pd.DataFrame([\n            {\n                'feature': feat,\n                'method': profile.method,\n                'strength': profile.strength,\n                'noise_level': profile.noise_level,\n                'outlier_score': profile.outlier_score\n            }\n            for feat, profile in compressor.compression_profiles.items()\n        ])\n        compression_summary.to_csv(f'compression_summary_{strategy_name}.csv', index=False)\n    \n    # Clean up memory\n    del X_train, y_train, X_test, train, test\n    gc.collect()\n    \n    return weighted_test_pred, scores_df, compressor, selected_columns\n\n# Main execution\nif __name__ == \"__main__\":\n    # Define compression grid with enhanced features\n    compression_grid = [\n        # (feature_compression, label_compression, use_adaptive, name)\n        (0.0, 0.0, False, \"baseline_enhanced_features\"),\n        (0.3, 0.0, False, \"features_30pct_enhanced\"),\n        (0.5, 0.0, False, \"features_50pct_enhanced\"),\n        (0.3, 0.1, False, \"both_30_10pct_enhanced\"),\n        (0.0, 0.0, True, \"adaptive_enhanced\"),\n        (0.0, 0.1, True, \"adaptive_label_10pct_enhanced\"),\n    ]\n    \n    # Store results\n    all_predictions = {}\n    all_scores = {}\n    all_compressors = {}\n    all_features_used = {}\n    \n    # Run grid search\n    for feature_comp, label_comp, use_adaptive, name in compression_grid:\n        try:\n            predictions, scores, compressor, features_used = run_pipeline_with_adaptive_compression(\n                feature_comp, label_comp, use_adaptive, name\n            )\n            \n            all_predictions[name] = predictions\n            all_scores[name] = scores\n            all_compressors[name] = compressor\n            all_features_used[name] = features_used\n            \n            # Save individual submission\n            sub = pd.read_csv(CFG.sample_sub_path)\n            sub[\"prediction\"] = predictions\n            sub.to_csv(f\"submission_{name}.csv\", index=False)\n            print(f\"\\nSaved submission_{name}.csv\")\n            \n        except Exception as e:\n            print(f\"\\nError in strategy {name}: {str(e)}\")\n            import traceback\n            traceback.print_exc()\n            continue\n    \n    # Create ensemble of all strategies\n    print(\"\\n\" + \"=\"*80)\n    print(\"Creating ensemble of all strategies...\")\n    print(\"=\"*80)\n    \n    if len(all_predictions) > 0:\n        ensemble_predictions = np.zeros(len(next(iter(all_predictions.values()))))\n        for strategy_name in all_predictions:\n            ensemble_predictions += all_predictions[strategy_name]\n        ensemble_predictions /= len(all_predictions)\n        \n        # Save ensemble submission\n        sub = pd.read_csv(CFG.sample_sub_path)\n        sub[\"prediction\"] = ensemble_predictions\n        sub.to_csv(\"submission_ensemble_enhanced.csv\", index=False)\n        print(\"\\nSaved submission_ensemble_enhanced.csv\")\n        \n        # Create comprehensive comparison visualization\n        fig, axes = plt.subplots(2, 3, figsize=(20, 12))\n        axes = axes.ravel()\n        \n        for idx, (name, scores_df) in enumerate(all_scores.items()):\n            if idx < len(axes):\n                mean_scores = scores_df.mean().sort_values(ascending=False)\n                \n                ax = axes[idx]\n                colors = plt.cm.viridis(np.linspace(0, 1, len(mean_scores)))\n                bars = ax.barh(range(len(mean_scores)), mean_scores.values, color=colors)\n                ax.set_yticks(range(len(mean_scores)))\n                ax.set_yticklabels(mean_scores.index)\n                ax.set_title(f\"{name}\", fontsize=10)\n                ax.set_xlabel(\"Pearson Correlation\")\n                \n                # Add value labels\n                for i, (value, bar) in enumerate(zip(mean_scores.values, bars)):\n                    ax.text(value + 0.001, bar.get_y() + bar.get_height()/2, \n                           f'{value:.4f}', va='center', fontsize=8)\n        \n        # Hide unused subplots\n        for idx in range(len(all_scores), len(axes)):\n            axes[idx].set_visible(False)\n        \n        plt.tight_layout()\n        plt.savefig(\"enhanced_compression_grid_results.png\", dpi=300, bbox_inches='tight')\n        plt.show()\n        \n        # Summary table\n        print(\"\\n\" + \"=\"*80)\n        print(\"ENHANCED COMPRESSION GRID SEARCH RESULTS\")\n        print(\"=\"*80)\n        \n        summary_data = []\n        for (feature_comp, label_comp, use_adaptive, name) in compression_grid:\n            if name in all_scores:\n                scores_df = all_scores[name]\n                mean_scores = scores_df.mean()\n                \n                summary_data.append({\n                    'Strategy': name,\n                    'Feature_Comp': feature_comp,\n                    'Label_Comp': label_comp,\n                    'Adaptive': use_adaptive,\n                    'Features_Used': len(all_features_used.get(name, [])),\n                    'LightGBM_gbdt': mean_scores.get('LightGBM (gbdt)', 0),\n                    'LightGBM_goss': mean_scores.get('LightGBM (goss)', 0),\n                    'XGBoost': mean_scores.get('XGBoost', 0),\n                    'Average': mean_scores.mean()\n                })\n        \n        summary_df = pd.DataFrame(summary_data)\n        summary_df = summary_df.sort_values('Average', ascending=False)\n        print(summary_df.to_string(index=False))\n        \n        # Save summary\n        summary_df.to_csv('enhanced_compression_grid_search_summary.csv', index=False)\n        \n        # Feature usage analysis\n        print(\"\\n\" + \"=\"*80)\n        print(\"FEATURE USAGE ANALYSIS\")\n        print(\"=\"*80)\n        \n        for name, features in all_features_used.items():\n            print(f\"\\n{name}: {len(features)} features\")\n            # Count feature types\n            top_features = [f for f in features if f in ['X759', 'X752', 'X298', 'X29', 'X586']]\n            microstructure = [f for f in features if any(x in f for x in ['qty', 'flow', 'liquidity', 'spread', 'volume'])]\n            print(f\"  - Top ranked features: {len(top_features)}\")\n            print(f\"  - Microstructure features: {len(microstructure)}\")\n        \n        print(\"\\n\" + \"=\"*80)\n        print(\"Enhanced pipeline completed!\")\n        print(\"Generated files:\")\n        print(\"- Individual submissions for each strategy\")\n        print(\"- submission_ensemble_enhanced.csv\")\n        print(\"- Feature importance and compression profiles\")\n        print(\"- Visualization files\")\n        print(\"- enhanced_compression_grid_search_summary.csv\")\n        print(\"=\"*80)","metadata":{"trusted":true},"outputs":[],"execution_count":null}]}