{"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":11418275,"sourceType":"competition"}],"dockerImageVersionId":31041,"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 numpy as np\nimport pandas as pd\nimport torch\nimport torch.nn as nn\nimport torch.optim as optim\nfrom torch.utils.data import Dataset, DataLoader, TensorDataset\nfrom sklearn.model_selection import KFold, TimeSeriesSplit\nfrom sklearn.preprocessing import StandardScaler, RobustScaler, QuantileTransformer\nfrom sklearn.linear_model import Ridge\nfrom sklearn.ensemble import HistGradientBoostingRegressor\nimport lightgbm as lgb\nimport time\nimport gc\nimport warnings\nfrom itertools import combinations\nwarnings.filterwarnings('ignore')\n\n# Check if CUDA is available\ndevice = torch.device('cuda' if torch.cuda.is_available() else 'cpu')\nprint(f\"Using device: {device}\")\n\nclass WideNeuralNetwork(nn.Module):\n    \"\"\"Wide neural network with skip connections and regularization\"\"\"\n    \n    def __init__(self, input_dim, hidden_dims=[2048, 1024, 512], dropout_rate=0.3):\n        super(WideNeuralNetwork, self).__init__()\n        \n        self.input_layer = nn.Linear(input_dim, hidden_dims[0])\n        self.input_bn = nn.BatchNorm1d(hidden_dims[0])\n        \n        # Build hidden layers\n        self.hidden_layers = nn.ModuleList()\n        self.batch_norms = nn.ModuleList()\n        self.dropouts = nn.ModuleList()\n        \n        for i in range(len(hidden_dims) - 1):\n            self.hidden_layers.append(nn.Linear(hidden_dims[i], hidden_dims[i+1]))\n            self.batch_norms.append(nn.BatchNorm1d(hidden_dims[i+1]))\n            self.dropouts.append(nn.Dropout(dropout_rate))\n        \n        # Skip connection from input to final hidden layer\n        self.skip_connection = nn.Linear(input_dim, hidden_dims[-1])\n        \n        # Output layer\n        self.output_layer = nn.Linear(hidden_dims[-1], 1)\n        \n        # Activation\n        self.activation = nn.GELU()\n        \n    def forward(self, x):\n        # Input layer\n        out = self.input_layer(x)\n        out = self.input_bn(out)\n        out = self.activation(out)\n        \n        # Hidden layers\n        for i, (layer, bn, dropout) in enumerate(zip(self.hidden_layers, self.batch_norms, self.dropouts)):\n            out = layer(out)\n            out = bn(out)\n            out = self.activation(out)\n            out = dropout(out)\n        \n        # Add skip connection\n        skip = self.skip_connection(x)\n        out = out + skip\n        \n        # Output\n        out = self.output_layer(out)\n        return out.squeeze()\n\nclass TabularResNet(nn.Module):\n    \"\"\"ResNet-style architecture for tabular data\"\"\"\n    \n    def __init__(self, input_dim, hidden_dim=256, num_blocks=6, dropout_rate=0.2):\n        super(TabularResNet, self).__init__()\n        \n        self.input_layer = nn.Linear(input_dim, hidden_dim)\n        self.input_bn = nn.BatchNorm1d(hidden_dim)\n        \n        # Residual blocks\n        self.blocks = nn.ModuleList()\n        for _ in range(num_blocks):\n            block = nn.Sequential(\n                nn.Linear(hidden_dim, hidden_dim),\n                nn.BatchNorm1d(hidden_dim),\n                nn.ReLU(),\n                nn.Dropout(dropout_rate),\n                nn.Linear(hidden_dim, hidden_dim),\n                nn.BatchNorm1d(hidden_dim)\n            )\n            self.blocks.append(block)\n        \n        self.activation = nn.ReLU()\n        self.output_layer = nn.Linear(hidden_dim, 1)\n        \n    def forward(self, x):\n        out = self.input_layer(x)\n        out = self.input_bn(out)\n        out = self.activation(out)\n        \n        for block in self.blocks:\n            residual = out\n            out = block(out)\n            out = out + residual\n            out = self.activation(out)\n        \n        out = self.output_layer(out)\n        return out.squeeze()\n\nclass CryptoNeuralNetworkPipeline:\n    \"\"\"Complete pipeline with feature engineering and neural networks for DRW crypto data\"\"\"\n    \n    def __init__(self):\n        self.feature_importance = {}\n        self.feature_stats = {}\n        self.models = []\n        self.nn_models = []\n        self.device = device\n        \n        # LightGBM parameters\n        self.lgb_params = {\n            'objective': 'regression',\n            'metric': 'rmse',\n            'boosting_type': 'gbdt',\n            'num_leaves': 20,\n            'learning_rate': 0.03,\n            'feature_fraction': 0.6,\n            'bagging_fraction': 0.6,\n            'bagging_freq': 5,\n            'min_child_samples': 50,\n            'reg_alpha': 0.5,\n            'reg_lambda': 0.5,\n            'min_gain_to_split': 0.01,\n            'max_depth': 6,\n            'verbose': -1,\n            'random_state': 42,\n            'n_jobs': -1\n        }\n    \n    def print_header(self, text, char=\"=\", width=80):\n        \"\"\"Print formatted header\"\"\"\n        print(f\"\\n{char*width}\")\n        print(f\"{text.center(width)}\")\n        print(f\"{char*width}\")\n    \n    def load_and_optimize_data(self, n_rows=None):\n        \"\"\"Load DRW crypto data with memory optimization\"\"\"\n        self.print_header(\"DATA LOADING AND OPTIMIZATION\")\n        \n        # Load data\n        print(f\"\\nLoading {'all' if n_rows is None else f'{n_rows:,}'} rows...\")\n        \n        if n_rows:\n            train = pd.read_parquet(\"/kaggle/input/drw-crypto-market-prediction/train.parquet\").tail(n_rows)\n        else:\n            train = pd.read_parquet(\"/kaggle/input/drw-crypto-market-prediction/train.parquet\")\n            \n        test = pd.read_parquet(\"/kaggle/input/drw-crypto-market-prediction/test.parquet\")\n        \n        # Reset index for train to avoid issues\n        train = train.reset_index(drop=True)\n        \n        # Memory optimization\n        print(\"\\nOptimizing memory usage...\")\n        for df in [train, test]:\n            float_cols = df.select_dtypes(include=['float64']).columns\n            df[float_cols] = df[float_cols].astype('float32')\n            \n            int_cols = df.select_dtypes(include=['int64']).columns\n            for col in int_cols:\n                if col != 'timestamp' and df[col].max() < 2**31:\n                    df[col] = df[col].astype('int32')\n        \n        # Display data info\n        print(f\"\\nDataset shapes:\")\n        print(f\"  Train: {train.shape} ({train.memory_usage().sum()/1024**2:.1f} MB)\")\n        print(f\"  Test: {test.shape} ({test.memory_usage().sum()/1024**2:.1f} MB)\")\n        \n        # Analyze target distribution\n        print(f\"\\nTarget statistics:\")\n        print(f\"  Mean: {train['label'].mean():.6f}\")\n        print(f\"  Std: {train['label'].std():.6f}\")\n        print(f\"  Skew: {train['label'].skew():.3f}\")\n        print(f\"  Min: {train['label'].min():.3f}\")\n        print(f\"  Max: {train['label'].max():.3f}\")\n        \n        return train, test\n    \n    def create_market_microstructure_features(self, df):\n        \"\"\"Create comprehensive market microstructure features\"\"\"\n        print(\"\\nCreating market microstructure features...\")\n        \n        eps = 1e-8\n        features_created = []\n        \n        # 1. Basic spreads\n        df['spread'] = df['ask_qty'] - df['bid_qty']\n        df['spread_ratio'] = df['spread'] / (df['bid_qty'] + df['ask_qty'] + eps)\n        df['spread_pct'] = df['spread'] / (df['bid_qty'] + eps)\n        df['normalized_spread'] = df['spread'] / (df['volume'] + eps)\n        features_created.extend(['spread', 'spread_ratio', 'spread_pct', 'normalized_spread'])\n        \n        # 2. Order flow metrics\n        df['order_flow'] = df['buy_qty'] - df['sell_qty']\n        df['order_flow_imbalance'] = df['order_flow'] / (df['buy_qty'] + df['sell_qty'] + eps)\n        df['order_pressure'] = df['order_flow'] / (df['volume'] + eps)\n        df['net_order_flow'] = df['order_flow'] / (df['bid_qty'] + df['ask_qty'] + eps)\n        df['order_flow_ratio'] = df['buy_qty'] / (df['sell_qty'] + eps)\n        features_created.extend(['order_flow', 'order_flow_imbalance', 'order_pressure', \n                               'net_order_flow', 'order_flow_ratio'])\n        \n        # 3. Volume analysis\n        df['volume_ratio'] = df['volume'] / (df['bid_qty'] + df['ask_qty'] + eps)\n        df['trade_intensity'] = (df['buy_qty'] + df['sell_qty']) / (df['volume'] + eps)\n        df['volume_imbalance'] = abs(df['buy_qty'] - df['sell_qty']) / (df['volume'] + eps)\n        df['relative_volume'] = df['volume'] / (df['bid_qty'] + df['ask_qty'] + df['buy_qty'] + df['sell_qty'] + eps)\n        df['volume_per_trade'] = df['volume'] / ((df['buy_qty'] + df['sell_qty']) / 2 + eps)\n        features_created.extend(['volume_ratio', 'trade_intensity', 'volume_imbalance', \n                               'relative_volume', 'volume_per_trade'])\n        \n        # 4. Liquidity metrics\n        df['liquidity'] = df['bid_qty'] + df['ask_qty']\n        df['liquidity_imbalance'] = (df['bid_qty'] - df['ask_qty']) / (df['liquidity'] + eps)\n        df['liquidity_consumption'] = df['volume'] / (df['liquidity'] + eps)\n        df['liquidity_ratio'] = df['liquidity'] / (df['volume'] + eps)\n        df['liquidity_adjusted_volume'] = df['volume'] * df['liquidity_ratio']\n        features_created.extend(['liquidity', 'liquidity_imbalance', 'liquidity_consumption', \n                               'liquidity_ratio', 'liquidity_adjusted_volume'])\n        \n        # 5. Market depth\n        df['bid_depth'] = df['bid_qty'] / (df['volume'] + eps)\n        df['ask_depth'] = df['ask_qty'] / (df['volume'] + eps)\n        df['market_depth'] = df['liquidity'] / (df['volume'] + eps)\n        df['depth_imbalance'] = (df['bid_depth'] - df['ask_depth']) / (df['bid_depth'] + df['ask_depth'] + eps)\n        df['weighted_depth'] = df['market_depth'] * df['volume_ratio']\n        features_created.extend(['bid_depth', 'ask_depth', 'market_depth', \n                               'depth_imbalance', 'weighted_depth'])\n        \n        # 6. Trading aggressiveness\n        df['buy_aggressiveness'] = df['buy_qty'] / (df['ask_qty'] + eps)\n        df['sell_aggressiveness'] = df['sell_qty'] / (df['bid_qty'] + eps)\n        df['net_aggressiveness'] = (df['buy_aggressiveness'] - df['sell_aggressiveness']) / \\\n                                   (df['buy_aggressiveness'] + df['sell_aggressiveness'] + eps)\n        df['aggressiveness_ratio'] = df['buy_aggressiveness'] / (df['sell_aggressiveness'] + eps)\n        features_created.extend(['buy_aggressiveness', 'sell_aggressiveness', \n                               'net_aggressiveness', 'aggressiveness_ratio'])\n        \n        # 7. Price pressure indicators\n        df['buy_pressure'] = df['buy_qty'] / (df['ask_qty'] + eps)\n        df['sell_pressure'] = df['sell_qty'] / (df['bid_qty'] + eps)\n        df['pressure_ratio'] = df['buy_pressure'] / (df['sell_pressure'] + eps)\n        df['net_pressure'] = (df['buy_pressure'] - df['sell_pressure']) / \\\n                            (df['buy_pressure'] + df['sell_pressure'] + eps)\n        features_created.extend(['buy_pressure', 'sell_pressure', 'pressure_ratio', 'net_pressure'])\n        \n        # 8. Advanced microstructure metrics\n        df['kyle_lambda'] = np.abs(df['order_pressure']) / (np.sqrt(df['volume']) + eps)\n        df['amihud_illiquidity'] = np.abs(df['order_flow_imbalance']) / (df['volume'] + eps)\n        df['realized_spread'] = np.abs(df['spread']) * df['volume']\n        df['effective_spread'] = df['spread_ratio'] * df['liquidity_ratio']\n        df['price_impact'] = df['kyle_lambda'] * df['net_aggressiveness']\n        features_created.extend(['kyle_lambda', 'amihud_illiquidity', 'realized_spread', \n                               'effective_spread', 'price_impact'])\n        \n        # 9. Log transforms for stability\n        df['log_volume'] = np.log1p(df['volume'])\n        df['log_liquidity'] = np.log1p(df['liquidity'])\n        df['log_spread'] = np.log1p(np.abs(df['spread']))\n        df['log_order_flow'] = np.sign(df['order_flow']) * np.log1p(np.abs(df['order_flow']))\n        df['log_bid_qty'] = np.log1p(df['bid_qty'])\n        df['log_ask_qty'] = np.log1p(df['ask_qty'])\n        features_created.extend(['log_volume', 'log_liquidity', 'log_spread', \n                               'log_order_flow', 'log_bid_qty', 'log_ask_qty'])\n        \n        # 10. Ratios and composite features\n        df['bid_ask_ratio'] = df['bid_qty'] / (df['ask_qty'] + eps)\n        df['buy_sell_ratio'] = df['buy_qty'] / (df['sell_qty'] + eps)\n        df['volume_liquidity_interaction'] = df['volume_ratio'] * df['liquidity_imbalance']\n        df['depth_pressure_interaction'] = df['market_depth'] * df['net_pressure']\n        features_created.extend(['bid_ask_ratio', 'buy_sell_ratio', \n                               'volume_liquidity_interaction', 'depth_pressure_interaction'])\n        \n        print(f\"  Created {len(features_created)} market microstructure features\")\n        return features_created\n    \n    def analyze_and_select_x_features(self, train, n_features=50):\n        \"\"\"Analyze and select best X features\"\"\"\n        self.print_header(\"X FEATURE ANALYSIS\", \"-\")\n        \n        # Get X columns\n        x_cols = [col for col in train.columns if col.startswith('X') and col[1:].isdigit()]\n        print(f\"\\nAnalyzing {len(x_cols)} X features...\")\n        \n        # Use sufficient sample for analysis\n        sample_size = min(30000, len(train))\n        train_sample = train.sample(n=sample_size, random_state=42)\n        target = train_sample['label']\n        \n        # Analyze features\n        feature_metrics = {}\n        \n        print(\"\\nCalculating feature metrics...\")\n        for i, col in enumerate(x_cols):\n            if i % 100 == 0 and i > 0:\n                print(f\"  Processed {i}/{len(x_cols)} features...\")\n            \n            try:\n                feature = train_sample[col]\n                \n                # Calculate various metrics\n                metrics = {\n                    'pearson_corr': abs(feature.corr(target)),\n                    'spearman_corr': abs(feature.corr(target, method='spearman')),\n                    'variance': feature.var(),\n                    'std': feature.std(),\n                    'unique_ratio': feature.nunique() / len(feature),\n                    'zero_ratio': (feature == 0).sum() / len(feature),\n                    'missing_ratio': feature.isna().sum() / len(feature)\n                }\n                \n                # Skip features with too many zeros or missing values\n                if metrics['zero_ratio'] > 0.95 or metrics['missing_ratio'] > 0.5:\n                    continue\n                \n                # Composite score emphasizing correlation\n                metrics['score'] = (\n                    metrics['pearson_corr'] * 0.45 +\n                    metrics['spearman_corr'] * 0.45 +\n                    min(metrics['variance'], 1) * 0.05 +\n                    (1 - metrics['zero_ratio']) * 0.05\n                )\n                \n                feature_metrics[col] = metrics\n                \n            except:\n                pass\n        \n        # Sort by score\n        sorted_features = sorted(feature_metrics.items(), key=lambda x: x[1]['score'], reverse=True)\n        \n        # Display top features\n        print(f\"\\nTop 20 X features by score:\")\n        print(f\"{'Rank':<6} {'Feature':<20} {'Score':<8} {'Pearson':<8} {'Spearman':<8}\")\n        print(\"-\" * 60)\n        \n        for i, (feat, metrics) in enumerate(sorted_features[:20]):\n            print(f\"{i+1:<6} {feat:<20} {metrics['score']:<8.4f} \"\n                  f\"{metrics['pearson_corr']:<8.4f} {metrics['spearman_corr']:<8.4f}\")\n        \n        # Select top features\n        selected_features = [feat for feat, _ in sorted_features[:n_features]]\n        \n        # Store metrics for later use\n        self.feature_stats = {feat: metrics for feat, metrics in feature_metrics.items()}\n        \n        print(f\"\\nSelected {len(selected_features)} X features\")\n        \n        return selected_features\n    \n    def create_feature_combinations(self, train, base_features, x_features, max_combinations=200):\n        \"\"\"Create powerful feature combinations\"\"\"\n        self.print_header(\"FEATURE COMBINATION DISCOVERY\", \"-\")\n        \n        # Sample for speed\n        sample_size = min(15000, len(train))\n        train_sample = train.sample(n=sample_size, random_state=42)\n        target = train_sample['label']\n        \n        all_combinations = []\n        \n        # Select best features for combinations\n        top_base = base_features[:15]\n        top_x = x_features[:15]\n        \n        # Test 2-way combinations\n        print(\"\\nTesting 2-way combinations...\")\n        tested = 0\n        \n        # Base × X combinations\n        for base_feat in top_base[:10]:\n            for x_feat in top_x[:10]:\n                if tested >= max_combinations // 2:\n                    break\n                    \n                try:\n                    # Test multiplication\n                    combo_mult = train_sample[base_feat] * train_sample[x_feat]\n                    corr_mult = abs(combo_mult.corr(target))\n                    \n                    # Test with log transform\n                    log_combo = np.log1p(np.abs(combo_mult))\n                    corr_log = abs(log_combo.corr(target))\n                    \n                    if max(corr_mult, corr_log) > 0.02:\n                        all_combinations.append({\n                            'name': f\"{base_feat}__X__{x_feat}\",\n                            'features': (base_feat, x_feat),\n                            'correlation': max(corr_mult, corr_log),\n                            'use_log': corr_log > corr_mult,\n                            'type': 'multiply'\n                        })\n                    \n                    tested += 1\n                    \n                except:\n                    pass\n        \n        # Base × Base combinations\n        for i, f1 in enumerate(top_base[:8]):\n            for f2 in top_base[i+1:12]:\n                if tested >= max_combinations:\n                    break\n                    \n                try:\n                    # Test multiplication\n                    combo = train_sample[f1] * train_sample[f2]\n                    corr = abs(combo.corr(target))\n                    \n                    if corr > 0.02:\n                        all_combinations.append({\n                            'name': f\"{f1}__X__{f2}\",\n                            'features': (f1, f2),\n                            'correlation': corr,\n                            'use_log': False,\n                            'type': 'multiply'\n                        })\n                    \n                    # Test division\n                    if train_sample[f2].min() > 0:\n                        combo_div = train_sample[f1] / train_sample[f2]\n                        corr_div = abs(combo_div.corr(target))\n                        \n                        if corr_div > 0.02:\n                            all_combinations.append({\n                                'name': f\"{f1}__DIV__{f2}\",\n                                'features': (f1, f2),\n                                'correlation': corr_div,\n                                'use_log': False,\n                                'type': 'divide'\n                            })\n                    \n                    tested += 2\n                    \n                except:\n                    pass\n        \n        # Sort by correlation\n        all_combinations.sort(key=lambda x: x['correlation'], reverse=True)\n        \n        print(f\"\\nTested ~{tested} combinations\")\n        print(f\"Found {len(all_combinations)} useful combinations\")\n        \n        # Display top combinations\n        print(f\"\\nTop 20 combinations:\")\n        print(f\"{'Rank':<6} {'Correlation':<12} {'Type':<10} {'Name'}\")\n        print(\"-\" * 60)\n        \n        for i, combo in enumerate(all_combinations[:20]):\n            log_marker = \"L\" if combo.get('use_log', False) else \" \"\n            print(f\"{i+1:<6} {combo['correlation']:<12.4f} {combo['type']:<10} \"\n                  f\"{log_marker} {combo['name'][:40]}\")\n        \n        return all_combinations[:20]\n    \n    def create_combination_features(self, df, combinations):\n        \"\"\"Create the selected combination features\"\"\"\n        print(f\"\\nCreating {len(combinations)} combination features...\")\n        \n        created_features = []\n        \n        for combo in combinations:\n            try:\n                if combo['type'] == 'multiply':\n                    f1, f2 = combo['features']\n                    values = df[f1] * df[f2]\n                    \n                    if combo.get('use_log', False):\n                        df[combo['name']] = np.log1p(np.abs(values))\n                    else:\n                        df[combo['name']] = values\n                        \n                elif combo['type'] == 'divide':\n                    f1, f2 = combo['features']\n                    df[combo['name']] = df[f1] / (df[f2] + 1e-8)\n                \n                created_features.append(combo['name'])\n                \n            except Exception as e:\n                continue\n        \n        print(f\"  Successfully created {len(created_features)} features\")\n        return created_features\n    \n    def train_lightgbm_model(self, X_train, y_train, X_val=None, y_val=None):\n        \"\"\"Train a LightGBM model\"\"\"\n        train_data = lgb.Dataset(X_train, label=y_train)\n        \n        if X_val is not None and y_val is not None:\n            val_data = lgb.Dataset(X_val, label=y_val, reference=train_data)\n            \n            model = lgb.train(\n                self.lgb_params,\n                train_data,\n                num_boost_round=500,\n                valid_sets=[val_data],\n                callbacks=[lgb.early_stopping(30), lgb.log_evaluation(0)]\n            )\n        else:\n            model = lgb.train(\n                self.lgb_params,\n                train_data,\n                num_boost_round=200\n            )\n        \n        return model\n    \n    def train_gb_ensemble(self, X_train, y_train):\n        \"\"\"Train an ensemble of gradient boosting models\"\"\"\n        print(\"\\nTraining gradient boosting ensemble...\")\n        \n        models = []\n        \n        # LightGBM configurations\n        lgb_configs = [\n            {**self.lgb_params, 'num_leaves': 15, 'learning_rate': 0.02},\n            {**self.lgb_params, 'num_leaves': 25, 'learning_rate': 0.01, 'reg_alpha': 1.0},\n            {**self.lgb_params, 'num_leaves': 10, 'learning_rate': 0.03, 'max_depth': 5},\n        ]\n        \n        for i, params in enumerate(lgb_configs):\n            print(f\"  Training LightGBM model {i+1}/{len(lgb_configs)}\")\n            train_data = lgb.Dataset(X_train, label=y_train)\n            model = lgb.train(params, train_data, num_boost_round=200)\n            models.append(('lgb', model))\n        \n        # HistGradientBoosting\n        print(\"  Training HistGradientBoosting\")\n        hgb = HistGradientBoostingRegressor(\n            max_iter=100,\n            learning_rate=0.02,\n            max_depth=5,\n            l2_regularization=1.0,\n            random_state=42\n        )\n        hgb.fit(X_train, y_train)\n        models.append(('hgb', hgb))\n        \n        print(f\"  Total GB ensemble size: {len(models)} models\")\n        return models\n    \n    def generate_gb_predictions(self, models, X_test):\n        \"\"\"Generate predictions from gradient boosting ensemble\"\"\"\n        predictions = []\n        \n        for model_type, model in models:\n            if model_type == 'lgb':\n                pred = model.predict(X_test, num_iteration=model.best_iteration)\n            elif model_type == 'hgb':\n                pred = model.predict(X_test)\n            \n            predictions.append(pred)\n        \n        # Simple average\n        final_predictions = np.mean(predictions, axis=0)\n        \n        return final_predictions\n    \n    def train_neural_network(self, model, train_loader, val_loader=None, \n                           epochs=100, lr=0.001, patience=15):\n        \"\"\"Train a neural network model\"\"\"\n        model = model.to(self.device)\n        optimizer = optim.AdamW(model.parameters(), lr=lr, weight_decay=1e-5)\n        scheduler = optim.lr_scheduler.ReduceLROnPlateau(optimizer, mode='min', \n                                                         factor=0.5, patience=5)\n        criterion = nn.MSELoss()\n        \n        best_val_loss = float('inf')\n        best_model_state = None\n        early_stop_counter = 0\n        \n        for epoch in range(epochs):\n            # Training\n            model.train()\n            train_loss = 0\n            for batch_X, batch_y in train_loader:\n                optimizer.zero_grad()\n                predictions = model(batch_X)\n                loss = criterion(predictions, batch_y)\n                loss.backward()\n                torch.nn.utils.clip_grad_norm_(model.parameters(), 1.0)\n                optimizer.step()\n                train_loss += loss.item()\n            \n            train_loss /= len(train_loader)\n            \n            # Validation\n            if val_loader:\n                model.eval()\n                val_loss = 0\n                val_predictions = []\n                val_targets = []\n                \n                with torch.no_grad():\n                    for batch_X, batch_y in val_loader:\n                        predictions = model(batch_X)\n                        loss = criterion(predictions, batch_y)\n                        val_loss += loss.item()\n                        val_predictions.extend(predictions.cpu().numpy())\n                        val_targets.extend(batch_y.cpu().numpy())\n                \n                val_loss /= len(val_loader)\n                val_corr = np.corrcoef(val_predictions, val_targets)[0, 1]\n                \n                scheduler.step(val_loss)\n                \n                # Early stopping\n                if val_loss < best_val_loss:\n                    best_val_loss = val_loss\n                    best_model_state = model.state_dict().copy()\n                    early_stop_counter = 0\n                else:\n                    early_stop_counter += 1\n                \n                if early_stop_counter >= patience:\n                    print(f\"    Early stopping at epoch {epoch+1}\")\n                    break\n                \n                if (epoch + 1) % 20 == 0:\n                    print(f\"    Epoch {epoch+1}: Train Loss={train_loss:.4f}, \"\n                          f\"Val Loss={val_loss:.4f}, Val Corr={val_corr:.4f}\")\n        \n        # Load best model\n        if best_model_state:\n            model.load_state_dict(best_model_state)\n        \n        return model\n    \n    def create_dataloaders(self, X_train, y_train, X_val=None, y_val=None, batch_size=2048):\n        \"\"\"Create PyTorch dataloaders\"\"\"\n        train_dataset = TensorDataset(\n            torch.FloatTensor(X_train).to(self.device),\n            torch.FloatTensor(y_train).to(self.device)\n        )\n        train_loader = DataLoader(train_dataset, batch_size=batch_size, shuffle=True)\n        \n        val_loader = None\n        if X_val is not None and y_val is not None:\n            val_dataset = TensorDataset(\n                torch.FloatTensor(X_val).to(self.device),\n                torch.FloatTensor(y_val).to(self.device)\n            )\n            val_loader = DataLoader(val_dataset, batch_size=batch_size, shuffle=False)\n        \n        return train_loader, val_loader\n    \n    def train_neural_network_ensemble(self, X_train, y_train):\n        \"\"\"Train final ensemble of neural networks\"\"\"\n        self.print_header(\"TRAINING NEURAL NETWORK ENSEMBLE\", \"-\")\n        \n        final_models = []\n        \n        # Different preprocessing strategies\n        scalers = [\n            StandardScaler(),\n            RobustScaler(),\n            QuantileTransformer(n_quantiles=100, output_distribution='normal')\n        ]\n        \n        # Train multiple models with different preprocessing\n        for scaler_idx, scaler in enumerate(scalers):\n            print(f\"\\nTraining with scaler {scaler_idx+1}/{len(scalers)}\")\n            \n            X_scaled = scaler.fit_transform(X_train)\n            train_loader, _ = self.create_dataloaders(X_scaled, y_train, batch_size=4096)\n            \n            # Create diverse models\n            if scaler_idx == 0:  # Full ensemble for first scaler\n                models = [\n                    WideNeuralNetwork(X_train.shape[1], hidden_dims=[4096, 2048, 1024], dropout_rate=0.25),\n                    WideNeuralNetwork(X_train.shape[1], hidden_dims=[3072, 1536, 768, 384], dropout_rate=0.3),\n                    TabularResNet(X_train.shape[1], hidden_dim=512, num_blocks=10, dropout_rate=0.2),\n                ]\n            elif scaler_idx == 1:  # Different architecture for second scaler\n                models = [\n                    WideNeuralNetwork(X_train.shape[1], hidden_dims=[2048, 1024, 512], dropout_rate=0.35),\n                    TabularResNet(X_train.shape[1], hidden_dim=256, num_blocks=8, dropout_rate=0.25),\n                ]\n            else:  # Simple model for third scaler\n                models = [\n                    WideNeuralNetwork(X_train.shape[1], hidden_dims=[1536, 768], dropout_rate=0.3),\n                ]\n            \n            for model in models:\n                print(f\"  Training {model.__class__.__name__}\")\n                trained_model = self.train_neural_network(\n                    model, train_loader, epochs=60, lr=0.001, patience=10\n                )\n                final_models.append((trained_model, scaler))\n        \n        print(f\"\\nTotal NN ensemble size: {len(final_models)} neural networks\")\n        \n        return final_models\n    \n    def generate_nn_predictions(self, models, X_test):\n        \"\"\"Generate predictions from neural network ensemble\"\"\"\n        predictions = []\n        \n        for model, scaler in models:\n            # Scale features\n            X_scaled = scaler.transform(X_test)\n            \n            model.eval()\n            with torch.no_grad():\n                X_tensor = torch.FloatTensor(X_scaled).to(self.device)\n                pred = model(X_tensor).cpu().numpy()\n            \n            predictions.append(pred)\n        \n        # Weighted average\n        final_predictions = np.mean(predictions, axis=0)\n        \n        return final_predictions\n    \n    def evaluate_models(self, train, features, n_folds=3):\n        \"\"\"Evaluate models using cross-validation\"\"\"\n        self.print_header(\"MODEL EVALUATION\", \"-\")\n        \n        X = train[features].values\n        y = train['label'].values\n        \n        kf = KFold(n_splits=n_folds, shuffle=True, random_state=42)\n        \n        gb_scores = []\n        nn_scores = []\n        \n        print(\"\\nCross-Validation:\")\n        for fold, (train_idx, val_idx) in enumerate(kf.split(X)):\n            print(f\"\\nFold {fold+1}/{n_folds}\")\n            \n            X_train, X_val = X[train_idx], X[val_idx]\n            y_train, y_val = y[train_idx], y[val_idx]\n            \n            # Evaluate gradient boosting\n            print(\"  Evaluating gradient boosting...\")\n            gb_model = self.train_lightgbm_model(X_train, y_train, X_val, y_val)\n            gb_pred = gb_model.predict(X_val, num_iteration=gb_model.best_iteration)\n            gb_score = np.corrcoef(y_val, gb_pred)[0, 1]\n            gb_scores.append(gb_score)\n            print(f\"    GB Score: {gb_score:.4f}\")\n            \n            # Evaluate neural network\n            print(\"  Evaluating neural network...\")\n            scaler = RobustScaler()\n            X_train_scaled = scaler.fit_transform(X_train)\n            X_val_scaled = scaler.transform(X_val)\n            \n            train_loader, val_loader = self.create_dataloaders(\n                X_train_scaled, y_train, X_val_scaled, y_val, batch_size=2048\n            )\n            \n            nn_model = WideNeuralNetwork(X_train.shape[1], hidden_dims=[2048, 1024, 512])\n            trained_nn = self.train_neural_network(\n                nn_model, train_loader, val_loader, epochs=40, lr=0.001, patience=10\n            )\n            \n            trained_nn.eval()\n            with torch.no_grad():\n                X_val_tensor = torch.FloatTensor(X_val_scaled).to(self.device)\n                nn_pred = trained_nn(X_val_tensor).cpu().numpy()\n            \n            nn_score = np.corrcoef(y_val, nn_pred)[0, 1]\n            nn_scores.append(nn_score)\n            print(f\"    NN Score: {nn_score:.4f}\")\n        \n        print(f\"\\nOverall GB CV Score: {np.mean(gb_scores):.4f} (±{np.std(gb_scores):.4f})\")\n        print(f\"Overall NN CV Score: {np.mean(nn_scores):.4f} (±{np.std(nn_scores):.4f})\")\n        \n        return np.mean(gb_scores), np.mean(nn_scores)\n    \n    def run_complete_pipeline(self, n_rows=None, blend_ratio=0.7):\n        \"\"\"Run the complete pipeline with DRW crypto data\"\"\"\n        self.print_header(\"DRW CRYPTO NEURAL NETWORK PIPELINE\", \"=\")\n        \n        start_time = time.time()\n        \n        # 1. Load and optimize data\n        train, test = self.load_and_optimize_data(n_rows)\n        \n        # 2. Create market microstructure features\n        self.print_header(\"FEATURE ENGINEERING\", \"-\")\n        market_features = self.create_market_microstructure_features(train)\n        self.create_market_microstructure_features(test)\n        \n        # 3. Analyze and select X features\n        x_features = self.analyze_and_select_x_features(train, n_features=50)\n        \n        # 4. Create feature combinations\n        base_features_for_combo = [\n            'order_flow_imbalance', 'net_order_flow', 'volume_ratio',\n            'liquidity_ratio', 'kyle_lambda', 'spread_ratio',\n            'net_pressure', 'buy_pressure'\n        ]\n        \n        combinations = self.create_feature_combinations(\n            train, base_features_for_combo, x_features, max_combinations=200\n        )\n        \n        # 5. Create combination features\n        combo_features = self.create_combination_features(train, combinations)\n        self.create_combination_features(test, combinations)\n        \n        # 6. Assemble final feature set\n        all_features = (\n            ['bid_qty', 'ask_qty', 'buy_qty', 'sell_qty', 'volume'] +\n            market_features +\n            x_features +\n            combo_features\n        )\n        \n        # Remove any duplicates\n        all_features = list(dict.fromkeys(all_features))\n        \n        print(f\"\\nFinal feature set: {len(all_features)} features\")\n        print(f\"  Original: 5\")\n        print(f\"  Market microstructure: {len(market_features)}\")\n        print(f\"  Selected X features: {len(x_features)}\")\n        print(f\"  Combinations: {len(combo_features)}\")\n        \n        # 7. Model evaluation\n        gb_score, nn_score = self.evaluate_models(train, all_features, n_folds=3)\n        \n        # 8. Train final models on all data\n        self.print_header(\"FINAL MODEL TRAINING\", \"-\")\n        \n        X_train = train[all_features].values\n        y_train = train['label'].values\n        X_test = test[all_features].values\n        \n        # Train gradient boosting ensemble\n        print(\"\\nTraining final gradient boosting ensemble...\")\n        gb_models = self.train_gb_ensemble(X_train, y_train)\n        \n        # Train neural network ensemble\n        nn_models = self.train_neural_network_ensemble(X_train, y_train)\n        \n        # 9. Generate predictions\n        print(\"\\nGenerating predictions...\")\n        gb_predictions = self.generate_gb_predictions(gb_models, X_test)\n        nn_predictions = self.generate_nn_predictions(nn_models, X_test)\n        \n        # 10. Blend predictions\n        print(f\"\\nBlending predictions: {blend_ratio*100:.0f}% NN, {(1-blend_ratio)*100:.0f}% GB\")\n        final_predictions = blend_ratio * nn_predictions + (1 - blend_ratio) * gb_predictions\n        \n        # 11. Create submission\n        submission = pd.DataFrame({\n            'id': test.index,\n            'prediction': final_predictions\n        })\n        \n        submission.to_csv('drw_crypto_nn_submission.csv', index=False)\n        \n        # 12. Final summary\n        self.print_header(\"PIPELINE COMPLETE\", \"=\")\n        \n        elapsed_time = (time.time() - start_time) / 60\n        \n        print(f\"\\nExecution Summary:\")\n        print(f\"  Time elapsed: {elapsed_time:.1f} minutes\")\n        print(f\"  GB CV Score: {gb_score:.4f}\")\n        print(f\"  NN CV Score: {nn_score:.4f}\")\n        print(f\"  Total features: {len(all_features)}\")\n        print(f\"  GB ensemble size: {len(gb_models)} models\")\n        print(f\"  NN ensemble size: {len(nn_models)} models\")\n        print(f\"  Blend ratio: {blend_ratio*100:.0f}% NN / {(1-blend_ratio)*100:.0f}% GB\")\n        print(f\"  Output file: drw_crypto_nn_submission.csv\")\n        \n        print(f\"\\nPrediction statistics:\")\n        print(f\"  GB - Mean: {gb_predictions.mean():.6f}, Std: {gb_predictions.std():.6f}\")\n        print(f\"  NN - Mean: {nn_predictions.mean():.6f}, Std: {nn_predictions.std():.6f}\")\n        print(f\"  Final - Mean: {final_predictions.mean():.6f}, Std: {final_predictions.std():.6f}\")\n        \n        # Cleanup\n        gc.collect()\n        torch.cuda.empty_cache() if torch.cuda.is_available() else None\n        \n        return submission, all_features\n\n\n# Main execution\nif __name__ == \"__main__\":\n    # Initialize pipeline\n    pipeline = CryptoNeuralNetworkPipeline()\n    \n    # Run with parameters\n    n_rows = 40000  # Use last 170k rows (about 9 months of data)\n    blend_ratio = 0.7  # 70% neural network, 30% gradient boosting\n    \n    try:\n        # Run complete pipeline\n        submission, features = pipeline.run_complete_pipeline(\n            n_rows=n_rows,\n            blend_ratio=blend_ratio\n        )\n        \n        print(\"\\n\" + \"=\"*80)\n        print(\"SUCCESS: Pipeline completed successfully!\")\n        print(\"=\"*80)\n        \n    except Exception as e:\n        print(f\"\\nError in pipeline: {str(e)}\")\n        import traceback\n        traceback.print_exc()","metadata":{"trusted":true},"outputs":[],"execution_count":null}]}