{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"pygments_lexer":"ipython3","nbconvert_exporter":"python","version":"3.6.4","file_extension":".py","codemirror_mode":{"name":"ipython","version":3},"name":"python","mimetype":"text/x-python"},"kaggle":{"accelerator":"none","dataSources":[{"sourceId":96164,"databundleVersionId":11418275,"sourceType":"competition"}],"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"# This Python 3 environment comes with many helpful analytics libraries installed\n# It is defined by the kaggle/python Docker image: https://github.com/kaggle/docker-python\n# For example, here's several helpful packages to load\n\nimport numpy as np # linear algebra\nimport pandas as pd # data processing, CSV file I/O (e.g. pd.read_csv)\n\n# Input data files are available in the read-only \"../input/\" directory\n# For example, running this (by clicking run or pressing Shift+Enter) will list all files under the input directory\n\nimport os\nfor dirname, _, filenames in os.walk('/kaggle/input'):\n    for filename in filenames:\n        print(os.path.join(dirname, filename))\n\n# You can write up to 20GB to the current directory (/kaggle/working/) that gets preserved as output when you create a version using \"Save & Run All\" \n# You can also write temporary files to /kaggle/temp/, but they won't be saved outside of the current session","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"#!/usr/bin/env python3\n\"\"\"\nDRW Crypto Market Prediction - Complete Pipeline\n================================================\n1. Progressive feature selection with XGBoost\n2. Training on full dataset with selected features\n3. Time-decay weighted ensemble\n4. Submission generation\n\"\"\"\n\nimport numpy as np\nimport pandas as pd\nimport matplotlib.pyplot as plt\nimport seaborn as sns\nimport xgboost as xgb\nfrom xgboost import XGBRegressor\nfrom sklearn.model_selection import train_test_split, KFold\nfrom sklearn.metrics import mean_squared_error, r2_score\nfrom scipy import stats\nfrom scipy.stats import pearsonr\nimport warnings\nimport os\nimport time\nfrom datetime import datetime\nwarnings.filterwarnings('ignore')\n\n# Set style for better visualizations\nplt.style.use('seaborn-v0_8-darkgrid')\nsns.set_palette(\"husl\")\n\n# Set random seed\nnp.random.seed(42)\n\n# Check if we're in Kaggle environment\nKAGGLE_INPUT_PATH = '/kaggle/input/drw-crypto-market-prediction'\nif os.path.exists(KAGGLE_INPUT_PATH):\n    DATA_PATH = KAGGLE_INPUT_PATH\nelse:\n    DATA_PATH = '.'\n\n# =========================\n# Part 1: Progressive Feature Selection\n# =========================\n\nclass ProgressiveXGBoostSelector:\n    \"\"\"\n    Progressive feature selection using XGBoost with increasing sample sizes\n    \"\"\"\n    \n    def __init__(self, sample_sizes=[50000, 100000, 200000], selection_ratio=0.5):\n        self.sample_sizes = sample_sizes\n        self.selection_ratio = selection_ratio\n        self.stage_results = []\n        self.selected_features = None\n        self.xgb_models = []\n        \n    def remove_bad_features(self, X, feature_names):\n        \"\"\"Remove features with no variance or all missing values\"\"\"\n        n_samples, n_features = X.shape\n        valid_mask = np.ones(n_features, dtype=bool)\n        \n        for i in range(n_features):\n            col = X[:, i]\n            finite_mask = np.isfinite(col)\n            \n            # Check if feature has any finite values\n            if not np.any(finite_mask):\n                valid_mask[i] = False\n                continue\n                \n            finite_vals = col[finite_mask]\n            \n            # Check for zero variance\n            if np.var(finite_vals) < 1e-10 or len(np.unique(finite_vals)) == 1:\n                valid_mask[i] = False\n                continue\n                \n            # Check if too sparse (>99% zeros)\n            if np.sum(finite_vals == 0) / len(finite_vals) > 0.99:\n                valid_mask[i] = False\n        \n        # Return cleaned data and convert feature names to list\n        cleaned_features = np.array(feature_names)[valid_mask]\n        return X[:, valid_mask], list(cleaned_features), valid_mask\n    \n    def rank_transform(self, X):\n        \"\"\"Apply rank transformation to all features\"\"\"\n        X_ranked = np.zeros_like(X)\n        \n        for i in range(X.shape[1]):\n            col = X[:, i]\n            finite_mask = np.isfinite(col)\n            \n            if np.any(finite_mask):\n                # Apply rank transformation to finite values\n                ranks = stats.rankdata(col[finite_mask], method='average')\n                percentiles = (ranks - 1) / (len(ranks) - 1)\n                X_ranked[finite_mask, i] = percentiles\n                X_ranked[~finite_mask, i] = 0.5\n            else:\n                X_ranked[:, i] = 0.5\n        \n        return X_ranked\n    \n    def train_xgboost_and_get_importance(self, X, y, feature_names, stage_num):\n        \"\"\"Train XGBoost and calculate feature importance\"\"\"\n        print(f\"\\nTraining XGBoost for Stage {stage_num}...\")\n        \n        # Split data\n        X_train, X_val, y_train, y_val = train_test_split(\n            X, y, test_size=0.2, random_state=42\n        )\n        \n        # Ensure feature_names is a list\n        if isinstance(feature_names, np.ndarray):\n            feature_names = list(feature_names)\n        \n        # Create DMatrix\n        dtrain = xgb.DMatrix(X_train, label=y_train, feature_names=feature_names)\n        dval = xgb.DMatrix(X_val, label=y_val, feature_names=feature_names)\n        \n        # XGBoost parameters for feature selection\n        params = {\n            'objective': 'reg:squarederror',\n            'max_depth': 6,\n            'learning_rate': 0.1,\n            'subsample': 0.8,\n            'colsample_bytree': 0.8,\n            'random_state': 42,\n            'n_jobs': -1,\n            'tree_method': 'hist'\n        }\n        \n        # Train model\n        evallist = [(dtrain, 'train'), (dval, 'eval')]\n        \n        model = xgb.train(\n            params,\n            dtrain,\n            num_boost_round=300,\n            evals=evallist,\n            early_stopping_rounds=20,\n            verbose_eval=50\n        )\n        \n        # Get predictions for evaluation\n        y_pred = model.predict(dval)\n        rmse = np.sqrt(mean_squared_error(y_val, y_pred))\n        r2 = r2_score(y_val, y_pred)\n        pearson_corr, _ = stats.pearsonr(y_val, y_pred)\n        \n        print(f\"Stage {stage_num} Performance:\")\n        print(f\"  RMSE: {rmse:.6f}\")\n        print(f\"  R²: {r2:.6f}\")\n        print(f\"  Pearson: {pearson_corr:.6f}\")\n        \n        # Get feature importance\n        importance_dict = model.get_score(importance_type='gain')\n        \n        # Create importance dataframe\n        importance_df = pd.DataFrame([\n            {'feature': f, 'importance': importance_dict.get(f, 0)}\n            for f in feature_names\n        ])\n        \n        # Add additional importance metrics\n        importance_df['importance_norm'] = importance_df['importance'] / importance_df['importance'].sum()\n        importance_df['rank'] = importance_df['importance'].rank(ascending=False)\n        importance_df = importance_df.sort_values('importance', ascending=False)\n        \n        # Store results\n        self.xgb_models.append(model)\n        \n        return importance_df, {\n            'rmse': rmse,\n            'r2': r2,\n            'pearson': pearson_corr,\n            'n_features': len(feature_names),\n            'sample_size': len(X)\n        }\n    \n    def select_top_features(self, importance_df, ratio=0.5):\n        \"\"\"Select top features based on importance\"\"\"\n        n_select = max(1, int(len(importance_df) * ratio))\n        top_features = importance_df.head(n_select)['feature'].values\n        return list(top_features)\n    \n    def run_progressive_selection(self, X_full, y_full, feature_names):\n        \"\"\"Run the progressive feature selection pipeline\"\"\"\n        print(\"=\"*80)\n        print(\"PROGRESSIVE FEATURE SELECTION WITH XGBOOST\")\n        print(\"=\"*80)\n        \n        current_features = list(feature_names)\n        stage_summaries = []\n        \n        for stage_num, sample_size in enumerate(self.sample_sizes, 1):\n            print(f\"\\n{'='*60}\")\n            print(f\"STAGE {stage_num}: {sample_size:,} SAMPLES\")\n            print('='*60)\n            \n            # Sample data\n            actual_sample_size = min(sample_size, len(X_full))\n            if actual_sample_size < len(X_full):\n                # Use most recent samples\n                sample_indices = np.arange(len(X_full) - actual_sample_size, len(X_full))\n            else:\n                sample_indices = np.arange(len(X_full))\n            \n            # Get current feature indices\n            feature_indices = [i for i, f in enumerate(feature_names) if f in current_features]\n            \n            # Extract data\n            X_stage = X_full[sample_indices][:, feature_indices]\n            y_stage = y_full[sample_indices]\n            \n            print(f\"Stage {stage_num} data shape: {X_stage.shape}\")\n            \n            # Remove bad features\n            X_clean, features_clean, valid_mask = self.remove_bad_features(\n                X_stage, current_features\n            )\n            \n            print(f\"After cleaning: {X_clean.shape} ({len(current_features) - len(features_clean)} features removed)\")\n            \n            # Rank transform\n            X_ranked = self.rank_transform(X_clean)\n            \n            # Train XGBoost and get importance\n            importance_df, metrics = self.train_xgboost_and_get_importance(\n                X_ranked, y_stage, features_clean, stage_num\n            )\n            \n            # Select top features\n            selected_features = self.select_top_features(importance_df, self.selection_ratio)\n            \n            print(f\"\\nStage {stage_num} Summary:\")\n            print(f\"  Features in: {len(features_clean)}\")\n            print(f\"  Features selected: {len(selected_features)}\")\n            print(f\"  Top 5 features: {selected_features[:5]}\")\n            \n            # Store stage results\n            stage_summary = {\n                'stage': stage_num,\n                'sample_size': actual_sample_size,\n                'features_in': len(features_clean),\n                'features_selected': len(selected_features),\n                'metrics': metrics,\n                'importance_df': importance_df,\n                'selected_features': selected_features,\n                'top_10_features': list(importance_df.head(10)['feature'].values)\n            }\n            stage_summaries.append(stage_summary)\n            self.stage_results.append(stage_summary)\n            \n            # Update features for next stage\n            current_features = selected_features\n            \n            # Early stopping if too few features\n            if len(current_features) < 10:\n                print(f\"\\nStopping early: Only {len(current_features)} features remaining\")\n                break\n        \n        self.selected_features = current_features\n        \n        return current_features, stage_summaries\n\n# =========================\n# Part 2: Full Dataset Training Configuration\n# =========================\n\nclass Config:\n    TRAIN_PATH = os.path.join(DATA_PATH, \"train.parquet\")\n    TEST_PATH = os.path.join(DATA_PATH, \"test.parquet\")\n    SUBMISSION_PATH = os.path.join(DATA_PATH, \"sample_submission.csv\")\n    \n    LABEL_COLUMN = \"label\"\n    N_FOLDS = 3\n    RANDOM_STATE = 42\n\n# Optimized XGBoost parameters for full dataset training\nXGB_PARAMS = {\n    \"tree_method\": \"hist\",\n    \"device\": \"cpu\",  # Change to \"gpu\" if GPU is available\n    \"colsample_bylevel\": 0.4778,\n    \"colsample_bynode\": 0.3628,\n    \"colsample_bytree\": 0.7107,\n    \"gamma\": 1.7095,\n    \"learning_rate\": 0.02213,\n    \"max_depth\": 20,\n    \"max_leaves\": 12,\n    \"min_child_weight\": 16,\n    \"n_estimators\": 1667,\n    \"subsample\": 0.06567,\n    \"reg_alpha\": 39.3524,\n    \"reg_lambda\": 75.4484,\n    \"verbosity\": 0,\n    \"random_state\": Config.RANDOM_STATE,\n    \"n_jobs\": -1\n}\n\nLEARNERS = [\n    {\"name\": \"xgb\", \"Estimator\": XGBRegressor, \"params\": XGB_PARAMS}\n]\n\n# =========================\n# Part 3: Utility Functions\n# =========================\n\ndef create_time_decay_weights(n: int, decay: float = 0.9) -> np.ndarray:\n    \"\"\"Create time decay weights for samples\"\"\"\n    positions = np.arange(n)\n    normalized = positions / (n - 1)\n    weights = decay ** (1.0 - normalized)\n    return weights * n / weights.sum()\n\ndef load_data_with_features(selected_features):\n    \"\"\"Load data with only selected features\"\"\"\n    # Load full data\n    train_df = pd.read_parquet(Config.TRAIN_PATH)\n    test_df = pd.read_parquet(Config.TEST_PATH)\n    submission_df = pd.read_csv(Config.SUBMISSION_PATH)\n    \n    # Filter to selected features\n    train_features = [f for f in selected_features if f in train_df.columns]\n    test_features = [f for f in selected_features if f in test_df.columns]\n    \n    print(f\"\\nUsing {len(train_features)} features for training\")\n    print(f\"Market features included: {[f for f in train_features if f in ['bid_qty', 'ask_qty', 'buy_qty', 'sell_qty', 'volume']]}\")\n    \n    # Select features and label\n    X_train = train_df[train_features]\n    y_train = train_df[Config.LABEL_COLUMN]\n    X_test = test_df[test_features]\n    \n    print(f\"\\nLoaded data - Train: {X_train.shape}, Test: {X_test.shape}, Submission: {submission_df.shape}\")\n    \n    return X_train, y_train, X_test, submission_df, train_features\n\ndef get_model_slices(n_samples: int):\n    \"\"\"Define model training slices\"\"\"\n    return [\n        {\"name\": \"full_data\", \"cutoff\": 0},\n        {\"name\": \"last_75pct\", \"cutoff\": int(0.25 * n_samples)},\n        {\"name\": \"last_50pct\", \"cutoff\": int(0.50 * n_samples)}\n    ]\n\n# =========================\n# Part 4: Training and Evaluation\n# =========================\n\ndef train_and_evaluate(X_train, y_train, X_test, feature_names):\n    \"\"\"Train models with time decay and multiple slices\"\"\"\n    n_samples = len(X_train)\n    model_slices = get_model_slices(n_samples)\n    \n    # Initialize prediction storage\n    oof_preds = {\n        learner[\"name\"]: {s[\"name\"]: np.zeros(n_samples) for s in model_slices}\n        for learner in LEARNERS\n    }\n    test_preds = {\n        learner[\"name\"]: {s[\"name\"]: np.zeros(len(X_test)) for s in model_slices}\n        for learner in LEARNERS\n    }\n    \n    # Create weights\n    full_weights = create_time_decay_weights(n_samples)\n    \n    # K-Fold cross-validation\n    kf = KFold(n_splits=Config.N_FOLDS, shuffle=False)\n    \n    for fold, (train_idx, valid_idx) in enumerate(kf.split(X_train), start=1):\n        print(f\"\\n--- Fold {fold}/{Config.N_FOLDS} ---\")\n        \n        X_valid = X_train.iloc[valid_idx]\n        y_valid = y_train.iloc[valid_idx]\n        \n        for s in model_slices:\n            cutoff = s[\"cutoff\"]\n            slice_name = s[\"name\"]\n            \n            # Get slice of data\n            subset_X = X_train.iloc[cutoff:].reset_index(drop=True)\n            subset_y = y_train.iloc[cutoff:].reset_index(drop=True)\n            \n            # Adjust indices for slice\n            rel_idx = train_idx[train_idx >= cutoff] - cutoff\n            \n            X_train_slice = subset_X.iloc[rel_idx]\n            y_train_slice = subset_y.iloc[rel_idx]\n            \n            # Get weights for slice\n            if cutoff > 0:\n                sw = create_time_decay_weights(len(subset_X))[rel_idx]\n            else:\n                sw = full_weights[train_idx]\n            \n            print(f\"  Training slice: {slice_name}, samples: {len(X_train_slice)}\")\n            \n            for learner in LEARNERS:\n                model = learner[\"Estimator\"](**learner[\"params\"])\n                \n                # Train model\n                model.fit(\n                    X_train_slice, \n                    y_train_slice, \n                    sample_weight=sw,\n                    eval_set=[(X_valid, y_valid)],\n                    verbose=False\n                )\n                \n                # Out-of-fold predictions\n                mask = valid_idx >= cutoff\n                if mask.any():\n                    idxs = valid_idx[mask]\n                    oof_preds[learner[\"name\"]][slice_name][idxs] = model.predict(X_train.iloc[idxs])\n                \n                # Handle predictions for samples before cutoff\n                if cutoff > 0 and (~mask).any():\n                    oof_preds[learner[\"name\"]][slice_name][valid_idx[~mask]] = \\\n                        oof_preds[learner[\"name\"]][\"full_data\"][valid_idx[~mask]]\n                \n                # Test predictions\n                test_preds[learner[\"name\"]][slice_name] += model.predict(X_test)\n    \n    # Normalize test predictions by number of folds\n    for learner_name in test_preds:\n        for slice_name in test_preds[learner_name]:\n            test_preds[learner_name][slice_name] /= Config.N_FOLDS\n    \n    return oof_preds, test_preds, model_slices\n\n# =========================\n# Part 5: Ensemble & Submission\n# =========================\n\ndef ensemble_and_submit(y_train, oof_preds, test_preds, submission_df):\n    \"\"\"Create ensemble predictions and generate submission\"\"\"\n    learner_ensembles = {}\n    \n    for learner_name in oof_preds:\n        # Calculate scores for each slice\n        scores = {}\n        for s in oof_preds[learner_name]:\n            scores[s] = pearsonr(y_train, oof_preds[learner_name][s])[0]\n        \n        total_score = sum(scores.values())\n        \n        # Simple average ensemble\n        oof_simple = np.mean(list(oof_preds[learner_name].values()), axis=0)\n        test_simple = np.mean(list(test_preds[learner_name].values()), axis=0)\n        score_simple = pearsonr(y_train, oof_simple)[0]\n        \n        # Weighted average ensemble\n        oof_weighted = sum(scores[s] / total_score * oof_preds[learner_name][s] for s in scores)\n        test_weighted = sum(scores[s] / total_score * test_preds[learner_name][s] for s in scores)\n        score_weighted = pearsonr(y_train, oof_weighted)[0]\n        \n        print(f\"\\n{learner_name.upper()} Model Scores:\")\n        for slice_name, score in scores.items():\n            print(f\"  {slice_name}: {score:.4f}\")\n        print(f\"  Simple Ensemble Pearson:   {score_simple:.4f}\")\n        print(f\"  Weighted Ensemble Pearson: {score_weighted:.4f}\")\n        \n        # Store ensemble predictions\n        learner_ensembles[learner_name] = {\n            \"oof_simple\": oof_simple,\n            \"test_simple\": test_simple,\n            \"oof_weighted\": oof_weighted,\n            \"test_weighted\": test_weighted\n        }\n    \n    # Final ensemble across all learners\n    final_oof = np.mean([le[\"oof_simple\"] for le in learner_ensembles.values()], axis=0)\n    final_test = np.mean([le[\"test_simple\"] for le in learner_ensembles.values()], axis=0)\n    final_score = pearsonr(y_train, final_oof)[0]\n    \n    print(f\"\\nFINAL ensemble across learners Pearson: {final_score:.4f}\")\n    \n    # Create submission\n    submission_df[\"prediction\"] = final_test\n    submission_df.to_csv(\"submission.csv\", index=False)\n    print(\"\\nSaved: submission.csv\")\n    \n    return final_score, learner_ensembles\n\n# =========================\n# Main Execution Pipeline\n# =========================\n\ndef main():\n    \"\"\"Main execution function\"\"\"\n    print(\"=\"*80)\n    print(\"DRW CRYPTO MARKET PREDICTION - COMPLETE PIPELINE\")\n    print(f\"Started at: {datetime.now()}\")\n    print(\"=\"*80)\n    \n    # Step 1: Load initial data for feature selection\n    print(\"\\n\" + \"=\"*60)\n    print(\"STEP 1: PROGRESSIVE FEATURE SELECTION\")\n    print(\"=\"*60)\n    \n    train_path = os.path.join(DATA_PATH, 'train.parquet')\n    df = pd.read_parquet(train_path)\n    print(f\"Total dataset size: {len(df):,} rows\")\n    \n    # Separate features and target\n    feature_cols = [col for col in df.columns if col not in ['timestamp', 'label']]\n    X = df[feature_cols].values\n    y = df['label'].values\n    \n    print(f\"Initial shape: {X.shape}\")\n    print(f\"Target statistics - Mean: {np.mean(y):.6f}, Std: {np.std(y):.6f}\")\n    \n    # Run progressive feature selection\n    selector = ProgressiveXGBoostSelector(\n        sample_sizes=[50000, 100000, 200000],\n        selection_ratio=0.5\n    )\n    \n    start_time = time.time()\n    final_features, stage_summaries = selector.run_progressive_selection(X, y, feature_cols)\n    selection_time = time.time() - start_time\n    \n    print(f\"\\nFeature selection completed in {selection_time:.2f} seconds\")\n    print(f\"Selected {len(final_features)} features from {len(feature_cols)}\")\n    \n    # Save selected features\n    pd.DataFrame({'feature': final_features}).to_csv('selected_features.csv', index=False)\n    print(\"Selected features saved to 'selected_features.csv'\")\n    \n    # Step 2: Train on full dataset with selected features\n    print(\"\\n\" + \"=\"*60)\n    print(\"STEP 2: FULL DATASET TRAINING WITH SELECTED FEATURES\")\n    print(\"=\"*60)\n    \n    # Load data with selected features\n    X_train, y_train, X_test, submission_df, train_features = load_data_with_features(final_features)\n    \n    # Train models\n    start_time = time.time()\n    oof_preds, test_preds, model_slices = train_and_evaluate(X_train, y_train, X_test, train_features)\n    training_time = time.time() - start_time\n    \n    print(f\"\\nTraining completed in {training_time:.2f} seconds ({training_time/60:.2f} minutes)\")\n    \n    # Step 3: Create ensemble and submission\n    print(\"\\n\" + \"=\"*60)\n    print(\"STEP 3: ENSEMBLE AND SUBMISSION\")\n    print(\"=\"*60)\n    \n    final_score, learner_ensembles = ensemble_and_submit(y_train, oof_preds, test_preds, submission_df)\n    \n    # Summary\n    print(\"\\n\" + \"=\"*60)\n    print(\"PIPELINE SUMMARY\")\n    print(\"=\"*60)\n    print(f\"Total processing time: {selection_time + training_time:.2f} seconds\")\n    print(f\"Feature selection time: {selection_time:.2f} seconds\")\n    print(f\"Model training time: {training_time:.2f} seconds\")\n    print(f\"Features used: {len(final_features)} (from {len(feature_cols)} original)\")\n    print(f\"Final CV Pearson score: {final_score:.4f}\")\n    print(f\"\\nTop 10 selected features:\")\n    for i, feat in enumerate(final_features[:10], 1):\n        print(f\"  {i:2d}. {feat}\")\n    \n    print(f\"\\nCompleted at: {datetime.now()}\")\n    print(\"=\"*80)\n    \n    return selector, final_features, final_score\n\nif __name__ == \"__main__\":\n    selector, final_features, final_score = main()","metadata":{"trusted":true},"outputs":[],"execution_count":null}]}