{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"pygments_lexer":"ipython3","nbconvert_exporter":"python","version":"3.6.4","file_extension":".py","codemirror_mode":{"name":"ipython","version":3},"name":"python","mimetype":"text/x-python"},"kaggle":{"accelerator":"none","dataSources":[{"sourceId":96164,"databundleVersionId":12993472,"sourceType":"competition"}],"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"# This Python 3 environment comes with many helpful analytics libraries installed\n# It is defined by the kaggle/python Docker image: https://github.com/kaggle/docker-python\n# For example, here's several helpful packages to load\n\nimport numpy as np # linear algebra\nimport pandas as pd # data processing, CSV file I/O (e.g. pd.read_csv)\n\n# Input data files are available in the read-only \"../input/\" directory\n# For example, running this (by clicking run or pressing Shift+Enter) will list all files under the input directory\n\nimport os\nfor dirname, _, filenames in os.walk('/kaggle/input'):\n    for filename in filenames:\n        print(os.path.join(dirname, filename))\n\n# You can write up to 20GB to the current directory (/kaggle/working/) that gets preserved as output when you create a version using \"Save & Run All\" \n# You can also write temporary files to /kaggle/temp/, but they won't be saved outside of the current session","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# Install required packages\n!pip install feature-engine category_encoders catboost shap optuna scikit-learn xgboost lightgbm -q\n\n# Memory-Optimized DRW Crypto Market Prediction Pipeline with James-Stein Encoding\n# Optimized for Kaggle's memory constraints while maintaining prediction quality\n\nimport numpy as np\nimport pandas as pd\nimport os\nimport sys\nimport warnings\nimport gc\nimport logging\nfrom datetime import datetime\nfrom typing import List, Tuple, Dict, Optional, Union\n\n# Core ML libraries\nfrom sklearn.model_selection import KFold, TimeSeriesSplit\nfrom sklearn.pipeline import Pipeline\nfrom sklearn.preprocessing import StandardScaler, RobustScaler, KBinsDiscretizer\nfrom sklearn.compose import ColumnTransformer\nfrom sklearn.feature_selection import VarianceThreshold, SelectKBest, f_regression\nfrom sklearn.metrics import mean_squared_error\nfrom sklearn.base import BaseEstimator, TransformerMixin\n\n# Boosting models\nfrom xgboost import XGBRegressor\nfrom lightgbm import LGBMRegressor\nfrom catboost import CatBoostRegressor\n\n# Statistical and feature engineering\nfrom scipy.stats import pearsonr, spearmanr\nfrom feature_engine.outliers import Winsorizer\nfrom feature_engine.transformation import YeoJohnsonTransformer\nfrom category_encoders import JamesSteinEncoder\n\n# Optimization and analysis\nimport optuna\nfrom optuna.samplers import TPESampler\nimport shap\nimport matplotlib.pyplot as plt\nimport seaborn as sns\n\n# Configure environment\nwarnings.filterwarnings('ignore')\nlogging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s')\nos.environ['PYTHONHASHSEED'] = '0'\nnp.random.seed(42)\n\n# Force garbage collection\ngc.collect()\n\n# ============================================================================================\n# Memory-Optimized Configuration\n# ============================================================================================\n\nclass Config:\n    \"\"\"Centralized configuration optimized for memory efficiency\"\"\"\n    \n    # File paths\n    TRAIN_PATH = \"/kaggle/input/drw-crypto-market-prediction/train.parquet\"\n    TEST_PATH = \"/kaggle/input/drw-crypto-market-prediction/test.parquet\"\n    SUBMISSION_PATH = \"/kaggle/input/drw-crypto-market-prediction/sample_submission.csv\"\n    \n    # Feature groups\n    MARKET_FEATURES = ['bid_qty', 'ask_qty', 'buy_qty', 'sell_qty', 'volume']\n    \n    # Core anonymous features identified through analysis (reduced set)\n    KEY_X_FEATURES = [\n        'X287', 'X446', 'X66', 'X123', 'X385', 'X594', 'X25', 'X3',\n        'X37', 'X174', 'X298', 'X168', 'X1', 'X76', 'X21', 'X19'\n    ]\n    \n    # Features to discretize for James-Stein encoding\n    FEATURES_TO_DISCRETIZE = ['volume', 'buy_qty', 'sell_qty']  # Reduced set\n    N_BINS = 8  # Reduced bins for memory efficiency\n    \n    # Model parameters\n    LABEL_COLUMN = \"label\"\n    RANDOM_STATE = 42\n    N_FOLDS = 3  # Reduced folds for memory\n    VALIDATION_SIZE = 0.15\n    GAP_SIZE = 0.02\n    \n    # Memory-optimized parameters\n    INITIAL_FEATURE_REDUCTION = 300  # Aggressively reduce features early\n    VARIANCE_THRESHOLD = 0.01  # Higher threshold to remove more features\n    N_SHAP_FEATURES = 50  # Reduced final features\n    SHAP_SAMPLE_SIZE = 2000  # Smaller sample for SHAP\n    USE_FLOAT32 = True  # Use float32 instead of float64\n    CHUNK_SIZE = 50000  # Process data in chunks\n    \n    # Simplified ensemble\n    ENSEMBLE_WEIGHTS = {\n        'xgboost': 1.0  # Use only XGBoost to save memory\n    }\n\n# ============================================================================================\n# Memory-Efficient Data Type Optimization\n# ============================================================================================\n\ndef optimize_dtypes(df):\n    \"\"\"Convert data types to save memory\"\"\"\n    for col in df.columns:\n        col_type = df[col].dtype\n        \n        if col_type != 'object':\n            c_min = df[col].min()\n            c_max = df[col].max()\n            \n            if str(col_type)[:3] == 'int':\n                if c_min > np.iinfo(np.int8).min and c_max < np.iinfo(np.int8).max:\n                    df[col] = df[col].astype(np.int8)\n                elif c_min > np.iinfo(np.int16).min and c_max < np.iinfo(np.int16).max:\n                    df[col] = df[col].astype(np.int16)\n                elif c_min > np.iinfo(np.int32).min and c_max < np.iinfo(np.int32).max:\n                    df[col] = df[col].astype(np.int32)\n            else:\n                if Config.USE_FLOAT32:\n                    df[col] = df[col].astype(np.float32)\n                    \n    return df\n\n# ============================================================================================\n# Early Feature Reduction\n# ============================================================================================\n\ndef reduce_features_early(df, target=None, n_features=300):\n    \"\"\"Aggressively reduce features early in the pipeline\"\"\"\n    \n    # First, remove constant features\n    constant_filter = VarianceThreshold(threshold=0)\n    df_filtered = constant_filter.fit_transform(df)\n    kept_cols = df.columns[constant_filter.get_support()].tolist()\n    df = df[kept_cols]\n    \n    logging.info(f\"Removed {len(df.columns) - len(kept_cols)} constant features\")\n    \n    # If we still have too many features and have a target, use correlation\n    if len(df.columns) > n_features and target is not None:\n        correlations = df.corrwith(target).abs()\n        top_features = correlations.nlargest(n_features).index.tolist()\n        df = df[top_features]\n        logging.info(f\"Reduced to top {n_features} correlated features\")\n    \n    return df\n\n# ============================================================================================\n# Memory-Efficient James-Stein Encoder\n# ============================================================================================\n\nclass MemoryEfficientJamesSteinEncoder:\n    \"\"\"Memory-optimized James-Stein encoder\"\"\"\n    \n    def __init__(self, categorical_features, min_samples_leaf=20, smoothing=1.0):\n        self.categorical_features = categorical_features\n        self.min_samples_leaf = min_samples_leaf\n        self.smoothing = smoothing\n        self.global_mean = None\n        self.encoders = {}\n        \n    def fit(self, X, y):\n        \"\"\"Fit the encoder on training data\"\"\"\n        if not self.categorical_features:\n            return self\n        \n        self.global_mean = float(y.mean())\n        \n        # Convert y to Series if needed\n        if not isinstance(y, pd.Series):\n            y = pd.Series(y, index=X.index)\n        \n        # Process each categorical feature\n        for col in self.categorical_features:\n            if col in X.columns:\n                try:\n                    encoder = JamesSteinEncoder(\n                        cols=[col],\n                        return_df=True,\n                        handle_missing='value',\n                        handle_unknown='value',\n                        model='independent',\n                        randomized=True,\n                        sigma=self.smoothing\n                    )\n                    \n                    encoder.fit(X[[col]], y)\n                    self.encoders[col] = encoder\n                    \n                    logging.info(f\"Fitted encoder for {col}\")\n                except Exception as e:\n                    logging.warning(f\"Failed to fit encoder for {col}: {str(e)}\")\n        \n        return self\n    \n    def transform(self, X):\n        \"\"\"Transform features using fitted encoders\"\"\"\n        X_transformed = X.copy()\n        \n        for col in self.categorical_features:\n            if col in X.columns and col in self.encoders:\n                try:\n                    encoded_values = self.encoders[col].transform(X[[col]])\n                    encoded_col_name = f'{col}_encoded'\n                    X_transformed[encoded_col_name] = encoded_values[col].astype(np.float32)\n                    X_transformed = X_transformed.drop(columns=[col])\n                except Exception as e:\n                    logging.warning(f\"Failed to transform {col}: {str(e)}\")\n        \n        # Fill missing values\n        if self.global_mean is not None:\n            X_transformed = X_transformed.fillna(self.global_mean)\n        \n        return X_transformed\n\n# ============================================================================================\n# Simplified Feature Engineering\n# ============================================================================================\n\ndef create_essential_features(df):\n    \"\"\"Create only the most essential engineered features\"\"\"\n    \n    # Order flow imbalance\n    if all(col in df.columns for col in ['buy_qty', 'sell_qty', 'volume']):\n        df['order_flow_imbalance'] = (df['buy_qty'] - df['sell_qty']) / (df['volume'] + 1e-10)\n        df['buy_ratio'] = df['buy_qty'] / (df['buy_qty'] + df['sell_qty'] + 1e-10)\n    \n    # Market depth\n    if all(col in df.columns for col in ['bid_qty', 'ask_qty']):\n        df['bid_ask_imbalance'] = (df['bid_qty'] - df['ask_qty']) / (df['bid_qty'] + df['ask_qty'] + 1e-10)\n    \n    # Log volume\n    if 'volume' in df.columns:\n        df['log_volume'] = np.log1p(df['volume'])\n    \n    # Convert to float32\n    if Config.USE_FLOAT32:\n        float_cols = df.select_dtypes(include=[np.float64]).columns\n        df[float_cols] = df[float_cols].astype(np.float32)\n    \n    return df\n\ndef create_categorical_features_memory_efficient(df, features_to_discretize, n_bins=8):\n    \"\"\"Memory-efficient categorical feature creation\"\"\"\n    categorical_features = []\n    \n    for feature in features_to_discretize:\n        if feature in df.columns:\n            try:\n                # Use pandas qcut for memory efficiency\n                cat_feature_name = f'{feature}_cat'\n                df[cat_feature_name] = pd.qcut(df[feature], q=n_bins, labels=False, duplicates='drop')\n                df[cat_feature_name] = df[cat_feature_name].fillna(-1).astype(np.int8)\n                categorical_features.append(cat_feature_name)\n                \n            except Exception as e:\n                logging.warning(f\"Failed to discretize {feature}: {str(e)}\")\n    \n    return df, categorical_features\n\n# ============================================================================================\n# Lightweight Feature Selection\n# ============================================================================================\n\nclass LightweightFeatureSelector:\n    \"\"\"Memory-efficient feature selection using simple methods\"\"\"\n    \n    def __init__(self, n_features=50):\n        self.n_features = n_features\n        self.selected_features = None\n        \n    def fit(self, X, y):\n        \"\"\"Select features using correlation and variance\"\"\"\n        \n        # Remove low variance features\n        variance_filter = VarianceThreshold(threshold=Config.VARIANCE_THRESHOLD)\n        X_filtered = variance_filter.fit_transform(X)\n        kept_features = X.columns[variance_filter.get_support()].tolist()\n        \n        # Use SelectKBest for efficiency\n        selector = SelectKBest(f_regression, k=min(self.n_features, len(kept_features)))\n        selector.fit(X_filtered, y)\n        \n        # Get selected feature names\n        selected_mask = selector.get_support()\n        self.selected_features = [kept_features[i] for i in range(len(kept_features)) if selected_mask[i]]\n        \n        logging.info(f\"Selected {len(self.selected_features)} features\")\n        \n        return self\n    \n    def transform(self, X):\n        \"\"\"Transform to selected features\"\"\"\n        return X[self.selected_features]\n\n# ============================================================================================\n# Memory-Efficient Model Training\n# ============================================================================================\n\ndef train_model_memory_efficient(X_train, y_train, X_val, y_val, X_test, categorical_features):\n    \"\"\"Memory-efficient model training pipeline\"\"\"\n    \n    # Apply encoding\n    encoder = MemoryEfficientJamesSteinEncoder(categorical_features=categorical_features)\n    X_train_encoded = encoder.fit(X_train, y_train).transform(X_train)\n    X_val_encoded = encoder.transform(X_val)\n    X_test_encoded = encoder.transform(X_test)\n    \n    # Clear original data\n    del X_train, X_val, X_test\n    gc.collect()\n    \n    # Feature selection\n    selector = LightweightFeatureSelector(n_features=Config.N_SHAP_FEATURES)\n    X_train_selected = selector.fit(X_train_encoded, y_train).transform(X_train_encoded)\n    X_val_selected = selector.transform(X_val_encoded)\n    X_test_selected = selector.transform(X_test_encoded)\n    \n    # Clear encoded data\n    del X_train_encoded, X_val_encoded, X_test_encoded\n    gc.collect()\n    \n    # Simple preprocessing\n    scaler = RobustScaler()\n    X_train_scaled = scaler.fit_transform(X_train_selected)\n    X_val_scaled = scaler.transform(X_val_selected)\n    X_test_scaled = scaler.transform(X_test_selected)\n    \n    # Train model with conservative parameters\n    model = XGBRegressor(\n        n_estimators=500,  # Reduced trees\n        max_depth=4,  # Shallower trees\n        learning_rate=0.05,\n        subsample=0.8,\n        colsample_bytree=0.8,\n        reg_alpha=10,\n        reg_lambda=10,\n        random_state=Config.RANDOM_STATE,\n        n_jobs=1,  # Single thread to save memory\n        verbosity=0\n    )\n    \n    model.fit(X_train_scaled, y_train)\n    \n    # Generate predictions\n    val_pred = model.predict(X_val_scaled)\n    test_pred = model.predict(X_test_scaled)\n    \n    # Calculate score\n    val_score = pearsonr(y_val, val_pred)[0]\n    \n    return val_pred, test_pred, val_score\n\n# ============================================================================================\n# Main Training Function with Memory Management\n# ============================================================================================\n\ndef train_and_predict_memory_efficient(train_df, test_df):\n    \"\"\"Memory-efficient training pipeline\"\"\"\n    \n    logging.info(\"Starting memory-optimized pipeline\")\n    logging.info(f\"Initial train shape: {train_df.shape}, test shape: {test_df.shape}\")\n    \n    # Optimize data types\n    train_df = optimize_dtypes(train_df)\n    test_df = optimize_dtypes(test_df)\n    gc.collect()\n    \n    # Extract target\n    y_train = train_df[Config.LABEL_COLUMN].values\n    \n    # Early feature reduction for X features\n    x_cols = [col for col in train_df.columns if col.startswith('X')]\n    other_cols = [col for col in train_df.columns if not col.startswith('X') and col != Config.LABEL_COLUMN]\n    \n    # Reduce X features\n    if len(x_cols) > Config.INITIAL_FEATURE_REDUCTION:\n        x_train = train_df[x_cols]\n        x_train = reduce_features_early(x_train, pd.Series(y_train), Config.INITIAL_FEATURE_REDUCTION)\n        reduced_x_cols = x_train.columns.tolist()\n        \n        # Update dataframes\n        train_df = pd.concat([train_df[other_cols + [Config.LABEL_COLUMN]], x_train], axis=1)\n        test_df = pd.concat([test_df[other_cols], test_df[reduced_x_cols]], axis=1)\n        \n        del x_train\n        gc.collect()\n    \n    # Create essential features\n    train_df = create_essential_features(train_df)\n    test_df = create_essential_features(test_df)\n    \n    # Create categorical features\n    train_df, categorical_features = create_categorical_features_memory_efficient(\n        train_df, Config.FEATURES_TO_DISCRETIZE, Config.N_BINS\n    )\n    test_df, _ = create_categorical_features_memory_efficient(\n        test_df, Config.FEATURES_TO_DISCRETIZE, Config.N_BINS\n    )\n    \n    # Prepare features\n    feature_cols = [col for col in train_df.columns \n                   if col not in ['timestamp', 'ID', Config.LABEL_COLUMN]]\n    \n    X_train = train_df[feature_cols]\n    X_test = test_df[feature_cols]\n    \n    logging.info(f\"Features after engineering: {len(feature_cols)}\")\n    \n    # Simple time series split\n    n_samples = len(X_train)\n    val_size = int(Config.VALIDATION_SIZE * n_samples)\n    \n    # Use only one validation split to save memory\n    train_end = n_samples - val_size\n    \n    X_train_fold = X_train.iloc[:train_end]\n    y_train_fold = y_train[:train_end]\n    X_val_fold = X_train.iloc[train_end:]\n    y_val_fold = y_train[train_end:]\n    \n    # Train model\n    logging.info(\"Training model...\")\n    val_pred, test_pred, val_score = train_model_memory_efficient(\n        X_train_fold, y_train_fold,\n        X_val_fold, y_val_fold,\n        X_test,\n        categorical_features\n    )\n    \n    logging.info(f\"Validation score: {val_score:.4f}\")\n    \n    # Clean up\n    del X_train, X_train_fold, X_val_fold\n    gc.collect()\n    \n    return test_pred, val_score\n\n# ============================================================================================\n# Main Execution\n# ============================================================================================\n\ndef main():\n    \"\"\"Main execution function with memory management\"\"\"\n    try:\n        # Load data\n        logging.info(\"Loading 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        # Validate data\n        assert Config.LABEL_COLUMN in train_df.columns, f\"Label column {Config.LABEL_COLUMN} not found\"\n        \n        # Remove missing labels\n        train_df = train_df.dropna(subset=[Config.LABEL_COLUMN])\n        \n        # Train and predict\n        test_predictions, val_score = train_and_predict_memory_efficient(train_df, test_df)\n        \n        # Create submission\n        submission_df['prediction'] = test_predictions\n        submission_df.to_csv('submission_memory_optimized.csv', index=False)\n        \n        logging.info(f\"\\nValidation Score: {val_score:.4f}\")\n        logging.info(\"Submission saved as submission_memory_optimized.csv\")\n        logging.info(\"Pipeline completed successfully!\")\n        \n    except Exception as e:\n        logging.error(f\"Pipeline failed: {str(e)}\")\n        raise\n    finally:\n        # Final cleanup\n        gc.collect()\n\nif __name__ == \"__main__\":\n    main()","metadata":{"trusted":true},"outputs":[],"execution_count":null}]}