{"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"},{"sourceId":190709,"sourceType":"modelInstanceVersion","modelInstanceId":158927,"modelId":181312}],"dockerImageVersionId":30787,"isInternetEnabled":false,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"markdown","source":"# An approach to the Jane Street competition using an LSTM model","metadata":{"execution":{"iopub.status.busy":"2024-10-28T02:28:02.376236Z","iopub.execute_input":"2024-10-28T02:28:02.376643Z","iopub.status.idle":"2024-10-28T02:28:02.403841Z","shell.execute_reply.started":"2024-10-28T02:28:02.376593Z","shell.execute_reply":"2024-10-28T02:28:02.402907Z"}}},{"cell_type":"markdown","source":"Missing values in the dataset were imputed with column means and the data was scaled between -1 and 1 based on the observation that most features are ranged from -k to k. The non-temporal feature `symbol_id` was one-hot encoded and integrated into LSTM following the method outlined in [this StackOverflow post](https://stackoverflow.com/a/65411307).","metadata":{}},{"cell_type":"markdown","source":"Initially, I tried training the model with the full dataset from one single partition, but using all symbols was computationally expensive. I stopped it after three epochs, resulting in a negative R2 score. By testing with fewer symbols, I am able to run more epochs and observed an improvement with a positive R2 after a few more epochs. So I assume the approach should work. I'm planning to reduce the size of the training data using some data reduction methods like PCA, though I have some doubt on PCA to effectively capture the non-linear trends in time series. Any suggestions or comments are welcome!","metadata":{}},{"cell_type":"markdown","source":"## Some Useful Materials","metadata":{}},{"cell_type":"markdown","source":"- [How to save RAM memory](https://www.kaggle.com/code/pavansanagapati/14-simple-tips-to-save-ram-memory-for-1-gb-dataset)\n- [Combining auxiliary features with sequences](https://stackoverflow.com/a/65411307)\n- [Implement cond-rnn in PyTorch](https://github.com/philipperemy/cond_rnn/issues/48)\n- [A post about how to handle irregular sequence](https://stats.stackexchange.com/q/312609)","metadata":{}},{"cell_type":"code","source":"## Import library\nimport gc\nimport numpy as np\nimport pandas as pd\nimport polars as pl\nimport polars.selectors as cs\nimport torch\nimport torch.nn.functional as F\n\nfrom tqdm.auto import tqdm\nfrom pathlib import Path\nfrom sklearn.compose import ColumnTransformer\nfrom sklearn.model_selection import train_test_split\nfrom sklearn.preprocessing import OneHotEncoder\nfrom sklearn.decomposition import KernelPCA, PCA\nfrom torch import nn\nfrom torch.utils.data import Dataset, DataLoader","metadata":{"_uuid":"7d4fc9e5-4f8a-440a-a9a4-7eb486e670a0","_cell_guid":"e68f98b1-8043-4d4a-9ecd-39cd1968ba71","trusted":true,"execution":{"iopub.status.busy":"2024-12-06T16:05:03.994418Z","iopub.execute_input":"2024-12-06T16:05:03.994895Z","iopub.status.idle":"2024-12-06T16:05:09.746320Z","shell.execute_reply.started":"2024-12-06T16:05:03.994855Z","shell.execute_reply":"2024-12-06T16:05:09.745061Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"## Configuration\npl.Config.set_tbl_cols(-1)\nDATA_DIR = Path(\"/kaggle/input/jane-street-real-time-market-data-forecasting\")\n\nDATE = \"date_id\"\nTIME = \"time_id\"\nSYMBOL = \"symbol_id\"\nWEIGHT = \"weight\"\nUNUSED_RESPONDERS = [f\"responder_{i}\" for i in range(9) if i != 6]\nFEATURES = [f\"feature_{i:02d}\" for i in range(79)]\nTARGET = \"responder_6\"\nSYMBOL_CATEGORIES = list(range(39))\nDEVICE = \"cuda\" if torch.cuda.is_available() else \"cpu\"","metadata":{"_uuid":"6ad75245-3aca-4932-8346-90f2f1734250","_cell_guid":"4fae04e7-f275-4429-a9b1-dce005d029a7","trusted":true,"execution":{"iopub.status.busy":"2024-12-06T16:05:09.748510Z","iopub.execute_input":"2024-12-06T16:05:09.749159Z","iopub.status.idle":"2024-12-06T16:05:09.757569Z","shell.execute_reply.started":"2024-12-06T16:05:09.749108Z","shell.execute_reply":"2024-12-06T16:05:09.755792Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"## Custom function\ndef load_data(partition):\n    \"\"\"\n    Loads data from multiple parquet files based on specified partitions.\n    \"\"\"\n    # Convert single integer partition input to a list\n    partition = [partition] if isinstance(partition, int) else partition\n    \n    out = pl.DataFrame()\n    for i in partition:\n        path = DATA_DIR / \"train.parquet\" / f\"partition_id={i}\" / \"part-0.parquet\"\n        df = pl.read_parquet(path)\n        df = reduce_mem_usage(df)\n        out = pl.concat([out, df], how=\"diagonal_relaxed\")\n    return out\n\n\n\ndef train_valid_split(df: pl.DataFrame, valid_prop=0.2):\n    \"\"\"\n    Splits DataFrame into training and validation sets based on a specified\n    proportion of the most recent dates.\n\n    Returns:\n        - Training set containing earlier dates.\n        - Validation set containing the most recent dates.\n    \"\"\"\n    # Ensure chronological order\n    df = df.sort(cs.by_name(\"symbol_id\", \"date_id\", \"time_id\"))\n    \n    # How many date are needed for validation\n    distinct_dates = df.select(cs.by_name(\"date_id\").unique())\n    num_date = len(distinct_dates)\n    num_valid_date = int(num_date * valid_prop)\n    \n    # Get the starting date for validation set\n    valid_start_date = distinct_dates[-num_valid_date, :]\n    train_df = df.filter(~(cs.by_name(\"date_id\") >= valid_start_date))\n    valid_df = df.filter(cs.by_name(\"date_id\") >= valid_start_date)    \n    return train_df, valid_df\n\n\n\ndef reduce_mem_usage(df: pl.DataFrame):\n    for col in df.columns:\n        if df[col].dtype.is_numeric():\n            c_min, c_max = df[col].min(), df[col].max()\n        if c_min is None or c_max is None:\n            continue\n        if df[col].dtype.is_integer():\n            if c_min >= 0:\n                if c_min >= np.iinfo(np.uint8).min and c_max <= np.iinfo(np.uint8).max:\n                    df = df.with_columns(pl.col(col).cast(pl.UInt8))\n                elif c_min >= np.iinfo(np.uint16).min and c_max <= np.iinfo(np.uint16).max:\n                    df = df.with_columns(pl.col(col).cast(pl.UInt16))\n                elif c_min >= np.iinfo(np.uint32).min and c_max <= np.iinfo(np.uint32).max:\n                    df = df.with_columns(pl.col(col).cast(pl.UInt32))\n                elif c_min >= np.iinfo(np.uint64).min and c_max <= np.iinfo(np.uint64).max:\n                    df = df.with_columns(pl.col(col).cast(pl.UInt64))\n            else:\n                if c_min >= np.iinfo(np.int8).min and c_max <= np.iinfo(np.int8).max:\n                    df = df.with_columns(pl.col(col).cast(pl.Int8))\n                elif c_min >= np.iinfo(np.int16).min and c_max <= np.iinfo(np.int16).max:\n                    df = df.with_columns(pl.col(col).cast(pl.Int16))\n                elif c_min >= np.iinfo(np.int32).min and c_max <= np.iinfo(np.int32).max:\n                    df = df.with_columns(pl.col(col).cast(pl.Int32))\n                elif c_min >= np.iinfo(np.int64).min and c_max <= np.iinfo(np.int64).max:\n                    df = df.with_columns(pl.col(col).cast(pl.Int64))\n        if df[col].dtype.is_float():\n            if c_min > np.finfo(np.float32).min and c_max < np.finfo(np.float32).max:\n                df = df.with_columns(pl.col(col).cast(pl.Float32))\n            elif c_min > np.finfo(np.float64).min and c_max < np.finfo(np.float64).max:\n                df = df.with_columns(pl.col(col).cast(pl.Float64))\n    return df\n\n\ndef r2_score(y_true, y_pred, weight=None):\n    weight = weight if weight is not None else 1\n    return 1 - np.sum(weight * (y_true - y_pred) ** 2) / np.sum(weight * y_true**2)\n\n\ndef impute_by_mean(df: pl.DataFrame):\n    global colmeans\n    for col in df.columns:\n        if colmeans[col] is not None:\n            df = df.with_columns(cs.by_name(col).fill_null(colmeans[col]))\n        else:\n            df = df.drop(col)\n    return df\n\n\ndef scale_by_minmax(df: pl.DataFrame, new_range):\n    global colminmaxs\n    new_min, new_max = new_range\n    for col in df.columns:\n        if colminmaxs[col][1] - colminmaxs[col][0] != 0:\n            df = df.with_columns((cs.by_name(col) - colminmaxs[col][0]) / (colminmaxs[col][1] - colminmaxs[col][0]) * (new_max - new_min) + new_min)\n        else:\n            df = df.drop(col)\n    return df\n\n\ndef encode_by_onehot(df: pl.DataFrame):\n    \"\"\"\n    Based on predefined categories (0, 1, ..., 38).\n    \"\"\"\n    one_hot_columns = [\n        pl.when(pl.col(SYMBOL) == i)\n        .then(1)\n        .otherwise(0)\n        .alias(f\"symbol_id_{i}\")\n        for i in SYMBOL_CATEGORIES\n    ]\n\n    return df.with_columns(one_hot_columns).select(\n        cs.exclude(SYMBOL)\n    )","metadata":{"_uuid":"b3f0e72c-0b6f-457c-9ca4-ea5f17496cf5","_cell_guid":"b053cf25-3358-4729-9e61-528f1f51fd9b","trusted":true,"execution":{"iopub.status.busy":"2024-12-06T16:05:09.759181Z","iopub.execute_input":"2024-12-06T16:05:09.759696Z","iopub.status.idle":"2024-12-06T16:05:09.790165Z","shell.execute_reply.started":"2024-12-06T16:05:09.759650Z","shell.execute_reply":"2024-12-06T16:05:09.789011Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"class JSDataset(Dataset):\n    def __init__(self, df, seq_len):\n        self.seq_len = seq_len\n        self.df = df\n        \n    def __len__(self):\n        return len(self.df) - self.seq_len\n    \n    def __getitem__(self, idx):\n        w = torch.tensor(self.df[idx + self.seq_len, WEIGHT])\n        y = torch.tensor(self.df[idx + self.seq_len, TARGET])\n        X_ = self.df[idx : (idx + self.seq_len), :].select(cs.exclude(TARGET))\n        X = X_.select(~cs.starts_with(\"symbol_id_\")).to_torch().type(torch.float32)\n        c = X_.select(cs.starts_with(\"symbol_id\")).to_torch().type(torch.float32)      \n        return X, y, c, w","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-06T16:05:09.792242Z","iopub.execute_input":"2024-12-06T16:05:09.792652Z","iopub.status.idle":"2024-12-06T16:05:09.807757Z","shell.execute_reply.started":"2024-12-06T16:05:09.792603Z","shell.execute_reply":"2024-12-06T16:05:09.806355Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ## Custom Dataset\n# class JSDataset(Dataset):\n#     '''\n#     For loading and processing time series data for model training.\n#     '''\n#     DATE = \"date_id\"\n#     TIME = \"time_id\"\n#     SYMBOL = \"symbol_id\"\n#     WEIGHT = \"weight\"\n#     TARGET = \"responder_6\"\n#     FEATURES = [f\"feature_{i:02d}\" for i in range(79)]\n#     SYMBOL_CATEGORIES = list(range(39))  # For one-hot encoding\n\n#     def __init__(self, df, seq_len):\n#         self.seq_len = seq_len  # Length of each sequence sample (Each sample is a sequence of time step)\n        \n#         # Execute the transformation in __init__ instead of __getitem__\n#         df_to_scale = df.select(cs.exclude(\n#             self.__class__.DATE,\n#             self.__class__.TIME,\n#             self.__class__.SYMBOL,\n#             self.__class__.WEIGHT\n#         ))\n#         df_to_scale = self._impute_by_mean(df_to_scale)\n#         df_to_scale = self._scale_by_minmax(df_to_scale, new_range=(-1, 1))\n#         df_to_scale = pl.concat(\n#             [\n#                 df.select(cs.by_name(self.__class__.SYMBOL, self.__class__.WEIGHT)),\n#                 df_to_scale\n#             ],\n#             how=\"horizontal\"\n#         )\n#         self.df = self._onehot_encode_symbol(df_to_scale)\n\n#     def __len__(self):\n#         return len(self.df) - self.seq_len\n\n#     def __getitem__(self, idx):\n#         '''\n#         Returns:\n#             X (torch.Tensor): The sequence of feature values.\n#             y (torch.Tensor): Target variable for the sequence.\n#             c (torch.Tensor): One-hot encoded symbol information.\n#             w (torch.Tensor): Weight for R2 score.\n#         '''\n#         w = torch.tensor(self.df[idx + self.seq_len, self.__class__.WEIGHT])\n#         y = torch.tensor(self.df[idx + self.seq_len, self.__class__.TARGET])\n#         X_ = self.df[idx : (idx + self.seq_len), :]\n#         X = X_.select(~cs.starts_with(\"symbol_id_\")).to_torch().type(torch.float32)\n#         c = X_.select(cs.starts_with(\"symbol_id\")).to_torch().type(torch.float32)      \n#         return X, y, c, w\n\n#     def _impute_by_mean(self, df):\n#         '''\n#         Based on precomputed column mean.\n#         '''\n#         global colmeans\n#         for col in df.columns:\n#             if colmeans[col] is not None:\n#                 df = df.with_columns(cs.by_name(col).fill_null(colmeans[col]))\n#             else:\n#                 df = df.drop(cs.by_name(col))\n#         return df\n\n#     def _scale_by_minmax(self, df, new_range):\n#         '''\n#         Based on precomputed column minimum and column maximum.\n#         '''\n#         global colminmaxs\n#         new_min, new_max = new_range\n#         for col in df.columns:\n#             if colminmaxs[col][1] - colminmaxs[col][0] != 0:\n#                 df = df.with_columns(\n#                     (cs.by_name(col) - colminmaxs[col][0])\n#                     / (colminmaxs[col][1] - colminmaxs[col][0])\n#                     * (new_max - new_min)\n#                     + new_min\n#                 )\n#             else:\n#                 df = df.drop(cs.by_name(col))\n#         return df\n\n#     def _onehot_encode_symbol(self, df):\n#         \"\"\"\n#         Based on predefined categories (0, 1, ..., 38).\n#         \"\"\"\n#         one_hot_columns = [\n#             pl.when(pl.col(self.__class__.SYMBOL) == i)\n#             .then(1)\n#             .otherwise(0)\n#             .alias(f\"symbol_id_{i}\")\n#             for i in self.__class__.SYMBOL_CATEGORIES\n#         ]\n\n#         return df.with_columns(one_hot_columns).select(\n#             cs.exclude(self.__class__.SYMBOL)\n#         )","metadata":{"_uuid":"3c38e465-1f6e-4c7a-b0fe-4804139d7c90","_cell_guid":"e6ff1888-a30f-4c01-b670-1c9efa2f71a0","jupyter":{"source_hidden":true},"trusted":true,"execution":{"iopub.status.busy":"2024-12-06T16:05:09.809185Z","iopub.execute_input":"2024-12-06T16:05:09.809578Z","iopub.status.idle":"2024-12-06T16:05:09.826461Z","shell.execute_reply.started":"2024-12-06T16:05:09.809540Z","shell.execute_reply":"2024-12-06T16:05:09.825217Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"class JSNet(nn.Module):\n    def __init__(self, input_size, hidden_size, num_layers, output_size, cond_size):\n        super().__init__()\n        self.cond_layer = nn.Linear(cond_size, cond_size)\n        self.cond_active = nn.ReLU()\n        self.lstm1 = nn.LSTM(input_size, hidden_size, num_layers, batch_first=True)\n        self.lstm2 = nn.LSTM(cond_size + hidden_size, hidden_size, num_layers, batch_first=True)\n        self.output_layer = nn.Linear(hidden_size, output_size)\n\n    def forward(self, X, cond):\n        cond = self.cond_layer(cond)\n        cond = self.cond_active(cond)\n        X, _ = self.lstm1(X)\n        X = torch.cat((X, cond), dim=-1)\n        X, _ = self.lstm2(X)\n        return self.output_layer(X[:, -1, :]).squeeze()","metadata":{"_uuid":"281fc786-f5b6-469a-9a08-5fe0366a5bc5","_cell_guid":"eee8be47-c140-4802-b27b-8bffa416bc5a","collapsed":false,"jupyter":{"outputs_hidden":false},"trusted":true,"execution":{"iopub.status.busy":"2024-12-06T16:05:09.827712Z","iopub.execute_input":"2024-12-06T16:05:09.828029Z","iopub.status.idle":"2024-12-06T16:05:09.844405Z","shell.execute_reply.started":"2024-12-06T16:05:09.828000Z","shell.execute_reply":"2024-12-06T16:05:09.843336Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"def train_step(model, loss_fn, optimiser, dataloader, device):\n    model.to(device)\n    model.train()\n    train_loss = 0\n    for X, y, c, w in tqdm(dataloader):\n        X, y, c, w = X.to(device), y.to(device), c.to(device), w.to(device)\n        y_hat = model(X, c)\n        loss = loss_fn(y_hat, y)\n        train_loss += loss.item()\n        optimiser.zero_grad()\n        loss.backward()\n        optimiser.step()\n    train_loss /= len(dataloader)\n    return train_loss\n\n\ndef valid_step(model, loss_fn, dataloader, device):\n    model.to(device)\n    model.eval()\n    y_hats = []  # To store all predictions\n    ys = []  # To store all true labels\n    ws = []  # To store all weights\n    valid_loss = 0\n    \n    # Disable gradient computation for validation\n    with torch.inference_mode():\n        for X, y, c, w in tqdm(dataloader):\n            X, y, c, w = X.to(device), y.to(device), c.to(device), w.to(device)\n            y_hat = model(X, c)\n            loss = loss_fn(y_hat, y)\n            valid_loss += loss.item()\n            y_hats.append(y_hat.cpu().numpy())\n            ys.append(y.cpu().numpy())\n            ws.append(w.cpu().numpy())\n        valid_loss /= len(dataloader)\n        \n    # Concat into single arrays\n    y_hats = np.concatenate(y_hats)\n    ys = np.concatenate(ys)\n    ws = np.concatenate(ws)\n    r2 = r2_score(ys, y_hats, ws)\n    \n    return valid_loss, r2","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-06T16:05:09.845653Z","iopub.execute_input":"2024-12-06T16:05:09.845984Z","iopub.status.idle":"2024-12-06T16:05:09.864979Z","shell.execute_reply.started":"2024-12-06T16:05:09.845953Z","shell.execute_reply":"2024-12-06T16:05:09.863875Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"if __name__ == \"__main__\":\n\n    df = load_data(partition=9)\n    df = df.select(cs.exclude(UNUSED_RESPONDERS))\n#     df = df.filter(cs.by_name(\"symbol_id\").is_in([0, 1, 2]))  # Just take a few symbols to test combining non-temporal features\n    df_train, df_valid = train_valid_split(df, valid_prop=0.2)\n    del df\n    gc.collect()\n\n    \n    cols_to_preprocess = FEATURES\n    cols_to_skip = [DATE, TIME, WEIGHT]\n    colmeans = df_train.select(cols_to_preprocess).mean()\n    colmeans = {col: colmeans[col].item() for col in colmeans.columns}\n    colminmaxs = pl.concat(\n        [\n            df_train.select(cols_to_preprocess).min(),\n            df_train.select(cols_to_preprocess).max(),\n        ],\n        how=\"vertical_relaxed\",\n    )\n    colminmaxs = {col: tuple(colminmaxs[col]) for col in colminmaxs.columns}\n    \n    \n    df_train_tmp = scale_by_minmax(impute_by_mean(df_train.select(cols_to_preprocess)), new_range=(-1, 1))\n    pca = PCA(n_components=0.95)\n    pca.fit(df_train_tmp)\n    df_train_tmp = pca.transform(df_train_tmp)\n    df_train = pl.concat(\n        [\n            df_train.select(cols_to_skip),\n            encode_by_onehot(df_train.select(pl.col(SYMBOL))),\n            pl.DataFrame(df_train_tmp, schema=[f\"e_vector_{i}\" for i in range(df_train_tmp.shape[1])]),\n            df_train.select(TARGET)\n        ], how=\"horizontal\"\n    )\n    del df_train_tmp\n    \n    \n    df_valid_tmp = scale_by_minmax(impute_by_mean(df_valid.select(cols_to_preprocess)), new_range=(-1, 1))\n    df_valid_tmp = pca.transform(df_valid_tmp)\n    df_valid = pl.concat(\n        [\n            df_valid.select(cols_to_skip),\n            encode_by_onehot(df_valid.select(pl.col(SYMBOL))),\n            pl.DataFrame(df_valid_tmp, schema=[f\"e_vector_{i}\" for i in range(df_valid_tmp.shape[1])]),\n            df_train.select(TARGET)\n        ], how=\"horizontal\"\n    )\n    del df_valid_tmp\n\n\n    seq_len = 1\n    js_train = JSDataset(df_train, seq_len)\n    js_valid = JSDataset(df_valid, seq_len)\n    \n    single_data = next(iter(js_train))\n    print(f\"Shape of a single sample\")\n    print(\"------------------------------------------------\")\n    print(f\"X: {single_data[0].shape}\")\n    print(f\"y: {single_data[1].shape}\")\n    print(f\"c: {single_data[2].shape}\")\n    print(f\"w: {single_data[3].shape}\")\n    print(\"------------------------------------------------\")\n\n    \n    batch_size = 64\n    train_dataloader = DataLoader(js_train, batch_size=batch_size, shuffle=False)\n    valid_dataloader = DataLoader(js_valid, batch_size=batch_size, shuffle=False)\n\n    input_size = next(iter(js_train))[0].shape[1]\n    hidden_size = 50\n    cond_size = next(iter(js_train))[2].shape[1]\n    output_size = 1\n    num_layers = 1\n    \n    # model = JSNet(input_size, hidden_size, num_layers, output_size, cond_size).to(DEVICE)\n    # loss_fn = nn.MSELoss()\n    # optimiser = torch.optim.Adam(params=model.parameters())\n    \n    # epochs = 1\n    # for epoch in tqdm(range(epochs)):\n    #     print(f\"Epoch: {epoch}\")\n    #     print(\"Inside training loop...\")\n    #     print(\"-----------------------------------\")\n    #     train_loss = train_step(model, loss_fn, optimiser, train_dataloader, DEVICE)\n    #     print(\"Inside testing loop...\")\n    #     print(\"-----------------------------------\")\n    #     valid_loss, r2 = valid_step(model, loss_fn, valid_dataloader, DEVICE)\n    #     print(f\"Train loss : {train_loss:.4f}\")\n    #     print(f\"Test loss  : {valid_loss:.4f}\")\n    #     print(f\"R2 Score   : {r2:.4f}\")","metadata":{"_uuid":"73215b6c-91d9-47e1-9452-4c9f8542da58","_cell_guid":"0883ea40-c591-41e6-ae35-cc02d8cfaf24","trusted":true,"execution":{"iopub.status.busy":"2024-12-06T16:08:01.722003Z","iopub.execute_input":"2024-12-06T16:08:01.722634Z","iopub.status.idle":"2024-12-06T16:09:13.029959Z","shell.execute_reply.started":"2024-12-06T16:08:01.722592Z","shell.execute_reply":"2024-12-06T16:09:13.028682Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"model = torch.load(\"/kaggle/input/js24/pytorch/default/2/checkpoint.pth\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-06T16:09:23.799958Z","iopub.execute_input":"2024-12-06T16:09:23.800420Z","iopub.status.idle":"2024-12-06T16:09:23.810934Z","shell.execute_reply.started":"2024-12-06T16:09:23.800384Z","shell.execute_reply":"2024-12-06T16:09:23.809549Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## Submission","metadata":{}},{"cell_type":"code","source":"import kaggle_evaluation.jane_street_inference_server\nimport os\n\nlags_ : pl.DataFrame | None = None\n\n\n# Replace this function with byour inference code.\n# You can return either a Pandas or Polars dataframe, though Polars is recommended.\n# Each batch of predictions (except the very first) must be returned within 1 minute of the batch features being provided.\ndef predict(test: pl.DataFrame, lags: pl.DataFrame | None) -> pl.DataFrame | pd.DataFrame:\n    \"\"\"Make a prediction.\"\"\"\n    # All the responders from the previous day are passed in at time_id == 0. We save them in a global variable for access at every time_id.\n    # Use them as extra features, if you like.\n    global lags_, cols_to_preprocess, cols_to_skip, colmeans, colminmaxs, pca, model, seq_len\n    if lags is not None:\n        lags_ = lags\n    \n    # Replace this section with your own predictions\n    predictions = test.select(\n        'row_id',\n        pl.lit(0.0).alias('responder_6'),\n    )\n\n    test = test.sort(cs.by_name(\"symbol_id\", \"date_id\", \"time_id\")).with_columns(pl.lit(0.0).alias('responder_6'))\n    test_tmp = scale_by_minmax(impute_by_mean(test.select(cols_to_preprocess)), new_range=(-1, 1))\n    test_tmp = pca.transform(test_tmp)\n    df_test = pl.concat(\n        [\n            test.select(cols_to_skip),\n            encode_by_onehot(test.select(pl.col(SYMBOL))),\n            pl.DataFrame(test_tmp, schema=[f\"e_vector_{i}\" for i in range(test_tmp.shape[1])]),\n            test.select(TARGET)\n        ], how=\"horizontal\"\n    )\n    del test_tmp\n\n    test_dataset = JSDataset(df_test, seq_len)\n    test_dataloader = DataLoader(test_dataset, batch_size=1, shuffle=False)\n\n    model.to(DEVICE)\n    model.eval()\n    y_hats = []\n    with torch.inference_mode():\n        for X, _, c, _ in test_dataloader:\n            X, c = X.to(DEVICE), c.to(DEVICE)\n            y_hat = model(X, c)\n            y_hats.append(y_hat.cpu().numpy())\n    y_hats = [float(item) for item in y_hats]\n    y_hats.insert(0, np.mean(y_hats))\n\n    predictions = predictions.with_columns(pl.Series(name=\"responder_6\", values=y_hats))\n    # print(predictions)\n\n    if isinstance(predictions, pl.DataFrame):\n        assert predictions.columns == ['row_id', 'responder_6']\n    elif isinstance(predictions, pd.DataFrame):\n        assert (predictions.columns == ['row_id', 'responder_6']).all()\n    else:\n        raise TypeError('The predict function must return a DataFrame')\n    # Confirm has as many rows as the test data.\n    assert len(predictions) == len(test)\n\n    return predictions\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-06T16:43:47.174303Z","iopub.execute_input":"2024-12-06T16:43:47.174718Z","iopub.status.idle":"2024-12-06T16:43:47.188352Z","shell.execute_reply.started":"2024-12-06T16:43:47.174684Z","shell.execute_reply":"2024-12-06T16:43:47.186815Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"inference_server = kaggle_evaluation.jane_street_inference_server.JSInferenceServer(predict)\n\nif os.getenv('KAGGLE_IS_COMPETITION_RERUN'):\n    inference_server.serve()\nelse:\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":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-06T16:43:56.493886Z","iopub.execute_input":"2024-12-06T16:43:56.494361Z","iopub.status.idle":"2024-12-06T16:43:56.694090Z","shell.execute_reply.started":"2024-12-06T16:43:56.494321Z","shell.execute_reply":"2024-12-06T16:43:56.693132Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# model_path = Path('/kaggle/working/models')\n# model_path.mkdir(parents=True, exist_ok=True)\n# checkpoint = {'epoch': epoch,\n#               'state_dict': model.state_dict(),\n#               'optimizer' : optimiser.state_dict()}\n\n# torch.save(checkpoint, model_path / \"checkpoint.pth\")","metadata":{"_uuid":"e5002a9f-cc89-4262-beb3-e5b13c9d3b28","_cell_guid":"804c8738-3d0a-4276-afea-0033837edd28","collapsed":false,"jupyter":{"outputs_hidden":false},"trusted":true,"execution":{"iopub.status.busy":"2024-12-03T15:56:49.636929Z","iopub.status.idle":"2024-12-03T15:56:49.637466Z","shell.execute_reply.started":"2024-12-03T15:56:49.637169Z","shell.execute_reply":"2024-12-03T15:56:49.637196Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# df.select(\"date_id\").min(), df.select(\"date_id\").max()\n\n# # Num. of missing values\n# null_counts = df.select(cs.all().is_null().sum())\n# null_prop = null_counts/len(df)\n\n# colmeans = df.mean()\n# colsds = df.std()\n\n# def impute_miss(df: pl.DataFrame, colmeans):\n#     if df.shape[1] != len(colmeans):\n#         raise ValueError(\"Column count in DataFrame does not match length of colmeans.\")\n#     for i, mean in enumerate(colmeans):\n#         if mean is not None:\n#             df = df.with_columns(cs.by_index(i).fill_null(mean))\n#     return df\n\n# def z_scale(df: pl.DataFrame, colmeans, colsds):\n#     for i, col in enumerate(df.columns):\n#         mean = colmeans[i]\n#         sd = colsds[i]\n#         if sd != 0:\n#             df = df.with_columns((cs.by_name(col)-mean)/sd)\n#     return df\n\n# df.describe()\n\n# impute_miss(df, colmeans.row(0))\n\n# class JSDataset(Dataset):\n#     DATA_DIR = Path(\"/kaggle/input/jane-street-real-time-market-data-forecasting\")\n#     PREDICTORS = [\"date_id\", \"time_id\", \"symbol_id\", \"weight\"] + [f\"feature_{i:02d}\"for i in range(79)] + [\"responder_6\"]\n#     TARGET = \"responder_6\"\n    \n#     def __init__(self, seq_len, partition):\n#         self.seq_len = seq_len\n#         self.partition = [partition] if isinstance(partition, int) else partition\n#         self.data = self._load_data()\n#         self.X = self.data.select(pl.col(self.__class__.PREDICTORS))\n#         self.y = self.data.select(pl.col(self.__class__.TARGET))\n        \n#     def __len__(self):\n#         return len(self.y) - self.seq_len\n    \n#     def __getitem__(self, idx):M\n#         return self.X[idx:(idx + self.seq_len), :].to_torch(), self.y[idx + self.seq_len].to_torch().squeeze()\n    \n#     def _load_data(self):\n#         out = pl.DataFrame()\n#         for i in self.partition:\n#             path = self.__class__.DATA_DIR / \"train.parquet\" / f\"partition_id={i}\" / \"part-0.parquet\"\n#             df = pl.read_parquet(path)\n#             df = self._reduce_mem_usage(df)\n#             out = pl.concat([out, df], how=\"diagonal_relaxed\")\n#         return out\n            \n#     def _reduce_mem_usage(self, df: pl.DataFrame):\n#         for col in df.columns:\n#             if df[col].dtype.is_numeric():\n#                 c_min, c_max = df[col].min(), df[col].max()\n#             if c_min is None or c_max is None:\n#                 continue\n#             if df[col].dtype.is_integer():\n#                 if c_min >= 0:\n#                     if c_min >= np.iinfo(np.uint8).min and c_max <= np.iinfo(np.uint8).max:\n#                         df = df.with_columns(pl.col(col).cast(pl.UInt8))\n#                     elif c_min >= np.iinfo(np.uint16).min and c_max <= np.iinfo(np.uint16).max:\n#                         df = df.with_columns(pl.col(col).cast(pl.UInt16))\n#                     elif c_min >= np.iinfo(np.uint32).min and c_max <= np.iinfo(np.uint32).max:\n#                         df = df.with_columns(pl.col(col).cast(pl.UInt32))\n#                     elif c_min >= np.iinfo(np.uint64).min and c_max <= np.iinfo(np.uint64).max:\n#                         df = df.with_columns(pl.col(col).cast(pl.UInt64))\n#                 else:\n#                     if c_min >= np.iinfo(np.int8).min and c_max <= np.iinfo(np.int8).max:\n#                         df = df.with_columns(pl.col(col).cast(pl.Int8))\n#                     elif c_min >= np.iinfo(np.int16).min and c_max <= np.iinfo(np.int16).max:\n#                         df = df.with_columns(pl.col(col).cast(pl.Int16))\n#                     elif c_min >= np.iinfo(np.int32).min and c_max <= np.iinfo(np.int32).max:\n#                         df = df.with_columns(pl.col(col).cast(pl.Int32))\n#                     elif c_min >= np.iinfo(np.int64).min and c_max <= np.iinfo(np.int64).max:\n#                         df = df.with_columns(pl.col(col).cast(pl.Int64))\n#             if df[col].dtype.is_float():\n#                 if c_min > np.finfo(np.float32).min and c_max < np.finfo(np.float32).max:\n#                     df = df.with_columns(pl.col(col).cast(pl.Float32))\n#                 elif c_min > np.finfo(np.float64).min and c_max < np.finfo(np.float64).max:\n#                     df = df.with_columns(pl.col(col).cast(pl.Float64))\n#         return df\n\n\n# ## '''\n# device = \"cuda\" if torch.cuda.is_available() else \"cpu\"\n# device\n\n# predictors = df_train.columns\n# target = \"responder_6\"\n\n# class JSDataset(Dataset):\n#     def __init__(self, data, seq_len):\n#         self.data = data\n#         self.seq_len = seq_len\n#         self.X = self.data.select(pl.col(predictors))\n#         self.y = self.data.select(pl.col(target))\n#         self.n_samples = len(self.data)\n\n#     def __getitem__(self, idx):\n#         X_seq = self.X[idx: idx+self.seq_len]\n#         y_seq = self.y[idx+self.seq_len]\n#         return X_seq.to_torch(), y_seq.to_torch()\n\n#     def __len__(self):\n#         return self.n_samples - self.seq_len\n    \n# seq_len = 50\n# train_dataset = JSDataset(df_train, seq_len)\n# test_dataset = JSDataset(df_test, seq_len)\n\n# batch_size = 64\n# train_dataloader = DataLoader(train_dataset, batch_size=batch_size)\n# test_dataloader = DataLoader(test_dataset, batch_size=batch_size)\n\n# class JSNet(nn.Module):\n#     def __init__(self, input_size, hidden_size, num_layers, output_size):\n#         super().__init__()\n#         self.lstm = nn.LSTM(input_size, hidden_size, num_layers, batch_first=True)\n# #         self.fc1 = nn.Linear(hidden_size, output_size)\n#         self.fc1 = nn.Linear(hidden_size, 512)\n#         self.fc2 = nn.Linear(512, output_size)\n\n#     def forward(self, x):\n#         x, _ = self.lstm(x)\n# #         return self.fc1(x[:, -1, :])\n#         return self.fc2(self.fc1(x[:, -1, :]))\n    \n# input_size = df_train.shape[1]\n# hidden_size = 256\n# num_layers = 1\n# output_size = 1\n\n# net = JSNet(input_size, hidden_size, num_layers, output_size).to(device)\n# loss_fn = nn.MSELoss()\n# optimiser = torch.optim.Adam(params=net.parameters(), lr=0.05)\n\n# epochs = 3\n\n# for epoch in range(epochs):\n#     print(f\"Running Epoch:{epoch}\")\n#     print(\"----------------------------\")\n#     train_loss = 0\n#     for batch, (X, y) in enumerate(tqdm(train_dataloader)):\n#         X, y = X.to(device), y.to(device)\n#         net.train()\n#         y_pred = net(X.type(torch.float32))\n#         loss = loss_fn(y_pred, y.squeeze(1))\n#         train_loss += loss\n#         optimiser.zero_grad()\n#         loss.backward()\n#         optimiser.step()\n#         if batch % 100 == 0:\n#             print(f\"No. used samples for training: {batch * len(X)}/{len(train_dataloader.dataset)}\")\n            \n#     train_loss /= len(train_dataloader)\n\n#     test_loss, test_r2_score = 0, 0\n#     net.eval()\n#     with torch.inference_mode():\n#         for X, y in test_dataloader:\n#             X, y = X.to(device), y.to(device)\n#             test_pred = net(X)\n#             test_loss += loss_fn(test_pred, y)\n#             test_r2_score += r2_score(y, test_pred)\n        \n#         test_loss /= len(test_dataloader)\n#         test_r2_score /= len(test_dataloader)\n        \n#     print(f\"Train loss: {train_loss:.4f} | Test loss: {test_loss:.4f} | R2: {test_r2_score:.4f}\")\n\n# print(train_loss)\n# print(test_loss)\n# print(test_r2_score)\n# '''","metadata":{"_uuid":"fb0af3f7-a76a-4e5a-8c94-8378440a94bb","_cell_guid":"c79ac993-5673-4a22-9731-b355ade00cb6","collapsed":false,"jupyter":{"outputs_hidden":false},"trusted":true,"execution":{"iopub.status.busy":"2024-12-03T15:56:49.639741Z","iopub.status.idle":"2024-12-03T15:56:49.640244Z","shell.execute_reply.started":"2024-12-03T15:56:49.639984Z","shell.execute_reply":"2024-12-03T15:56:49.640010Z"}},"outputs":[],"execution_count":null}]}