{"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":"gpu","dataSources":[{"sourceId":84493,"databundleVersionId":9871156,"sourceType":"competition"},{"sourceId":10178844,"sourceType":"datasetVersion","datasetId":6287307}],"dockerImageVersionId":30823,"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":true}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"# This Python 3 environment comes with many helpful analytics libraries installed\n# It is defined by the kaggle/python Docker image: https://github.com/kaggle/docker-python\n# For example, here's several helpful packages to load\n\nimport numpy as np # linear algebra\nimport pandas as pd # data processing, CSV file I/O (e.g. pd.read_csv)\n\n# Input data files are available in the read-only \"../input/\" directory\n# For example, running this (by clicking run or pressing Shift+Enter) will list all files under the input directory\nfrom os.path import join,exists\nimport os\nfor dirname, _, filenames in os.walk('/kaggle/input'):\n    for filename in filenames:\n        print(os.path.join(dirname, filename))\n\n# You can write up to 20GB to the current directory (/kaggle/working/) that gets preserved as output when you create a version using \"Save & Run All\" \n# You can also write temporary files to /kaggle/temp/, but they won't be saved outside of the current session","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","trusted":true,"execution":{"iopub.status.busy":"2025-01-12T09:49:46.075235Z","iopub.execute_input":"2025-01-12T09:49:46.075488Z","iopub.status.idle":"2025-01-12T09:49:46.435792Z","shell.execute_reply.started":"2025-01-12T09:49:46.075444Z","shell.execute_reply":"2025-01-12T09:49:46.435040Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"import os\nimport torch\nimport torch.nn as nn\nimport torch.nn.functional as F\nimport pandas as pd\nimport numpy as np\nimport pickle\nimport polars as pl\nfrom torch.utils.data import DataLoader, Dataset\nimport matplotlib.pyplot as plt\nfrom sklearn.metrics import mean_squared_error, r2_score\n\n\nCONFIG = {\n    \"seq_length\": 20,  # Length of input sequences for the model\n    \"batch_size\": 512,  # Number of samples per gradient update\n    \"learning_rate\": 0.0001,  # Learning rate for the optimizer\n    \"num_epochs\": 12,  # Number of training epochs\n    \"num_workers\": 4,  # Number of subprocesses to use for data loading\n    \"num_heads\": 8, \n    \"num_hash_buckets\": 64,\n    \"pin_memory\": True,  # Whether to pin memory for faster data transfer to GPU\n    \"prefetch_factor\": 16,  # Number of batches to prefetch\n    \"dropout_rate\": 0.2,  # Dropout rate for regularization\n    \"input_channels\": 79,  # Number of input features\n    \"output_features\": 1,  # Number of output features after CNN\n    \"lstm_hidden_size\": 128,  # Hidden size for LSTM\n    \"embedding_dim\": 16,  # Dimension of the embedding for symbol_id\n    \"weight_decay\": 0.0001, \n    # \"num_symbols\": 39,  # Number of unique symbol_ids (0 to 38)\n    \"feature_columns\": [f'feature_{i:02d}' for i in range(79)],  # List of feature column names\n    \"target_column\": \"responder_6\",  # Name of the target variable\n    \"data_path\": \"/kaggle/input/jane-street-real-time-market-data-forecasting/train.parquet/\",  # Path to the dataset\n    \"dataset_save_path\": \"/kaggle/working/processed_dataset.parquet\",  # Path to save the processed dataset\n    \"checkpoint_dir\": \"/kaggle/working/models\",  # Path to save the trained model\n}","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-01-12T09:49:46.436559Z","iopub.execute_input":"2025-01-12T09:49:46.436951Z","iopub.status.idle":"2025-01-12T09:49:50.340959Z","shell.execute_reply.started":"2025-01-12T09:49:46.436922Z","shell.execute_reply":"2025-01-12T09:49:50.340293Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"class AddNorm(nn.Module):\n    def __init__(self, size, dropout=0.1):\n        super(AddNorm, self).__init__()\n        self.norm = nn.LayerNorm(size)\n        self.dropout = nn.Dropout(dropout)\n\n    def forward(self, x, residual):\n        return self.norm(x + self.dropout(residual))\n\nclass GatedResidualNetwork(nn.Module):\n    def __init__(self, input_size):\n        super(GatedResidualNetwork, self).__init__()\n        self.linear1 = nn.Linear(input_size, input_size)\n        self.linear2 = nn.Linear(input_size, input_size)\n        self.activation = nn.ReLU()\n\n    def forward(self, x):\n        gate = torch.sigmoid(self.linear2(x))  # Gating mechanism\n        residual = self.linear1(x)  # Linear transformation\n        return gate * residual + x  # Residual connection\n\nclass AttCNNBiLSTM(nn.Module):\n    def __init__(self, input_channels, lstm_hidden_size, output_features, num_hash_buckets, embedding_dim, num_heads, dropout_rate=0.2):\n        super(AttCNNBiLSTM, self).__init__()\n        \n        # Embedding layer for symbol_id using hashing trick\n        self.num_hash_buckets = num_hash_buckets\n        self.embedding = nn.Embedding(num_hash_buckets, embedding_dim)  \n        \n        # Convolutional layers\n        self.conv1 = nn.Conv1d(in_channels=input_channels + embedding_dim, out_channels=128, kernel_size=3, padding=1)\n        self.bn1 = nn.BatchNorm1d(128)\n        #self.pool1 = nn.MaxPool1d(kernel_size=2)\n        self.dropout1 = nn.Dropout(dropout_rate)\n\n        self.conv2 = nn.Conv1d(in_channels=128, out_channels=64, kernel_size=3, padding=1)\n        self.bn2 = nn.BatchNorm1d(64)\n        #self.pool2 = nn.MaxPool1d(kernel_size=2)\n        self.dropout2 = nn.Dropout(dropout_rate)\n        \n        self.conv3 = nn.Conv1d(in_channels=64, out_channels=32, kernel_size=3, padding=1)\n        self.bn3 = nn.BatchNorm1d(32)\n        #self.pool3 = nn.MaxPool1d(kernel_size=2)\n        self.dropout3 = nn.Dropout(dropout_rate)\n        \n        # Additional GRN after CNN layers\n        #self.grn_cnn_to_lstm = GatedResidualNetwork(32)  # Input size matches the output of the last CNN layer\n        #self.add_norm_cnn_to_lstm = AddNorm(32)  # Add & Norm after GRN        \n\n        # LSTM layers\n        self.lstm = nn.LSTM(input_size=32, hidden_size=lstm_hidden_size, num_layers=2, batch_first=True, bidirectional=True)\n        self.dropout_lstm = nn.Dropout(dropout_rate)\n\n        # Multi-Head Attention\n        self.multihead_attn = nn.MultiheadAttention(embed_dim=lstm_hidden_size * 2, num_heads=num_heads, dropout=0.2, add_zero_attn=False)\n\n        # Gated Residual Network after Attention\n        self.grn1 = GatedResidualNetwork(lstm_hidden_size * 2)  # Input size matches the output of the attention layer\n        self.add_norm_attn = AddNorm(lstm_hidden_size * 2)  # Add & Norm after GRN\n\n        # Fully connected layer for prediction\n        #self.fc1 = nn.Linear(lstm_hidden_size * 2, 64)  # First fully connected layer\n        #self.fc2 = nn.Linear(64, 16)  # Second fully connected layer\n        self.fc_output = nn.Linear(lstm_hidden_size * 2, output_features)  # Final output layer\n\n    def forward(self, x, symbol_ids, return_features=False):\n        # Hashing trick to map symbol_ids to indices\n        hashed_indices = torch.fmod(symbol_ids, self.num_hash_buckets)  # Hashing trick\n        \n        # Embedding for hashed symbol_ids\n        embedded_symbols = self.embedding(hashed_indices)  # Shape: (batch_size, seq_length, embedding_dim)\n        embedded_symbols = embedded_symbols.permute(0, 2, 1)  # Shape: (batch_size, embedding_dim, seq_length)\n\n        # Concatenate embeddings with features\n        x = torch.cat((x, embedded_symbols), dim=1)  # Shape: (batch_size, input_channels + embedding_dim, seq_length)\n\n        # Apply Convolutional Layers\n        x = F.relu(self.bn1(self.conv1(x)))\n        #x = self.pool1(x)\n        x = self.dropout1(x)\n\n        x = F.relu(self.bn2(self.conv2(x)))\n        #x = self.pool2(x)\n        x = self.dropout2(x)\n        \n        x = F.relu(self.bn3(self.conv3(x)))\n        #x = self.pool3(x)\n        x = self.dropout3(x)\n        \n        # Apply Gated Residual Network\n        #x = self.grn_cnn_to_lstm(x.permute(0, 2, 1))  # Shape: (batch_size, seq_length, num_channels)\n        #x = x.permute(0, 2, 1)  # Shape: (batch_size, num_channels, seq_length)\n        \n        # Prepare for LSTM\n        x = x.permute(0, 2, 1)  # Shape: (batch_size, seq_length, num_channels)\n        \n        # Apply Add & Norm\n        #x = self.add_norm_cnn_to_lstm(x, x)  # Use the same x for residual connection\n\n        # Apply LSTM\n        lstm_out, _ = self.lstm(x)\n        lstm_out = self.dropout_lstm(lstm_out)\n\n        # Multi-Head Attention\n        attn_out, _ = self.multihead_attn(lstm_out, lstm_out, lstm_out)\n\n        # Apply Gated Residual Network and Add & Norm\n        grn_out = self.grn1(attn_out.mean(dim=1))  # Shape: (batch_size, 128)\n        grn_out = self.add_norm_attn(grn_out, attn_out.mean(dim=1))  # Apply Add & Norm\n\n        if return_features:\n            return grn_out\n\n        # Fully connected layers\n        #fc_out = F.relu(self.fc1(grn_out))  # First fully connected layer\n        #fc_out = F.relu(self.fc2(fc_out))  # Second fully connected layer\n\n        # Fully connected layer for prediction\n        output = self.fc_output(grn_out)  # Final output layer\n        return output","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-01-12T09:49:50.341689Z","iopub.execute_input":"2025-01-12T09:49:50.341997Z","iopub.status.idle":"2025-01-12T09:49:50.355086Z","shell.execute_reply.started":"2025-01-12T09:49:50.341977Z","shell.execute_reply":"2025-01-12T09:49:50.354090Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"class TimeSeriesDataset(Dataset):\n    def __init__(self, data, seq_length, feature_columns, target_column):\n        #unique_symbol_ids = data['symbol_id'].unique()  # Get unique symbol_ids\n        #sorted_symbol_ids = np.sort(unique_symbol_ids)\n        self.data = torch.tensor(data.select(feature_columns).to_numpy(), dtype=torch.float32)\n        self.targets = torch.tensor(data.select(target_column).to_numpy().flatten(), dtype=torch.float32)\n        self.symbol_ids = torch.tensor(data.select('symbol_id').to_numpy().flatten(), dtype=torch.long)\n        # self.time_ids = torch.tensor(data.select('time_id').to_numpy().flatten(), dtype=torch.long)\n        # self.date_ids = torch.tensor(data.select('date_id').to_numpy().flatten(), dtype=torch.long)\n        \n        # Ensure that the sequence length does not exceed the available data length\n        if seq_length > len(self.data):\n            raise ValueError(\"Sequence length must be less than or equal to the length of the data.\")\n\n        self.seq_length = seq_length\n\n    def __len__(self):\n        return len(self.data) - self.seq_length + 1\n\n    def __getitem__(self, idx):\n        seq = self.data[idx:idx + self.seq_length]  # Shape: (seq_length, num_features)\n        target = self.targets[idx + self.seq_length - 1]  # Shape: (1,)\n        symbol_ids = self.symbol_ids[idx:idx + self.seq_length]  # Get corresponding symbol_ids\n        # Ensure that seq and symbol_ids are of the same length\n        # if len(seq) != self.seq_length or len(symbol_ids) != self.seq_length:\n        #     raise ValueError(f\"Sequence length mismatch: seq length {len(seq)}, symbol_ids length {len(symbol_ids)}\")\n        return seq.permute(1, 0), target.unsqueeze(0), symbol_ids  # Return (num_features, seq_length), (1,), (seq_length,)\n\ndef weighted_r2_score(y_true, y_pred, weights):\n    # Convert to numpy arrays if they are PyTorch tensors\n    if isinstance(y_true, torch.Tensor):\n        y_true = y_true.cpu().detach().numpy()\n    if isinstance(y_pred, torch.Tensor):\n        y_pred = y_pred.cpu().detach().numpy()\n    if isinstance(weights, torch.Tensor):\n        weights = weights.cpu().detach().numpy()\n\n    # Ensure all arrays are 1D\n    y_true = y_true.flatten()\n    y_pred = y_pred.flatten()\n    weights = weights.flatten()\n\n    # Calculate the weighted mean of y_true\n    weighted_mean = np.sum(weights * y_true) / np.sum(weights) if np.sum(weights) != 0 else 0\n\n    # Calculate the weighted residuals\n    residuals = y_true - y_pred\n    weighted_residual_sum = np.sum(weights * (residuals ** 2))\n    weighted_total_sum = np.sum(weights * ((y_true - weighted_mean) ** 2))\n\n    # Calculate R²\n    r2 = 1 - (weighted_residual_sum / weighted_total_sum) if weighted_total_sum != 0 else 0\n    return r2","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-01-12T09:50:54.530833Z","iopub.execute_input":"2025-01-12T09:50:54.531185Z","iopub.status.idle":"2025-01-12T09:50:54.539086Z","shell.execute_reply.started":"2025-01-12T09:50:54.531162Z","shell.execute_reply":"2025-01-12T09:50:54.538178Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"data_folder = \"/kaggle/input/updated-janestreet-data/\"\nfeature_columns = CONFIG['feature_columns']\ntarget_column = CONFIG['target_column']\nseq_length = CONFIG['seq_length']\n\ndataframe_0 = pl.read_parquet(\"/kaggle/input/updated-janestreet-data/dataframe_(0,).parquet\")\ndata = TimeSeriesDataset(dataframe_0, seq_length, feature_columns, target_column)\n# Load a sample from each file\nfor i in range(1,39):\n    file_path = join(data_folder, f\"dataframe_({i},).parquet\")\n    chunk = pl.read_parquet(file_path)\n    temp_dataset = TimeSeriesDataset(chunk, seq_length, feature_columns, target_column)\n    data.data = torch.cat((data.data, temp_dataset.data),dim=0)\n    data.targets = torch.cat((data.targets, temp_dataset.targets),dim=0)\n    data.symbol_ids = torch.cat((data.symbol_ids, temp_dataset.symbol_ids),dim=0)\n\n# # Concatenate all samples into one DataFrame if needed\n# train_df = pl.concat(samples)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-01-12T09:50:58.342066Z","iopub.execute_input":"2025-01-12T09:50:58.342415Z","iopub.status.idle":"2025-01-12T09:52:33.667222Z","shell.execute_reply.started":"2025-01-12T09:50:58.342389Z","shell.execute_reply":"2025-01-12T09:52:33.666553Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"\n\n# Create dataset and dataloader\n# dataset = TimeSeriesDataset(data, seq_length, feature_columns, target_column)\ntrain_loader = DataLoader(data, batch_size=CONFIG['batch_size'], pin_memory=True, num_workers=4, shuffle=False)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-01-12T09:52:33.668158Z","iopub.execute_input":"2025-01-12T09:52:33.668366Z","iopub.status.idle":"2025-01-12T09:52:33.672267Z","shell.execute_reply.started":"2025-01-12T09:52:33.668348Z","shell.execute_reply":"2025-01-12T09:52:33.671529Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# Step 3: Train the AttCNNBiLSTM Model\ninput_channels = CONFIG['input_channels']\noutput_features = CONFIG['output_features']\nlstm_hidden_size = CONFIG['lstm_hidden_size']\nnum_hash_buckets = CONFIG['num_hash_buckets']\nembedding_dim = CONFIG['embedding_dim']\n# num_symbols = CONFIG['num_symbols']\nnum_heads = CONFIG['num_heads']\n\ndevice = torch.device(\"cuda\" if torch.cuda.is_available() else \"cpu\")  # Check if multiple GPUs are available\n\nmodel = AttCNNBiLSTM(input_channels, lstm_hidden_size, output_features, num_hash_buckets, embedding_dim, num_heads).to(device)\nif torch.cuda.device_count() > 1:\n    model = torch.nn.DataParallel(model)\n\ncriterion = torch.nn.SmoothL1Loss()\n#criterion = torch.nn.MSELoss()\noptimizer = torch.optim.Adam(model.parameters(), lr=CONFIG['learning_rate'])\n#optimizer = torch.optim.AdamW(model.parameters(), lr=CONFIG['learning_rate'], weight_decay=CONFIG['weight_decay'])\n#optimizer = torch.optim.SGD(model.parameters(), lr=CONFIG['learning_rate'], momentum=0.9, weight_decay=CONFIG['weight_decay'])\nfrom tqdm import tqdm  # Import tqdm for progress bar\n# Initialize variables for checkpointing\nbest_val_loss = float('inf')  # Set to infinity initially\ncheckpoint_dir = CONFIG['checkpoint_dir']  # Directory to save checkpoints\nos.makedirs(checkpoint_dir, exist_ok=True)  # Ensure the directory exists\nnum_epochs = CONFIG['num_epochs']\nfor epoch in range(num_epochs):\n    model.train()\n    epoch_loss = 0\n    total_batches = len(train_loader)\n    print(f'Epoch [{epoch + 1}/{num_epochs}], Total Batches: {total_batches}')\n    \n    # Training loop\n    for batch_X, batch_y, symbol_ids in tqdm(train_loader, desc=\"Processing Batches\", total=total_batches):\n        batch_X, batch_y, symbol_ids = batch_X.to(device), batch_y.to(device), symbol_ids.to(device)  # Move data to GPU\n        optimizer.zero_grad()\n        \n        outputs = model(batch_X, symbol_ids)  # Pass symbol_ids to the model\n        loss = criterion(outputs, batch_y)  # Compute loss\n        \n        if torch.isnan(loss).any():\n            print(\"Loss is NaN. Stopping training.\")\n            break\n        \n        loss.backward()  # Backward pass\n        torch.nn.utils.clip_grad_norm_(model.parameters(), max_norm=1.0)  # Gradient clipping\n        optimizer.step()  # Update weights\n        \n        epoch_loss += loss.item()\n    \n    train_loss = epoch_loss / len(train_loader)\n    print(f'Epoch [{epoch + 1}/{num_epochs}], Training Loss: {train_loss:.4f}')\n    \n    # # Validation loop\n    # model.eval()\n    # val_loss = 0\n    # with torch.no_grad():\n    #     for batch_X, batch_y, symbol_ids in val_loader:\n    #         batch_X, batch_y, symbol_ids = batch_X.to(device), batch_y.to(device), symbol_ids.to(device)\n    #         outputs = model(batch_X, symbol_ids)\n    #         loss = criterion(outputs, batch_y)\n    #         val_loss += loss.item()\n    \n    # val_loss /= len(val_loader)\n    # print(f'Epoch [{epoch + 1}/{num_epochs}], Validation Loss: {val_loss:.4f}')\n    \n    # # Checkpointing\n    # if val_loss < best_val_loss:\n    #     best_val_loss = val_loss\n    #     checkpoint_path = os.path.join(checkpoint_dir, f\"best_model_epoch_{epoch + 1}.pth\")\n    #     torch.save(model.state_dict(), checkpoint_path)\n    #     print(f'Checkpoint saved at {checkpoint_path}')\n    \n    # Save model after every epoch (optional)\n    epoch_checkpoint_path = os.path.join(checkpoint_dir, f\"model_epoch_{epoch + 1}.pth\")\n    torch.save(model.state_dict(), epoch_checkpoint_path)\n    print(f'Epoch checkpoint saved at {epoch_checkpoint_path}')","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-01-12T10:21:59.309037Z","iopub.execute_input":"2025-01-12T10:21:59.309386Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# Function to evaluate the model\ndef evaluate_model(model, data_loader, device):\n    model.eval()\n    model.to(device)\n    all_predictions = []\n    all_targets = []\n    \n    with torch.no_grad():\n        for batch_X, batch_y, symbol_ids in data_loader:\n            batch_X = batch_X.to(device)\n            batch_y = batch_y.to(device)\n            symbol_ids = symbol_ids.to(device)\n            outputs = model(batch_X, symbol_ids)  # Get model predictions\n            all_predictions.append(outputs.cpu().numpy())\n            all_targets.append(batch_y.cpu().numpy())\n    \n    # Concatenate all predictions and targets\n    all_predictions = np.concatenate(all_predictions)\n    all_targets = np.concatenate(all_targets)\n\n    return all_predictions, all_targets\n\n# Evaluate the model\npredictions, targets = evaluate_model(model, train_loader, device)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-01-12T09:49:51.152863Z","iopub.status.idle":"2025-01-12T09:49:51.153209Z","shell.execute_reply":"2025-01-12T09:49:51.153054Z"}},"outputs":[],"execution_count":null}]}