{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"name":"python","version":"3.10.14","mimetype":"text/x-python","codemirror_mode":{"name":"ipython","version":3},"pygments_lexer":"ipython3","nbconvert_exporter":"python","file_extension":".py"},"kaggle":{"accelerator":"none","dataSources":[{"sourceId":84493,"databundleVersionId":9871156,"sourceType":"competition"}],"dockerImageVersionId":30804,"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"markdown","source":"# Introduction\n\nI started a proof of concept of online learning using ARIMA as a model. Now, moving towards the usage of a more powerful model I created a evaluation pipeline using LGBM to verifiy what R2 score should I expect when submit the solution and whether the training limit constraint (1 minute) has been met.\n\nCheck out previous work:\n- [JRTSMDF - ARIMA online Learning](https://www.kaggle.com/code/serjhenrique/jrtsmdf-arima-online-learning)\n- [Blog post explaining ARIMA Solution in detail](https://serjhenrique.com/arima-and-online-learning-in-financial-forecasting/)\n\n\nYou can use this notebook to play with feature engineering and model parameters (Let me know what changes worked in your experimentation ). The submission notebook is on its way.\n\n\nFrom my initial experiments the R2 score improves when we increase the training window. \n\n**OBS**: Change the date_id filter to a date_id like 1000 or any prior to that to have a better evaluation.\n\n**If you find it useful, please, upvote. I hope you enjoy!!**","metadata":{}},{"cell_type":"code","source":"import pandas as pd\nimport polars as pl\nimport numpy as np\nimport lightgbm as lgb\nimport time\n\nfrom typing import List, Union\n\nimport wandb\nfrom kaggle_secrets import UserSecretsClient","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","trusted":true,"execution":{"iopub.status.busy":"2024-12-05T13:07:20.706801Z","iopub.execute_input":"2024-12-05T13:07:20.707270Z","iopub.status.idle":"2024-12-05T13:07:24.408118Z","shell.execute_reply.started":"2024-12-05T13:07:20.707219Z","shell.execute_reply":"2024-12-05T13:07:24.407058Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"class FeatureEngineering:\n\n    def __init__(\n        self, \n        data: Union[pl.DataFrame, pl.LazyFrame], \n        lags: Union[pl.DataFrame, pl.LazyFrame, None] = None\n    ):\n        # Check if data is LazyFrame or DataFrame\n        if isinstance(data, pl.DataFrame):\n            self.data = data.lazy()  # Convert to LazyFrame if it's a DataFrame\n        else:\n            self.data = data  # Keep as LazyFrame\n\n        # Check if lags is LazyFrame or DataFrame\n        if isinstance(lags, pl.DataFrame):\n            self.lags = lags.lazy()  # Convert to LazyFrame if it's a DataFrame\n        else:\n            self.lags = lags  # Keep as LazyFrame\n\n    def add_lag_responder(self, data: pl.LazyFrame, columns: List[str]) -> pl.LazyFrame:\n        if self.lags is None:\n            last_responder_df = data.group_by(['symbol_id','date_id']).agg(\n                (pl.col(c).last().cast(pl.Float32).alias(f'{c}_lag_1') for c in columns)\n            )\n        else:\n            last_responder_df = self.lags.group_by(['symbol_id','date_id']).agg(\n                (pl.col(f'{c}_lag_1').last().cast(pl.Float32).alias(f'{c}_lag_1') for c in columns)\n            )\n    \n        last_responder_df = last_responder_df.with_columns(\n            pl.col('date_id') + 1\n        )\n\n        data = data.join(\n            last_responder_df,\n            on=['symbol_id','date_id'],\n            how='left'\n        ) \n\n        return data\n\n    def drop_responders(self, data: pl.LazyFrame) -> pl.LazyFrame:\n        return data.drop([f'responder_{i}' for i in range(9) if i != 6])\n\n    def drop_partition_id(self, data: pl.LazyFrame) -> pl.LazyFrame:\n        return data.drop(['partition_id']) if 'partition_id' in data.collect_schema().names() else data\n\n    def differencing_transform(self, data: pl.LazyFrame, cols: List[str]) -> pl.LazyFrame:\n\n        return (\n            data.sort([\"symbol_id\", \"date_id\", \"time_id\"])\n            .with_columns(\n                [\n                    (pl.col(col).diff(1))\n                    .over(\"symbol_id\")\n                    .alias(col)\n                    for col in cols\n                ]\n            )\n        )\n\n    def create_rolling_features(self, data: pl.LazyFrame, period: int, group_col: Union[str, List[str]], agg_col: Union[str, List[str]]) -> pl.LazyFrame:\n\n        rolling_df = data.with_columns(\n            [\n                pl.col(c)\n                .rolling_mean(window_size=period, min_periods=1)\n                .over(group_col)\n                .cast(pl.Float32)\n                .alias(f\"{c}_rolling_{str(period)}_mean\")\n                for c in (agg_col if isinstance(agg_col, list) else [agg_col])\n            ]\n        )\n\n        return rolling_df\n\n    def create_daily_rolling_features(self, data: pl.LazyFrame, period: int , group_col: Union[str, List[str]], agg_col: Union[str, List[str]]) -> pl.LazyFrame:\n\n        daily_df = data.select(group_col+agg_col).unique().sort(group_col)\n\n        rolling_df = daily_df.with_columns(\n            [\n                pl.col(c)\n                .rolling_mean(window_size=period, min_periods=0)\n                .over('symbol_id')\n                .cast(pl.Float32)\n                .alias(f\"{c}_rolling_{str(period)}_mean\")\n                for c in (agg_col if isinstance(agg_col, list) else [agg_col])\n            ]\n        ).select(group_col + [ f\"{c}_rolling_{str(period)}_mean\" for c in (agg_col if isinstance(agg_col, list) else [agg_col]) ] )\n\n        data = data.join(\n            rolling_df,\n            on=group_col,\n            how='left'\n        ) \n\n        return data\n\n    def create_sin_features(self, data: pl.LazyFrame) -> pl.LazyFrame:\n        return data.with_columns(\n            (pl.col('date_id') * (2 * np.pi / 5)).sin().cast(pl.Float32).alias('date_id_5_sin_feature'),\n            (pl.col('date_id') * (2 * np.pi / 21)).sin().cast(pl.Float32).alias('date_id_21_sin_feature'),\n             pl.when(pl.col('date_id') < 677)\n            .then((pl.col('time_id') * (2 * np.pi / 849)).sin().cast(pl.Float32))\n            .otherwise((pl.col('time_id') * (2 * np.pi / 968)).sin().cast(pl.Float32))\n            .alias('time_id_sin_feature')\n        )\n\n    def fill_null(self, data: pl.LazyFrame) -> pl.LazyFrame:\n        return data.fill_null(strategy=\"zero\") \n\n    def filter_from_date(self, data: pl.LazyFrame, day: int) -> pl.LazyFrame:\n        return data.filter(\n            (pl.col('date_id') >= day)\n        )\n    \n\n    def run(self, is_train=True):\n        \n        if is_train:\n\n            result = self.data.pipe(\n                self.fill_null\n            ).pipe(\n                self.add_lag_responder, [f'responder_{i}' for i in range(9)]\n            ).pipe(\n                self.create_daily_rolling_features, 5, ['date_id','symbol_id'], [f'responder_{i}_lag_1' for i in range(9)]\n            ).pipe(\n                self.drop_responders\n            ).pipe(\n                self.drop_partition_id\n            ).pipe(\n                self.create_sin_features\n            )\n        else:\n            result = self.data.pipe(\n                self.fill_null \n            ).pipe(\n                self.add_lag_responder, [f'responder_{i}' for i in range(9)]\n            ).pipe(\n                self.drop_partition_id\n            ).pipe(\n                self.create_sin_features\n            )\n        \n        return result if isinstance(result, pl.DataFrame) else result.collect()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-05T13:07:24.410060Z","iopub.execute_input":"2024-12-05T13:07:24.410427Z","iopub.status.idle":"2024-12-05T13:07:24.434879Z","shell.execute_reply.started":"2024-12-05T13:07:24.410394Z","shell.execute_reply":"2024-12-05T13:07:24.433605Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"class LGBMTrainer():\n    def __init__(\n        self, \n        data: pl.LazyFrame,\n        y: str,\n        train_window: int,\n        forecast_window: int,\n        lgbm_params: dict\n    ):\n        self.data = data\n        self.exog = []\n        self.y = y\n        self.train_window = train_window\n        self.forecast_window = forecast_window\n        self.lgbm_params = lgbm_params\n\n\n    def get_last_n_dates_per_symbol(self, data: pl.LazyFrame, train_size: int) -> pl.LazyFrame:\n        # Get unique dates per symbol\n        unique_dates = (\n            data\n            .select(['symbol_id', 'date_id'])\n            .unique()\n        )\n        \n        # Sort and get top N dates per symbol\n        top_dates = (\n            unique_dates\n            .sort(['symbol_id', 'date_id'], descending=[False, True])\n            .group_by('symbol_id')\n            .head(train_size)\n        )\n        \n        # Filter original data using these dates\n        final_result = (\n            data\n            .join(\n                top_dates,\n                on=['symbol_id', 'date_id'],\n                how='inner'\n            )\n            .sort(['symbol_id', 'date_id', 'time_id'])\n        )\n        \n        return final_result\n\n    def _calculate_r2(self, y_true, y_pred, weights):\n        \"\"\"\n        Calculate the sample weighted zero-mean R-squared score (R2).\n    \n        Parameters:\n        - y_true (pd.Series or np.array): Ground truth values.\n        - y_pred (pd.Series or np.array): Predicted values.\n        - weights (pd.Series or np.array): Sample weights.\n    \n        Returns:\n        - float: R2 score.\n        \"\"\"\n        numerator = np.sum(weights * (y_true - y_pred) ** 2)\n        denominator = np.sum(weights * (y_true ** 2))\n        r2_score = 1 - (numerator / denominator)\n        return r2_score\n\n    def train(self, train_df: pl.DataFrame, forecast_df: Union[pl.DataFrame,None]=None) -> lgb.Booster:\n        \n        self.exog = [c for c in train_df.columns if c not in ['date_id','row_id','weight','responder_6']]\n\n        train_ds = lgb.Dataset(\n            data = train_df[self.exog],\n            label = train_df[self.y].to_numpy()\n        )\n\n        forecast_ds = lgb.Dataset(\n            data = forecast_df[self.exog],\n            label = forecast_df[self.y].to_numpy()\n        )\n\n        booster = lgb.train(\n            params = self.lgbm_params,\n            train_set = train_ds,\n            valid_sets = [forecast_ds],\n        )\n\n        return booster\n        \n\n    def evaluate(self):\n        unique_date_ids = data.select('date_id').unique().sort('date_id').collect().to_numpy().flatten()\n\n        step = self.train_window + self.forecast_window\n        metric_list = []\n        r2_score_list = []\n        time_list = []\n        for i in range(0, len(unique_date_ids) - step):\n            start_time = time.perf_counter()\n            \n            forecast_dates = unique_date_ids[i+self.train_window:i+step]\n\n            forecast_df = self.data.filter(pl.col('date_id').is_in(forecast_dates))\n            train_df = self.data.pipe(\n                lambda x: x.filter(pl.col('date_id') < forecast_dates[0])\n            ).pipe(\n                self.get_last_n_dates_per_symbol, self.train_window\n            )\n\n            transform = FeatureEngineering(\n                train_df.collect()\n            )\n            train_df = transform.run()\n\n            transform = FeatureEngineering(\n                forecast_df.collect()\n            )\n            forecast_df = transform.run()\n            \n            booster = self.train(train_df, forecast_df)\n\n            metric = booster.best_score['valid_0']['rmse']\n\n            y_valid_pred = booster.predict(forecast_df[self.exog])\n            r2_score = self._calculate_r2(forecast_df[self.y].to_numpy(), y_valid_pred, forecast_df['weight'].to_numpy()) \n\n            metric_list.append(metric)\n            r2_score_list.append(r2_score)\n\n            end_time = time.perf_counter()\n            delta_time = end_time - start_time\n            time_list.append(delta_time)\n            print(f\"Elapsed time: {delta_time:.6f} seconds. Metric: {metric}. R2 Score: {r2_score}\")\n        \n\n        return metric_list, r2_score_list, time_list","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-05T13:07:24.436806Z","iopub.execute_input":"2024-12-05T13:07:24.437287Z","iopub.status.idle":"2024-12-05T13:07:24.462684Z","shell.execute_reply.started":"2024-12-05T13:07:24.437241Z","shell.execute_reply":"2024-12-05T13:07:24.461397Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"data = pl.scan_parquet('/kaggle/input/jane-street-real-time-market-data-forecasting/train.parquet')\ndata = data.filter(pl.col('date_id') > 1650) # CHANGE HERE TO 1000 OR ANY PRIOR TO THAT TO HAVE A BETTER EVALUATION\n\nlgbm_params = {   \n    'verbose':-1,\n    'boosting_type': 'gbdt',\n    'objective': 'regression_l2',\n    'metric': 'rmse',\n    'n_estimators': 40, \n    'learning_rate': 0.022167654221948718, \n    'max_depth': 5, \n    'num_leaves': 42, \n    'min_child_samples': 15, \n    'subsample': 0.8658662737679551, \n    'colsample_bytree': 0.6316624163538186, \n    'reg_alpha': 7.313049180654483, \n    'reg_lambda': 8.113493177057705,\n    'random_state': 42,\n    #'device': 'gpu'\n}\n\ntrainer = LGBMTrainer(\n    data = data,\n    y = 'responder_6',\n    train_window = 21,\n    forecast_window = 1,\n    lgbm_params = lgbm_params\n)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-05T13:07:24.465537Z","iopub.execute_input":"2024-12-05T13:07:24.466022Z","iopub.status.idle":"2024-12-05T13:07:24.512016Z","shell.execute_reply.started":"2024-12-05T13:07:24.465973Z","shell.execute_reply":"2024-12-05T13:07:24.511163Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"metric, r2_score, time = trainer.evaluate()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-05T13:07:24.512944Z","iopub.execute_input":"2024-12-05T13:07:24.513258Z","iopub.status.idle":"2024-12-05T13:13:22.465063Z","shell.execute_reply.started":"2024-12-05T13:07:24.513228Z","shell.execute_reply":"2024-12-05T13:13:22.464059Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"print(f'Mean R2 Score {np.mean(r2_score)}')\nprint(f'Mean Metric {np.mean(metric)}')\nprint(f'Mean Elapsed time {np.mean(time)}')","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-05T13:13:22.466310Z","iopub.execute_input":"2024-12-05T13:13:22.466716Z","iopub.status.idle":"2024-12-05T13:13:22.473121Z","shell.execute_reply.started":"2024-12-05T13:13:22.466673Z","shell.execute_reply":"2024-12-05T13:13:22.472163Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"config = {\n    'model': 'LGBM',\n    'dataset': '',\n    'save_code': False,\n    'note': ''\n}\n\nuser_secrets = UserSecretsClient()\nwandb_key = user_secrets.get_secret(\"wandb_key\")\n\nwandb.login(key=wandb_key)\nrun = wandb.init(project=\"JSRTMDF Online Learning\", job_type='Model training', config=config)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-05T13:14:03.534756Z","iopub.execute_input":"2024-12-05T13:14:03.535219Z","iopub.status.idle":"2024-12-05T13:14:10.337550Z","shell.execute_reply.started":"2024-12-05T13:14:03.535182Z","shell.execute_reply":"2024-12-05T13:14:10.336270Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"run.log(lgbm_params)\nrun.log({'R2':np.mean(r2_score)})\nrun.log({'RMSE':np.mean(metric)})\nrun.log({'Time Elapsed':np.mean(time)})\nwandb.finish()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-05T13:14:15.551845Z","iopub.execute_input":"2024-12-05T13:14:15.552691Z","iopub.status.idle":"2024-12-05T13:14:20.786639Z","shell.execute_reply.started":"2024-12-05T13:14:15.552650Z","shell.execute_reply":"2024-12-05T13:14:20.785350Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"import seaborn as sns \nimport matplotlib.pyplot as plt\n\ndf = pd.DataFrame({'R2_Score': r2_score, 'Metric': metric, 'Time': time})\n\n# Set up the figure and axes\nfig, axes = plt.subplots(1, 3, figsize=(15, 5))\n\n# Plot distributions\nsns.histplot(df['R2_Score'], ax=axes[0], kde=True)\naxes[0].set_title('R2 Score Distribution')\n\nsns.histplot(df['Metric'], ax=axes[1], kde=True)\naxes[1].set_title('Metric Distribution')\n\nsns.histplot(df['Time'], ax=axes[2], kde=True)\naxes[2].set_title('Time Distribution')\n\n# Adjust layout\nplt.tight_layout()\nplt.show()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-05T13:14:34.629295Z","iopub.execute_input":"2024-12-05T13:14:34.629690Z","iopub.status.idle":"2024-12-05T13:14:35.885800Z","shell.execute_reply.started":"2024-12-05T13:14:34.629654Z","shell.execute_reply":"2024-12-05T13:14:35.884776Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"plt.figure(figsize=(21, 5)) \nplt.plot(r2_score, marker='o')\nplt.title('R-squared Score Trend')\nplt.xlabel('X-axis Label')\nplt.ylabel('R2 Score')\nplt.grid(True)\nplt.show()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-05T13:14:35.887304Z","iopub.execute_input":"2024-12-05T13:14:35.887870Z","iopub.status.idle":"2024-12-05T13:14:36.175223Z","shell.execute_reply.started":"2024-12-05T13:14:35.887836Z","shell.execute_reply":"2024-12-05T13:14:36.174157Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null}]}