{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"name":"python","version":"3.10.12","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":30823,"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"!pip install polars[gpu] torch --extra-index-url=https://pypi.nvidia.com\n\n# Imports\nimport glob\nimport gc\nimport polars as pl\nimport numpy as np\nimport torch\nimport torch.nn as nn\nimport torch.optim as optim\nimport joblib\nfrom torch.utils.data import Dataset, DataLoader\nfrom sklearn.preprocessing import StandardScaler\nfrom tqdm import tqdm\n\n# Loading the training dataset\nclass LoadData:\n    def __init__(self, file_paths, partition_ids=None):\n        self.file_paths = file_paths\n        self.partition_ids = partition_ids\n\n    def load_and_concat(self):\n        if self.partition_ids is not None:\n            selected_files = [\n                fp for fp in self.file_paths\n                if any(f'partition_id={pid}' in fp for pid in self.partition_ids)\n            ]\n        else:\n            selected_files = self.file_paths\n\n        partitioned_data = [pl.scan_parquet(file_path) for file_path in selected_files]\n        df = pl.concat(partitioned_data, rechunk=False)\n        \n        del partitioned_data\n        gc.collect()\n        \n        return df\n\npartition_ids = [5, 6, 7, 8, 9]\nfile_paths_all = sorted(glob.glob('/kaggle/input/jane-street-real-time-market-data-forecasting/train.parquet/*/*.parquet'))\nloader = LoadData(file_paths=file_paths_all, partition_ids=partition_ids)\ndf_selected = loader.load_and_concat()\n\n# Sort the dataframe by date and time\ndf_selected = df_selected.sort([\"date_id\", \"time_id\"])\n\n# Split into Train/ Val based on date_id\nmax_date_id_df = df_selected.select(pl.col(\"date_id\").max()).collect()\nmax_date_id = max_date_id_df[\"date_id\"][0]\nsplit_index = max_date_id - 120\n\ndf_train = df_selected.filter(pl.col(\"date_id\") <= split_index)\ndf_val = df_selected.filter(pl.col(\"date_id\") > split_index)\n\ndel max_date_id_df\ngc.collect()\n\n# Shifting responders by 1 day and joining with the train and val datasets\nall_cols = df_train.collect_schema().names()\ntarget_cols = [col for col in all_cols if 'responder' in col]\ntrain_lags = (\n    df_train.select(\n        ['date_id', 'time_id', 'symbol_id'] + \n        [pl.col(col).shift().over(['symbol_id', 'time_id'])\n         .alias(f'lag_1_{col}') for col in target_cols]\n    )\n)\n\ndf_train = df_train.join(train_lags, on=['date_id', 'time_id', 'symbol_id'])\n\nval_lags = (\n    df_val.select(\n        ['date_id', 'time_id', 'symbol_id'] + \n        [pl.col(col).shift().over(['symbol_id', 'time_id'])\n         .alias(f'lag_1_{col}') for col in target_cols]\n    )\n)\n\ndf_val = df_val.join(val_lags, on=['date_id', 'time_id', 'symbol_id'])\n\n# Fill NA's with rolling mean values, set the rest of NA's to 0\n\nexcluded_features = [col for col in df_train.collect_schema().names() if col.startswith('responder_')] + ['weight']\ncols_to_fill = [col for col in df_train.collect_schema().names() if col not in excluded_features]\n\ndef filling_na(cols_to_fill, df: pl.DataFrame):\n    df = df.fill_null(0)\n    return df\n\ndf_train = filling_na(cols_to_fill, df_train)\ndf_val = filling_na(cols_to_fill, df_val)\n\ndf_train = df_train.collect(engine='gpu')\ndf_val = df_val.collect(engine='gpu')\n\nbatch_size = 1_000_000\nscaler = StandardScaler()\nn_train = df_train.height\n\ny_train = df_train.select(\"responder_6\").to_numpy()\nw_train = df_train.select(\"weight\").to_numpy()\ny_val = df_val.select(\"responder_6\").to_numpy()\nw_val = df_val.select(\"weight\").to_numpy()\n\nstart_idx = 0\nwhile start_idx < n_train:\n    end_idx = min(start_idx + batch_size, n_train)\n    \n    chunk_df = df_train.slice(start_idx, end_idx - start_idx)\n    chunk_array = chunk_df.select(cols_to_fill).to_numpy()\n    \n    scaler.partial_fit(chunk_array)\n    \n    del chunk_df, chunk_array\n    gc.collect()\n    \n    start_idx = end_idx\n\nn_features = len(cols_to_fill)\nX_train_scaled = np.empty((n_train, n_features), dtype=np.float32)\n\nstart_idx = 0\nwhile start_idx < n_train:\n    end_idx = min(start_idx + batch_size, n_train)\n    chunk_df = df_train.slice(start_idx, end_idx - start_idx)\n    \n    chunk_array = chunk_df.select(cols_to_fill).to_numpy()\n    scaled_array = scaler.transform(chunk_array)\n    \n    X_train_scaled[start_idx:end_idx] = scaled_array.astype(np.float32, copy=False)\n    \n    del chunk_df, chunk_array, scaled_array\n    gc.collect()\n    start_idx = end_idx\n\ndel df_train\ngc.collect()\n\nn_val = df_val.height\n\nstart_idx = 0\nX_val_scaled = np.empty((n_val, n_features), dtype=np.float32)\n\nwhile start_idx < n_val:\n    end_idx = min(start_idx + batch_size, n_val)\n    chunk_df = df_val.slice(start_idx, end_idx - start_idx)\n    \n    chunk_array = chunk_df.select(cols_to_fill).to_numpy()\n    scaled_array = scaler.transform(chunk_array)\n    \n    X_val_scaled[start_idx:end_idx] = scaled_array.astype(np.float32, copy=False)\n    \n    del chunk_df, chunk_array, scaled_array\n    gc.collect()\n    start_idx = end_idx\n\ndel df_val\ngc.collect()\n\njoblib.dump(scaler, \"scaler.joblib\")\n\n# Defining PyTorch Dataset and DataLoader\nclass NumPySequenceDataset(Dataset):\n    def __init__(self, X, y, w, seq_length):\n        self.X = X\n        self.y = y.reshape(-1)\n        self.w = w.reshape(-1)\n        self.seq_length = seq_length\n\n    def __len__(self):\n        return len(self.X) - self.seq_length + 1\n\n    def __getitem__(self, idx):\n        x_slice = self.X[idx : idx + self.seq_length]\n        y_label = self.y[idx + self.seq_length - 1]\n        w_label = self.w[idx + self.seq_length - 1]\n\n        x_slice = torch.tensor(x_slice, dtype=torch.float32)\n        y_label = torch.tensor(y_label, dtype=torch.float32)\n        w_label = torch.tensor(w_label, dtype=torch.float32)\n\n        return x_slice, y_label, w_label\n\ndef train_one_epoch(model, loader, optimizer, device):\n    model.train()\n    running_loss = 0.0\n\n    for batch_idx, (batch_x, batch_y, batch_w) in enumerate(tqdm(loader, desc=\"Train\", leave=False)):\n        batch_x, batch_y, batch_w = batch_x.to(device), batch_y.to(device), batch_w.to(device)\n        optimizer.zero_grad()\n        y_pred = model(batch_x).view(-1)\n        loss = weighted_mse_loss(y_pred, batch_y, batch_w)\n        loss.backward()\n        optimizer.step()\n        running_loss += loss.item() * len(batch_x)\n\n    train_loss = running_loss / len(loader.dataset)\n    return train_loss\n\nseq_length = 1\n\ntrain_dataset = NumPySequenceDataset(X_train_scaled, y_train, w_train, seq_length)\nval_dataset = NumPySequenceDataset(X_val_scaled, y_val, w_val, seq_length)\n\n# Weighted MSE + Weighted R²\ndef weighted_mse_loss(y_pred, y_true, w):\n    \"\"\"\n    Weighted MSE for training/backprop.\n    y_pred, y_true, w all shape: [batch_size].\n    \"\"\"\n    y_pred = y_pred.view(-1)\n    sq_err = (y_true - y_pred)**2\n    weighted_sq_err = w * sq_err\n    return torch.mean(weighted_sq_err)\n\ndef compute_weighted_r2(model, loader, device):\n    \"\"\"\n    Compute Weighted R² over an entire dataset (loader).\n    R² = 1 - (Sum(w*(y - pred)^2) / Sum(w*y^2))\n    \"\"\"\n    model.eval()\n    numerator = 0.0\n    denominator = 0.0\n    with torch.no_grad():\n        for batch_x, batch_y, batch_w in loader:\n            batch_x, batch_y, batch_w = batch_x.to(device), batch_y.to(device), batch_w.to(device)\n            y_pred = model(batch_x).view(-1)\n            numerator += torch.sum(batch_w * (batch_y - y_pred)**2).item()\n            denominator += torch.sum(batch_w * (batch_y**2)).item()\n    if denominator == 0:\n        return 0.0\n    return 1.0 - (numerator / denominator)\n\n# Two-Layer LSTM Model\nclass TwoLayerLSTM(nn.Module):\n    def __init__(self, input_dim, hidden_dim, output_dim=1, dropout=0.2):\n        super(TwoLayerLSTM, self).__init__()\n        self.lstm = nn.LSTM(\n            input_size=input_dim,\n            hidden_size=hidden_dim,\n            num_layers=2,  # two LSTM layers\n            batch_first=True,\n            dropout=dropout\n        )\n        self.fc = nn.Linear(hidden_dim, output_dim)\n\n    def forward(self, x):\n        # x shape: [batch_size, seq_length, input_dim]\n        out, (hn, cn) = self.lstm(x)\n        # out shape: [batch_size, seq_length, hidden_dim]\n        # Take the last time step\n        out = out[:, -1, :]\n        out = self.fc(out)  # shape: [batch_size, output_dim]\n        return out\n\n# Validation Step (similar to train_one_epoch, but no backprop)\ndef validate_one_epoch(model, loader, device):\n    model.eval()\n    running_loss = 0.0\n    with torch.no_grad():\n        for batch_x, batch_y, batch_w in loader:\n            batch_x, batch_y, batch_w = batch_x.to(device), batch_y.to(device), batch_w.to(device)\n            y_pred = model(batch_x).view(-1)\n            loss = weighted_mse_loss(y_pred, batch_y, batch_w)\n            running_loss += loss.item() * len(batch_x)\n\n    val_loss = running_loss / len(loader.dataset)\n    return val_loss\n    \n# Create DataLoaders + Model + Full Training Loop\nbatch_size = 1024\ntrain_loader = DataLoader(train_dataset, batch_size=batch_size, shuffle=True, num_workers=2)\nval_loader   = DataLoader(val_dataset,   batch_size=batch_size, shuffle=False, num_workers=2)\n\n# Set device (GPU if available)\ndevice = torch.device(\"cuda\")\nprint(\"Using device:\", device)\n\n# Model hyperparams\ninput_dim = len(cols_to_fill)\nhidden_dim = 32\ndropout = 0.4\nweight_decay = 1e-3\nepochs = 100\nlr = 5e-5\n\n# Instantiate the model\nmodel = TwoLayerLSTM(\n    input_dim=input_dim,\n    hidden_dim=hidden_dim,\n    output_dim=1,\n    dropout=dropout\n).to(device)\n\n# Define optimizer\noptimizer = optim.Adam(model.parameters(), lr=lr, weight_decay=weight_decay)\n\n# Additional variables for early stopping\nbest_val_r2 = float('-inf')\nbest_epoch = 0\nno_improvement_count = 0\npatience = 5\nmax_epochs = 150\n    \nfor epoch in range(1, max_epochs+1):\n    print(f\"\\nEpoch {epoch}/{max_epochs}\")\n    \n    # Train for one epoch\n    train_loss = train_one_epoch(model, train_loader, optimizer, device)\n    train_r2 = compute_weighted_r2(model, train_loader, device)\n    \n    # Validate\n    val_loss = validate_one_epoch(model, val_loader, device)\n    val_r2 = compute_weighted_r2(model, val_loader, device)\n\n    print(f\"  Train Loss: {train_loss:.4f}, R²: {train_r2:.4f}\")\n    print(f\"  Val   Loss: {val_loss:.4f}, R²: {val_r2:.4f}\")\n\n    # Check if Val R² improved\n    if val_r2 > best_val_r2:\n        best_val_r2 = val_r2\n        best_epoch = epoch\n        no_improvement_count = 0\n        torch.save(model.state_dict(), \"best_model_so_far.pth\")  # checkpoint\n        print(f\"  (New best Val R²: {best_val_r2:.4f} at epoch {epoch})\")\n    else:\n        no_improvement_count += 1\n        print(f\"  (No improvement, {no_improvement_count} epochs in a row)\")\n        \n        if no_improvement_count >= patience:\n            print(f\"\\nEarly stopping at epoch {epoch} due to no improvement in Val R² for {patience} epochs.\")\n            break\n\n# After training loop, roll back to best checkpoint\nprint(f\"Loading best model from epoch {best_epoch} (Val R²={best_val_r2:.4f}).\")\nmodel.load_state_dict(torch.load(\"best_model_so_far.pth\"))","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-01-12T23:39:28.086158Z","iopub.execute_input":"2025-01-12T23:39:28.086444Z","iopub.status.idle":"2025-01-13T03:12:12.598045Z","shell.execute_reply.started":"2025-01-12T23:39:28.086419Z","shell.execute_reply":"2025-01-13T03:12:12.597179Z"}},"outputs":[],"execution_count":null}]}