{"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":"\"\"\"\nDRW Crypto Market Prediction - Complete Working Pipeline\n========================================================\nMemory-optimized neural network solution for crypto prediction\n\"\"\"\n\nimport numpy as np\nimport pandas as pd\nimport torch\nimport torch.nn as nn\nimport torch.optim as optim\nimport torch.nn.functional as F\nfrom torch.utils.data import DataLoader, TensorDataset\nimport gc\nimport warnings\nfrom pathlib import Path\nimport lightgbm as lgb\nfrom sklearn.preprocessing import StandardScaler, QuantileTransformer\nfrom sklearn.model_selection import KFold\nfrom scipy.stats import pearsonr\nimport logging\nimport joblib\n\nwarnings.filterwarnings(\"ignore\")\n\n# Setup logging\nlogging.basicConfig(level=logging.INFO, format='%(asctime)s - %(message)s')\nlogger = logging.getLogger(__name__)\n\n# =========================\n# Configuration\n# =========================\nclass Config:\n    # Paths\n    TRAIN_PATH = \"/kaggle/input/drw-crypto-market-prediction/train.parquet\"\n    TEST_PATH = \"/kaggle/input/drw-crypto-market-prediction/test.parquet\"\n    OUTPUT_DIR = Path(\"/kaggle/working\")\n    \n    # Data settings\n    N_FOLDS = 3\n    TOP_N_FEATURES = 150\n    SAMPLE_SIZE = 50000\n    RANDOM_STATE = 42\n    \n    # Neural Network\n    USE_GPU = torch.cuda.is_available()\n    DEVICE = torch.device(\"cuda\" if USE_GPU else \"cpu\")\n    BATCH_SIZE = 2048\n    HIDDEN_DIMS = [512, 256, 128]\n    DROPOUT_RATE = 0.3\n    LEARNING_RATE = 0.001\n    EPOCHS = 30\n    EARLY_STOPPING_PATIENCE = 10\n\nlogger.info(f\"Using device: {Config.DEVICE}\")\n\n# =========================\n# Memory Optimization\n# =========================\ndef optimize_memory(df):\n    \"\"\"Optimize DataFrame memory usage\"\"\"\n    for col in df.columns:\n        if col != 'timestamp':\n            df[col] = df[col].astype(np.float32)\n    return df\n\n# =========================\n# Feature Engineering\n# =========================\nclass FeatureEngineer:\n    \"\"\"Simple but effective feature engineering\"\"\"\n    \n    def __init__(self):\n        self.feature_cols = None\n        \n    def create_features(self, df):\n        \"\"\"Create market microstructure features\"\"\"\n        logger.info(\"Creating features...\")\n        \n        # Basic market features\n        if all(col in df.columns for col in [\"bid_qty\", \"ask_qty\", \"buy_qty\", \"sell_qty\", \"volume\"]):\n            eps = 1e-10\n            \n            # Core features\n            df['book_imbalance'] = (df['bid_qty'] - df['ask_qty']) / (df['bid_qty'] + df['ask_qty'] + eps)\n            df['net_order_flow'] = df['buy_qty'] - df['sell_qty']\n            df['order_flow_ratio'] = df['net_order_flow'] / (df['volume'] + eps)\n            df['kyle_lambda'] = np.abs(df['net_order_flow']) / (df['volume'] + eps)\n            df['trade_intensity'] = df['volume'] / (df['bid_qty'] + df['ask_qty'] + eps)\n            df['total_depth'] = df['bid_qty'] + df['ask_qty']\n            df['depth_ratio'] = df['bid_qty'] / (df['ask_qty'] + eps)\n            df['vpin_proxy'] = np.abs(df['buy_qty'] - df['sell_qty']) / (df['buy_qty'] + df['sell_qty'] + eps)\n            \n            # Convert to float32\n            for col in ['book_imbalance', 'net_order_flow', 'order_flow_ratio', 'kyle_lambda',\n                       'trade_intensity', 'total_depth', 'depth_ratio', 'vpin_proxy']:\n                df[col] = df[col].astype(np.float32)\n        \n        # Simple X features statistics\n        x_cols = [col for col in df.columns if col.startswith('X_')]\n        if len(x_cols) > 0:\n            # Process in chunks\n            chunk_size = 100\n            for i in range(0, len(x_cols), chunk_size):\n                chunk_cols = x_cols[i:i+chunk_size]\n                df[f'X_mean_{i//chunk_size}'] = df[chunk_cols].mean(axis=1).astype(np.float32)\n                df[f'X_std_{i//chunk_size}'] = df[chunk_cols].std(axis=1).astype(np.float32)\n        \n        # Clean up\n        df = df.replace([np.inf, -np.inf], 0).fillna(0)\n        \n        # Store feature columns\n        self.feature_cols = [col for col in df.columns if col not in ['label', 'timestamp']]\n        \n        logger.info(f\"Created {len(self.feature_cols)} features\")\n        return df\n\n# =========================\n# Feature Selection\n# =========================\nclass FeatureSelector:\n    \"\"\"Select top features using LightGBM importance\"\"\"\n    \n    def __init__(self, n_features=150):\n        self.n_features = n_features\n        self.selected_features = None\n        \n    def fit(self, X, y):\n        \"\"\"Fit the feature selector\"\"\"\n        logger.info(f\"Selecting top {self.n_features} features...\")\n        \n        # Sample data for speed\n        n_samples = min(Config.SAMPLE_SIZE, len(X))\n        idx = np.random.choice(len(X), n_samples, replace=False)\n        \n        # Train LightGBM\n        lgb_data = lgb.Dataset(X.iloc[idx], label=y[idx])\n        \n        params = {\n            'objective': 'regression',\n            'metric': 'rmse',\n            'num_leaves': 31,\n            'learning_rate': 0.1,\n            'feature_fraction': 0.8,\n            'verbose': -1,\n            'force_col_wise': True\n        }\n        \n        model = lgb.train(\n            params,\n            lgb_data,\n            num_boost_round=100,\n            callbacks=[lgb.log_evaluation(0)]\n        )\n        \n        # Get feature importance\n        importance = model.feature_importance(importance_type='gain')\n        feature_importance = pd.DataFrame({\n            'feature': X.columns,\n            'importance': importance\n        }).sort_values('importance', ascending=False)\n        \n        # Select top features\n        self.selected_features = feature_importance.head(self.n_features)['feature'].tolist()\n        \n        logger.info(f\"Selected {len(self.selected_features)} features\")\n        \n        # Clean up\n        del model, lgb_data\n        gc.collect()\n        \n        return self\n    \n    def transform(self, X):\n        \"\"\"Transform data to selected features\"\"\"\n        return X[self.selected_features]\n    \n    def fit_transform(self, X, y):\n        \"\"\"Fit and transform in one step\"\"\"\n        self.fit(X, y)\n        return self.transform(X)\n\n# =========================\n# Neural Network Model\n# =========================\nclass CryptoMLP(nn.Module):\n    \"\"\"Multi-layer perceptron for crypto prediction\"\"\"\n    \n    def __init__(self, input_dim, hidden_dims=[512, 256, 128], dropout_rate=0.3):\n        super().__init__()\n        \n        # Build layers\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(inplace=True),\n                nn.Dropout(dropout_rate)\n            ])\n            prev_dim = hidden_dim\n        \n        # Output layer\n        layers.append(nn.Linear(prev_dim, 1))\n        \n        self.model = nn.Sequential(*layers)\n        \n        # Skip connection\n        self.skip = nn.Linear(input_dim, 1)\n        \n        # Initialize weights\n        self._init_weights()\n        \n    def _init_weights(self):\n        \"\"\"Initialize weights\"\"\"\n        for m in self.modules():\n            if isinstance(m, nn.Linear):\n                nn.init.kaiming_normal_(m.weight, mode='fan_out', nonlinearity='relu')\n                if m.bias is not None:\n                    nn.init.constant_(m.bias, 0)\n            elif isinstance(m, nn.BatchNorm1d):\n                nn.init.constant_(m.weight, 1)\n                nn.init.constant_(m.bias, 0)\n    \n    def forward(self, x):\n        # Main path\n        out = self.model(x)\n        \n        # Add skip connection\n        skip = self.skip(x)\n        out = out + 0.1 * skip\n        \n        return out\n\n# =========================\n# Model Trainer\n# =========================\nclass ModelTrainer:\n    \"\"\"Train and evaluate neural network models\"\"\"\n    \n    def __init__(self, config):\n        self.config = config\n        self.models = []\n        self.scalers = []\n        self.feature_selector = None\n        self.feature_engineer = None\n        self.selected_features = None\n        \n    def train(self, train_df):\n        \"\"\"Train models with cross-validation\"\"\"\n        logger.info(\"Starting model training...\")\n        \n        # Feature engineering\n        self.feature_engineer = FeatureEngineer()\n        train_df = self.feature_engineer.create_features(train_df)\n        \n        # Prepare data\n        X = train_df[self.feature_engineer.feature_cols]\n        y = train_df['label'].values.astype(np.float32)\n        \n        # Clean up\n        del train_df\n        gc.collect()\n        \n        # Feature selection\n        self.feature_selector = FeatureSelector(n_features=self.config.TOP_N_FEATURES)\n        X = self.feature_selector.fit_transform(X, y)\n        self.selected_features = X.columns.tolist()\n        \n        # Cross-validation\n        kf = KFold(n_splits=self.config.N_FOLDS, shuffle=True, random_state=self.config.RANDOM_STATE)\n        cv_scores = []\n        \n        for fold, (train_idx, val_idx) in enumerate(kf.split(X)):\n            logger.info(f\"\\nTraining fold {fold + 1}/{self.config.N_FOLDS}\")\n            \n            # Split data\n            X_train = X.iloc[train_idx].values\n            X_val = X.iloc[val_idx].values\n            y_train = y[train_idx]\n            y_val = y[val_idx]\n            \n            # Scale features\n            scaler = StandardScaler()\n            X_train = scaler.fit_transform(X_train).astype(np.float32)\n            X_val = scaler.transform(X_val).astype(np.float32)\n            \n            # Train model\n            model, score = self._train_fold(X_train, y_train, X_val, y_val, fold)\n            \n            # Store\n            self.models.append(model)\n            self.scalers.append(scaler)\n            cv_scores.append(score)\n            \n            # Clean up\n            del X_train, X_val, y_train, y_val\n            gc.collect()\n            \n        logger.info(f\"\\nCV Score: {np.mean(cv_scores):.4f} (+/- {np.std(cv_scores):.4f})\")\n        \n    def _train_fold(self, X_train, y_train, X_val, y_val, fold):\n        \"\"\"Train a single fold\"\"\"\n        # Create datasets\n        train_dataset = TensorDataset(\n            torch.tensor(X_train, dtype=torch.float32),\n            torch.tensor(y_train, dtype=torch.float32).unsqueeze(1)\n        )\n        val_dataset = TensorDataset(\n            torch.tensor(X_val, dtype=torch.float32),\n            torch.tensor(y_val, dtype=torch.float32).unsqueeze(1)\n        )\n        \n        train_loader = DataLoader(\n            train_dataset,\n            batch_size=self.config.BATCH_SIZE,\n            shuffle=True\n        )\n        val_loader = DataLoader(\n            val_dataset,\n            batch_size=self.config.BATCH_SIZE * 2,\n            shuffle=False\n        )\n        \n        # Create model\n        model = CryptoMLP(\n            input_dim=X_train.shape[1],\n            hidden_dims=self.config.HIDDEN_DIMS,\n            dropout_rate=self.config.DROPOUT_RATE\n        ).to(self.config.DEVICE)\n        \n        # Optimizer\n        optimizer = optim.AdamW(model.parameters(), lr=self.config.LEARNING_RATE, weight_decay=1e-5)\n        scheduler = optim.lr_scheduler.ReduceLROnPlateau(optimizer, mode='min', patience=5, factor=0.5)\n        criterion = nn.MSELoss()\n        \n        # Training\n        best_score = -np.inf\n        best_model_state = None\n        patience_counter = 0\n        \n        for epoch in range(self.config.EPOCHS):\n            # Train\n            model.train()\n            train_loss = 0\n            \n            for batch_x, batch_y in train_loader:\n                batch_x = batch_x.to(self.config.DEVICE)\n                batch_y = batch_y.to(self.config.DEVICE)\n                \n                optimizer.zero_grad()\n                output = model(batch_x)\n                loss = criterion(output, batch_y)\n                loss.backward()\n                torch.nn.utils.clip_grad_norm_(model.parameters(), 1.0)\n                optimizer.step()\n                \n                train_loss += loss.item()\n            \n            # Validate\n            model.eval()\n            val_preds = []\n            val_targets = []\n            val_loss = 0\n            \n            with torch.no_grad():\n                for batch_x, batch_y in val_loader:\n                    batch_x = batch_x.to(self.config.DEVICE)\n                    batch_y = batch_y.to(self.config.DEVICE)\n                    \n                    output = model(batch_x)\n                    loss = criterion(output, batch_y)\n                    \n                    val_loss += loss.item()\n                    val_preds.extend(output.cpu().numpy().flatten())\n                    val_targets.extend(batch_y.cpu().numpy().flatten())\n            \n            # Metrics\n            avg_train_loss = train_loss / len(train_loader)\n            avg_val_loss = val_loss / len(val_loader)\n            val_score = pearsonr(val_targets, val_preds)[0]\n            \n            scheduler.step(avg_val_loss)\n            \n            # Early stopping\n            if val_score > best_score:\n                best_score = val_score\n                best_model_state = model.state_dict().copy()\n                patience_counter = 0\n            else:\n                patience_counter += 1\n            \n            if (epoch + 1) % 10 == 0:\n                logger.info(f\"Epoch {epoch+1}: Train Loss = {avg_train_loss:.4f}, \"\n                          f\"Val Loss = {avg_val_loss:.4f}, Val Score = {val_score:.4f}\")\n            \n            if patience_counter >= self.config.EARLY_STOPPING_PATIENCE:\n                logger.info(f\"Early stopping at epoch {epoch+1}\")\n                break\n        \n        # Load best model\n        model.load_state_dict(best_model_state)\n        logger.info(f\"Fold {fold+1} best score: {best_score:.4f}\")\n        \n        return model, best_score\n    \n    def predict(self, test_df):\n        \"\"\"Make predictions on test data\"\"\"\n        logger.info(\"Making predictions...\")\n        \n        # Apply feature engineering\n        test_df = self.feature_engineer.create_features(test_df)\n        \n        # Select features\n        X_test = test_df[self.selected_features].values.astype(np.float32)\n        \n        # Clean up\n        del test_df\n        gc.collect()\n        \n        # Get predictions from each fold\n        all_predictions = []\n        \n        for i, (model, scaler) in enumerate(zip(self.models, self.scalers)):\n            logger.info(f\"Predicting with fold {i+1}/{len(self.models)}\")\n            \n            # Scale\n            X_scaled = scaler.transform(X_test).astype(np.float32)\n            \n            # Predict\n            model.eval()\n            predictions = []\n            \n            with torch.no_grad():\n                dataset = TensorDataset(torch.tensor(X_scaled, dtype=torch.float32))\n                loader = DataLoader(dataset, batch_size=self.config.BATCH_SIZE * 4, shuffle=False)\n                \n                for batch in loader:\n                    batch_x = batch[0].to(self.config.DEVICE)\n                    output = model(batch_x)\n                    predictions.extend(output.cpu().numpy().flatten())\n            \n            all_predictions.append(np.array(predictions))\n            \n            # Clean up\n            del X_scaled\n            gc.collect()\n        \n        # Average predictions\n        final_predictions = np.mean(all_predictions, axis=0)\n        \n        return final_predictions\n\n# =========================\n# Main Pipeline\n# =========================\ndef main():\n    \"\"\"Main execution pipeline\"\"\"\n    logger.info(\"Starting DRW Crypto Prediction Pipeline\")\n    logger.info(\"=\"*60)\n    \n    # Load training data\n    logger.info(\"Loading training data...\")\n    train_df = pd.read_parquet(Config.TRAIN_PATH)\n    train_df = optimize_memory(train_df)\n    logger.info(f\"Training data shape: {train_df.shape}\")\n    \n    # Use recent data if too large\n    if len(train_df) > 600000:\n        logger.info(\"Using recent 600k rows...\")\n        train_df = train_df.tail(600000).reset_index(drop=True)\n    \n    # Initialize trainer\n    trainer = ModelTrainer(Config)\n    \n    # Train models\n    trainer.train(train_df)\n    \n    # Clean up\n    del train_df\n    gc.collect()\n    \n    # Load test data\n    logger.info(\"\\nLoading test data...\")\n    test_df = pd.read_parquet(Config.TEST_PATH)\n    test_df = optimize_memory(test_df)\n    logger.info(f\"Test data shape: {test_df.shape}\")\n    \n    # Generate predictions\n    predictions = trainer.predict(test_df)\n    \n    # Create submission\n    submission = pd.DataFrame({\n        'row_id': range(len(predictions)),\n        'label': predictions\n    })\n    \n    submission.to_csv(Config.OUTPUT_DIR / 'submission.csv', index=False)\n    logger.info(f\"\\nSubmission saved to: {Config.OUTPUT_DIR / 'submission.csv'}\")\n    \n    # Clean up\n    del test_df, predictions\n    gc.collect()\n    \n    logger.info(\"\\n\" + \"=\"*60)\n    logger.info(\"Pipeline completed successfully!\")\n    logger.info(\"=\"*60)\n\nif __name__ == \"__main__\":\n    main()","metadata":{"trusted":true},"outputs":[],"execution_count":null}]}