{"metadata":{"kaggle":{"accelerator":"nvidiaTeslaT4","dataSources":[{"sourceId":84493,"databundleVersionId":9849268,"sourceType":"competition"}],"dockerImageVersionId":30786,"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":true},"kernelspec":{"display_name":"Python 3","language":"python","name":"python3"},"language_info":{"codemirror_mode":{"name":"ipython","version":3},"file_extension":".py","mimetype":"text/x-python","name":"python","nbconvert_exporter":"python","pygments_lexer":"ipython3","version":"3.10.14"},"papermill":{"default_parameters":{},"duration":4.669361,"end_time":"2024-10-10T13:05:46.686069","environment_variables":{},"exception":null,"input_path":"__notebook__.ipynb","output_path":"__notebook__.ipynb","parameters":{},"start_time":"2024-10-10T13:05:42.016708","version":"2.6.0"}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"# pip install polars pandas loguru scikit-learn xgboost\n\nimport os\nimport numpy as np\nimport pandas as pd\nimport polars as pl\nfrom typing import List, Tuple, Optional, Dict\nfrom loguru import logger\nfrom sklearn.model_selection import TimeSeriesSplit\nfrom sklearn.preprocessing import StandardScaler\nfrom xgboost import XGBRegressor\nfrom kaggle_evaluation.jane_street_inference_server import JSInferenceServer\n\n# Set up logging\nlogger.add(\"model.log\", rotation=\"500 MB\")\n\nclass XGBoostSwarm:\n    \"\"\"\n    An improved swarm of XGBoost models for the Jane Street Market Prediction challenge.\n    \"\"\"\n\n    def __init__(self, n_models: int = 1, bag_fraction: float = 0.8):\n        \"\"\"\n        Initialize the XGBoostSwarm.\n\n        Args:\n            n_models (int): Number of XGBoost models in the swarm.\n            bag_fraction (float): Fraction of data to use for each model in bagging.\n        \"\"\"\n        self.n_models = n_models\n        self.bag_fraction = bag_fraction\n        self.models: List[XGBRegressor] = []\n        self.scalers: List[StandardScaler] = []\n        self.feature_columns: List[str] = []\n        self.feature_importances: Dict[str, float] = {}\n\n    def preprocess_data(self, data: pl.DataFrame) -> pl.DataFrame:\n        \"\"\"\n        Preprocess the input data and engineer new features.\n\n        Args:\n            data (pl.DataFrame): Input data.\n\n        Returns:\n            pl.DataFrame: Preprocessed data with engineered features.\n        \"\"\"\n        # Existing features\n        df = data.select(self.feature_columns)\n\n        # 1. Exponential Moving Average (EMA) of key features\n        for col in self.feature_columns[:5]:  # Assuming first 5 features are key\n            df = df.with_columns(pl.col(col).ewm_mean(span=10).alias(f\"{col}_ema\"))\n\n        # 2. Relative Strength Index (RSI) of responder_6\n        if 'responder_6' in data.columns:\n            df = df.with_columns([\n                pl.col('responder_6').diff().alias('delta'),\n                pl.col('responder_6').diff().map(lambda x: max(x, 0)).ewm_mean(span=14).alias('gain'),\n                pl.col('responder_6').diff().map(lambda x: max(-x, 0)).ewm_mean(span=14).alias('loss')\n            ])\n            df = df.with_columns([\n                ((100 - (100 / (1 + pl.col('gain') / pl.col('loss'))))).alias('rsi')\n            ])\n\n        # 3. Volatility measure\n        if 'responder_6' in data.columns:\n            df = df.with_columns([\n                pl.col('responder_6').rolling_std(window_size=10).alias('volatility')\n            ])\n\n        # 4. Cross-feature interactions\n        df = df.with_columns([\n            (pl.col('feature_0') / pl.col('feature_1')).alias('feature_ratio_0_1')\n        ])\n\n        # 5. Time-based features\n        df = df.with_columns([\n            (pl.col('date_id') % 7).alias('day_of_week'),\n            (pl.col('time_id') % 24).alias('hour_of_day')\n        ])\n\n        return df\n\n    def train(self, train_data: pl.DataFrame):\n        \"\"\"\n        Train the swarm of XGBoost models using bagging and time-series cross-validation.\n\n        Args:\n            train_data (pl.DataFrame): Training data.\n        \"\"\"\n        logger.info(\"Starting model training...\")\n\n        self.feature_columns = [col for col in train_data.columns if col.startswith(\"feature_\")]\n        X = self.preprocess_data(train_data)\n        y = train_data[\"responder_6\"].to_numpy()\n\n        tscv = TimeSeriesSplit(n_splits=5)\n\n        for i in range(self.n_models):\n            logger.info(f\"Training model {i+1}/{self.n_models}\")\n            \n            # Bagging: randomly sample a subset of the data\n            sample_indices = np.random.choice(len(X), size=int(self.bag_fraction * len(X)), replace=False)\n            X_sample = X.sample(sample_indices)\n            y_sample = y[sample_indices]\n\n            scaler = StandardScaler()\n            X_scaled = scaler.fit_transform(X_sample)\n            self.scalers.append(scaler)\n\n            model = XGBRegressor(\n                n_estimators=100,\n                learning_rate=0.1,\n                max_depth=5,\n                random_state=42 + i,\n                n_jobs=-1,\n                device=\"cuda\"\n            )\n\n            # Time-series cross-validation\n            for train_index, val_index in tscv.split(X_scaled):\n                X_train, X_val = X_scaled[train_index], X_scaled[val_index]\n                y_train, y_val = y_sample[train_index], y_sample[val_index]\n\n                model.fit(\n                    X_train, y_train,\n                    eval_set=[(X_val, y_val)],\n                    early_stopping_rounds=10,\n                    verbose=False,\n                    callbacks=[self.adaptive_learning_rate_callback()]\n                )\n\n            self.models.append(model)\n\n            # Update feature importances\n            importances = model.feature_importances_\n            for feat, imp in zip(X.columns, importances):\n                self.feature_importances[feat] = self.feature_importances.get(feat, 0) + imp / self.n_models\n\n        # Sort and log feature importances\n        sorted_importances = sorted(self.feature_importances.items(), key=lambda x: x[1], reverse=True)\n        logger.info(\"Feature importances:\")\n        for feat, imp in sorted_importances[:20]:  # Log top 20 features\n            logger.info(f\"{feat}: {imp}\")\n\n        logger.info(\"Model training completed.\")\n\n    def adaptive_learning_rate_callback(self):\n        \"\"\"\n        Custom callback to adjust learning rate based on validation performance.\n        \"\"\"\n        def callback(env):\n            if len(env.evaluation_result_list) > 1 and env.iteration > 0:\n                prev_score = env.evaluation_result_list[-2][2]\n                curr_score = env.evaluation_result_list[-1][2]\n                if curr_score > prev_score:\n                    env.model.learning_rate *= 0.95\n        return callback\n\n    def predict(self, features: pl.DataFrame) -> np.ndarray:\n        \"\"\"\n        Make predictions using the swarm of XGBoost models.\n\n        Args:\n            features (pl.DataFrame): Input features.\n\n        Returns:\n            np.ndarray: Averaged predictions from all models.\n        \"\"\"\n        predictions = np.zeros(len(features))\n        for model, scaler in zip(self.models, self.scalers):\n            X_scaled = scaler.transform(features)\n            predictions += model.predict(X_scaled)\n        return predictions / self.n_models\n\n# Global variables\nmodel: Optional[XGBoostSwarm] = None\nlags_: Optional[pl.DataFrame] = None\n\n\ndef load_and_train_model():\n    \"\"\"\n    Load the training data and train the XGBoostSwarm model.\n    \"\"\"\n    global model\n    logger.info(\"Loading training data...\")\n    train_data = pl.read_parquet(\"/kaggle/input/jane-street-real-time-market-data-forecasting/train.parquet\")\n    \n    logger.info(\"Initializing and training XGBoostSwarm...\")\n    model = XGBoostSwarm(n_models=5, bag_fraction=0.8)\n    model.train(train_data)\n    logger.info(\"Model training completed.\")\n\ndef predict(test: pl.DataFrame, lags: Optional[pl.DataFrame]) -> pl.DataFrame:\n    \"\"\"\n    Make predictions for the test data.\n\n    Args:\n        test (pl.DataFrame): Test data.\n        lags (Optional[pl.DataFrame]): Lagged data.\n\n    Returns:\n        pl.DataFrame: Predictions for responder_6.\n    \"\"\"\n    global model, lags_\n\n    if model is None:\n        logger.error(\"Model not loaded. Please run load_and_train_model() before making predictions.\")\n        raise RuntimeError(\"Model not loaded\")\n\n    if lags is not None:\n        lags_ = lags\n\n    # Preprocess test data and engineer features\n    test_features = model.preprocess_data(test)\n\n    # Make predictions\n    predictions = model.predict(test_features)\n\n    # Create and return the predictions DataFrame\n    return pl.DataFrame({\n        'row_id': test['row_id'],\n        'responder_6': predictions\n    })\n\n# # Set up the inference server\n# inference_server = JSInferenceServer(predict)\n\n# if os.getenv('KAGGLE_IS_COMPETITION_RERUN'):\n#     inference_server.serve()\n# else:\n#     inference_server.run_local_gateway(\n#         (\n#             '/kaggle/input/jane-street-real-time-market-data-forecasting/test.parquet',\n#             '/kaggle/input/jane-street-real-time-market-data-forecasting/lags.parquet',\n#         )\n#     )\n\n\n# Load and train the model explicitly\nload_and_train_model()\n","metadata":{"_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","papermill":{"duration":1.223703,"end_time":"2024-10-10T13:05:45.825911","exception":false,"start_time":"2024-10-10T13:05:44.602208","status":"completed"},"tags":[],"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"The evaluation API requires that you set up a server which will respond to inference requests. We have already defined the server; you just need write the predict function. When we evaluate your submission on the hidden test set the client defined in `jane_street_gateway` will run in a different container with direct access to the hidden test set and hand off the data timestep by timestep.\n\n\n\nYour code will always have access to the published copies of the files.","metadata":{"papermill":{"duration":0.002051,"end_time":"2024-10-10T13:05:45.83073","exception":false,"start_time":"2024-10-10T13:05:45.828679","status":"completed"},"tags":[]}},{"cell_type":"code","source":"import os\nimport numpy as np\nimport pandas as pd\nimport cudf\nimport cupy as cp\nfrom typing import List, Tuple, Optional, Dict\nfrom loguru import logger\nfrom sklearn.model_selection import TimeSeriesSplit\nfrom sklearn.preprocessing import StandardScaler\nfrom xgboost import XGBRegressor\nfrom kaggle_evaluation.jane_street_inference_server import JSInferenceServer\n\n# Set up logging\nlogger.add(\"model.log\", rotation=\"500 MB\")\n\nclass XGBoostSwarm:\n    \"\"\"\n    A memory-efficient, GPU-accelerated swarm of XGBoost models for the Jane Street Market Prediction challenge.\n    \"\"\"\n\n    def __init__(self, n_models: int = 1, bag_fraction: float = 0.4):\n        self.n_models = n_models\n        self.bag_fraction = bag_fraction\n        self.models: List[XGBRegressor] = []\n        self.scalers: List[StandardScaler] = []\n        self.feature_columns: List[str] = []\n        self.feature_importances: Dict[str, float] = {}\n\n    def preprocess_data(self, data: cudf.DataFrame) -> cudf.DataFrame:\n        df = data[self.feature_columns].copy()\n\n        # 1. Exponential Moving Average (EMA) of key features\n        for col in self.feature_columns[:5]:\n            df[f\"{col}_ema\"] = df[col].ewm(span=10).mean()\n\n        # 2. Relative Strength Index (RSI) of responder_6\n        if 'responder_6' in data.columns:\n            delta = data['responder_6'].diff()\n            gain = delta.clip(lower=0)\n            loss = -delta.clip(upper=0)\n            avg_gain = gain.ewm(span=14).mean()\n            avg_loss = loss.ewm(span=14).mean()\n            rs = avg_gain / avg_loss\n            df['rsi'] = 100 - (100 / (1 + rs))\n\n        # 3. Volatility measure\n        if 'responder_6' in data.columns:\n            df['volatility'] = data['responder_6'].rolling(window=10).std()\n\n        # 4. Cross-feature interactions\n        df['feature_ratio_0_1'] = df['feature_0'] / df['feature_1']\n\n        # 5. Time-based features\n        df['day_of_week'] = data['date_id'] % 7\n        df['hour_of_day'] = data['time_id'] % 24\n\n        return df\n\n    def train(self, train_data: cudf.DataFrame):\n        logger.info(\"Starting model training...\")\n\n        self.feature_columns = [col for col in train_data.columns if col.startswith(\"feature_\")]\n        X = self.preprocess_data(train_data)\n        y = train_data[\"responder_6\"].values\n\n        tscv = TimeSeriesSplit(n_splits=5)\n\n        for i in range(self.n_models):\n            logger.info(f\"Training model {i+1}/{self.n_models}\")\n            \n            # Bagging: randomly sample a subset of the data\n            sample_indices = cp.random.choice(len(X), size=int(self.bag_fraction * len(X)), replace=False)\n            X_sample = X.iloc[sample_indices]\n            y_sample = y[sample_indices]\n\n            scaler = StandardScaler()\n            X_scaled = scaler.fit_transform(X_sample.to_pandas())\n            self.scalers.append(scaler)\n\n            model = XGBRegressor(\n                n_estimators=100,\n                learning_rate=0.1,\n                max_depth=5,\n                random_state=42 + i,\n                n_jobs=-1,\n                tree_method='gpu_hist',  # Use GPU acceleration\n                gpu_id=0\n            )\n\n            # Time-series cross-validation\n            for train_index, val_index in tscv.split(X_scaled):\n                X_train, X_val = X_scaled[train_index], X_scaled[val_index]\n                y_train, y_val = y_sample[train_index], y_sample[val_index]\n\n                model.fit(\n                    X_train, y_train,\n                    eval_set=[(X_val, y_val)],\n                    early_stopping_rounds=10,\n                    verbose=False\n                )\n\n            self.models.append(model)\n\n            # Update feature importances\n            importances = model.feature_importances_\n            for feat, imp in zip(X.columns, importances):\n                self.feature_importances[feat] = self.feature_importances.get(feat, 0) + imp / self.n_models\n\n        # Sort and log feature importances\n        sorted_importances = sorted(self.feature_importances.items(), key=lambda x: x[1], reverse=True)\n        logger.info(\"Feature importances:\")\n        for feat, imp in sorted_importances[:20]:  # Log top 20 features\n            logger.info(f\"{feat}: {imp}\")\n\n        logger.info(\"Model training completed.\")\n\n    def predict(self, features: cudf.DataFrame) -> cp.ndarray:\n        predictions = cp.zeros(len(features))\n        for model, scaler in zip(self.models, self.scalers):\n            X_scaled = scaler.transform(features.to_pandas())\n            predictions += cp.asarray(model.predict(X_scaled))\n        return predictions / self.n_models\n\n# Global variables\nmodel: Optional[XGBoostSwarm] = None\nlags_: Optional[cudf.DataFrame] = None\n\ndef load_and_train_model():\n    global model\n    logger.info(\"Loading training data...\")\n    train_data = cudf.read_parquet(\"/kaggle/input/jane-street-real-time-market-data-forecasting/train.parquet\")\n    \n    logger.info(\"Initializing and training XGBoostSwarm...\")\n    model = XGBoostSwarm(n_models=5, bag_fraction=0.8)\n    model.train(train_data)\n    logger.info(\"Model training completed.\")\n\ndef predict(test: cudf.DataFrame, lags: Optional[cudf.DataFrame]) -> pd.DataFrame:\n    global model, lags_\n\n    if model is None:\n        logger.error(\"Model not loaded. Please run load_and_train_model() before making predictions.\")\n        raise RuntimeError(\"Model not loaded\")\n\n    if lags is not None:\n        lags_ = lags\n\n    # Preprocess test data and engineer features\n    test_features = model.preprocess_data(test)\n\n    # Make predictions\n    predictions = model.predict(test_features)\n\n    # Create and return the predictions DataFrame\n    return pd.DataFrame({\n        'row_id': test['row_id'].to_pandas(),\n        'responder_6': predictions.get()\n    })\n\n# Load and train the model explicitly\nload_and_train_model()\n\n# # Set up the inference server\n# inference_server = JSInferenceServer(predict)\n\n# if os.getenv('KAGGLE_IS_COMPETITION_RERUN'):\n#     inference_server.serve()\n# else:\n#     inference_server.run_local_gateway(\n#         (\n#             '/kaggle/input/jane-street-real-time-market-data-forecasting/test.parquet',\n#             '/kaggle/input/jane-street-real-time-market-data-forecasting/lags.parquet',\n#         )\n#     )","metadata":{"execution":{"iopub.status.busy":"2024-10-16T05:48:47.530668Z","iopub.execute_input":"2024-10-16T05:48:47.530972Z","iopub.status.idle":"2024-10-16T05:48:54.662514Z","shell.execute_reply.started":"2024-10-16T05:48:47.530935Z","shell.execute_reply":"2024-10-16T05:48:54.661153Z"},"trusted":true},"execution_count":1,"outputs":[{"traceback":["\u001b[0;31m---------------------------------------------------------------------------\u001b[0m","\u001b[0;31mModuleNotFoundError\u001b[0m                       Traceback (most recent call last)","Cell \u001b[0;32mIn[1], line 7\u001b[0m\n\u001b[1;32m      5\u001b[0m \u001b[38;5;28;01mimport\u001b[39;00m \u001b[38;5;21;01mcupy\u001b[39;00m \u001b[38;5;28;01mas\u001b[39;00m \u001b[38;5;21;01mcp\u001b[39;00m\n\u001b[1;32m      6\u001b[0m \u001b[38;5;28;01mfrom\u001b[39;00m \u001b[38;5;21;01mtyping\u001b[39;00m \u001b[38;5;28;01mimport\u001b[39;00m List, Tuple, Optional, Dict\n\u001b[0;32m----> 7\u001b[0m \u001b[38;5;28;01mfrom\u001b[39;00m \u001b[38;5;21;01mloguru\u001b[39;00m \u001b[38;5;28;01mimport\u001b[39;00m logger\n\u001b[1;32m      8\u001b[0m \u001b[38;5;28;01mfrom\u001b[39;00m \u001b[38;5;21;01msklearn\u001b[39;00m\u001b[38;5;21;01m.\u001b[39;00m\u001b[38;5;21;01mmodel_selection\u001b[39;00m \u001b[38;5;28;01mimport\u001b[39;00m TimeSeriesSplit\n\u001b[1;32m      9\u001b[0m \u001b[38;5;28;01mfrom\u001b[39;00m \u001b[38;5;21;01msklearn\u001b[39;00m\u001b[38;5;21;01m.\u001b[39;00m\u001b[38;5;21;01mpreprocessing\u001b[39;00m \u001b[38;5;28;01mimport\u001b[39;00m StandardScaler\n","\u001b[0;31mModuleNotFoundError\u001b[0m: No module named 'loguru'"],"ename":"ModuleNotFoundError","evalue":"No module named 'loguru'","output_type":"error"}]}]}