{"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":"nvidiaTeslaT4","dataSources":[{"sourceId":84493,"databundleVersionId":9871156,"sourceType":"competition"}],"dockerImageVersionId":30823,"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":true}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"!pip install polars[gpu] torch --extra-index-url=https://pypi.nvidia.com\n\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# 1) Loading the full training dataset (including what was previously val)\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# 2) Sort the entire dataset by date/time (no separate validation split)\ndf_selected = df_selected.sort([\"date_id\", \"time_id\"])\ndf_train = df_selected\n\n# 3) Create lag columns for responders\nall_cols = df_train.collect_schema().names()\ntarget_cols = [col for col in all_cols if 'responder' in col]\n\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\n# 4) Fill NA's (simple fill with 0). Exclude 'responder_*' and 'weight' from training\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_train = df_train.collect(engine='gpu')\n\n# 5) Scale the data in chunks\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()\n\n# Partial fit the scaler\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    scaler.partial_fit(chunk_array)\n    \n    del chunk_df, chunk_array\n    gc.collect()\n    start_idx = end_idx\n\n# Transform the data fully\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\n# Save the scaler\njoblib.dump(scaler, \"scaler.joblib\")\n\n# 6) Create a PyTorch Dataset with seq_length=1\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\nseq_length = 1\ntrain_dataset = NumPySequenceDataset(X_train_scaled, y_train, w_train, seq_length)\n\ndef weighted_mse_loss(y_pred, y_true, w):\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    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# 7) Define the LSTM model (2-layer, seq_length=1)\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,\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: [batch_size, seq_length=1, input_dim]\n        out, (hn, cn) = self.lstm(x)\n        out = out[:, -1, :]\n        out = self.fc(out)\n        return out\n\n# 8) DataLoader + Training Loop\nbatch_size = 1024\ntrain_loader = DataLoader(train_dataset, batch_size=batch_size, shuffle=True, num_workers=2)\n\ndevice = torch.device(\"cuda\")\nprint(\"Using device:\", device)\n\ninput_dim = len(cols_to_fill)\nhidden_dim = 32\ndropout = 0.4\nweight_decay = 1e-3\nmax_epochs = 50\nlr = 5e-5\n\nmodel = TwoLayerLSTM(\n    input_dim=input_dim,\n    hidden_dim=hidden_dim,\n    output_dim=1,\n    dropout=dropout\n).to(device)\n\noptimizer = optim.Adam(model.parameters(), lr=lr, weight_decay=weight_decay)\n\ndef train_one_epoch(model, loader, optimizer, device):\n    model.train()\n    running_loss = 0.0\n    for batch_x, batch_y, batch_w in 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    return running_loss / len(loader.dataset)\n\n##################################################################\n# Early stopping based on *training* Weighted R²\n##################################################################\nbest_train_r2 = float('-inf')\nno_improvement_count = 0\npatience = 3  # early stop if no improvement for 3 consecutive epochs\n\nfor epoch in range(1, max_epochs + 1):\n    print(f\"\\nEpoch {epoch}/{max_epochs}\")\n    train_loss = train_one_epoch(model, train_loader, optimizer, device)\n    train_r2 = compute_weighted_r2(model, train_loader, device)\n    print(f\"  Train Loss: {train_loss:.4f}, Weighted R²: {train_r2:.4f}\")\n\n    # Check if train R² improved\n    if train_r2 > best_train_r2:\n        best_train_r2 = train_r2\n        no_improvement_count = 0\n        torch.save(model.state_dict(), \"best_model_so_far.pth\")\n        print(f\"  (New best training R²: {best_train_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 \"\n                  f\"in training R² for {patience} consecutive epochs.\")\n            break\n\nprint(f\"Loading best model from final best training R²={best_train_r2:.4f}\")\nmodel.load_state_dict(torch.load(\"best_model_so_far.pth\"))\nprint(\"Final model ready!\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-01-13T22:10:08.517860Z","iopub.execute_input":"2025-01-13T22:10:08.518160Z"}},"outputs":[],"execution_count":null}]}