{"metadata":{"kernelspec":{"display_name":"Python 3","language":"python","name":"python3"},"language_info":{"codemirror_mode":{"name":"ipython","version":3},"file_extension":".py","mimetype":"text/x-python","name":"python","nbconvert_exporter":"python","pygments_lexer":"ipython3","version":"3.10.14"},"kaggle":{"accelerator":"none","dataSources":[{"sourceId":84493,"databundleVersionId":9871156,"sourceType":"competition"},{"sourceId":213604,"sourceType":"modelInstanceVersion","isSourceIdPinned":true,"modelInstanceId":182075,"modelId":204305}],"dockerImageVersionId":30786,"isInternetEnabled":false,"language":"python","sourceType":"notebook","isGpuEnabled":false},"papermill":{"default_parameters":{},"duration":7.594014,"end_time":"2024-10-10T11:58:36.355301","environment_variables":{},"exception":null,"input_path":"__notebook__.ipynb","output_path":"__notebook__.ipynb","parameters":{},"start_time":"2024-10-10T11:58:28.761287","version":"2.6.0"}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"import os\n\nimport pandas as pd\nimport polars as pl\n\nimport kaggle_evaluation.jane_street_inference_server\n\nimport numpy as np\nimport pandas as pd\nfrom sklearn.model_selection import train_test_split\nimport math\nfrom scipy.stats import entropy\nimport matplotlib.pyplot as plt\nfrom matplotlib.dates import MonthLocator, DateFormatter\nimport torch\nimport torch.nn as nn\nimport torch.nn.functional as F\nfrom sklearn.preprocessing import StandardScaler\nfrom torch.utils.data import TensorDataset, DataLoader\nimport torch.optim as optim\nimport time\n\nimport os\nimport polars as pl\nimport sys\nfrom typing import Callable, Dict, List, Optional","metadata":{"_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","execution":{"iopub.status.busy":"2025-01-03T09:01:58.379388Z","iopub.execute_input":"2025-01-03T09:01:58.379823Z","iopub.status.idle":"2025-01-03T09:02:04.411753Z","shell.execute_reply.started":"2025-01-03T09:01:58.379761Z","shell.execute_reply":"2025-01-03T09:02:04.410770Z"},"papermill":{"duration":1.594977,"end_time":"2024-10-10T11:58:33.569648","exception":false,"start_time":"2024-10-10T11:58:31.974671","status":"completed"},"tags":[],"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"The evaluation API requires that you set up a server which will respond to inference requests. We have already defined the server; you just need write the predict function. When we evaluate your submission on the hidden test set the client defined in `jane_street_gateway` will run in a different container with direct access to the hidden test set and hand off the data timestep by timestep.\n\n\n\nYour code will always have access to the published copies of the files.","metadata":{"papermill":{"duration":0.00263,"end_time":"2024-10-10T11:58:33.575645","exception":false,"start_time":"2024-10-10T11:58:33.573015","status":"completed"},"tags":[]}},{"cell_type":"code","source":"\nlags_ : pl.DataFrame | None = None\n\nimport polars as pl\n\ndef lambda_init_fn(depth):\n    return 0.8 - 0.6 * math.exp(-0.3 * depth)\n\n\n\n\nclass DifferentialAttention(nn.Module):\n    def __init__(self, dim_model: int, head_nums: int, depth: int):\n        super().__init__()\n        \n        self.head_dim = dim_model // head_nums\n        self.Q = nn.Linear(dim_model, 2 * self.head_dim, bias=False)\n        self.K = nn.Linear(dim_model, 2 * self.head_dim, bias=False)\n        self.V = nn.Linear(dim_model, 2 * self.head_dim, bias=False)\n        self.scale = self.head_dim ** -0.5\n        self.depth = depth\n        self.lambda_q1 = nn.Parameter(torch.zeros(self.head_dim, dtype=torch.float32).normal_(mean=0,std=0.1))\n        self.lambda_q2 = nn.Parameter(torch.zeros(self.head_dim, dtype=torch.float32).normal_(mean=0,std=0.1))\n        self.lambda_k1 = nn.Parameter(torch.zeros(self.head_dim, dtype=torch.float32).normal_(mean=0,std=0.1))\n        self.lambda_k2 = nn.Parameter(torch.zeros(self.head_dim, dtype=torch.float32).normal_(mean=0,std=0.1))\n        self.rotary_emb = RotaryEmbedding(self.head_dim * 2)\n\n    def forward(self, x):\n        lambda_init = lambda_init_fn(self.depth)\n        Q = self.Q(x)\n        K = self.K(x)\n\n        seq_len = x.shape[1]\n        cos, sin = self.rotary_emb(seq_len, device=x.device)\n        Q, K = apply_rotary_pos_emb(Q, K, cos, sin)\n    \n        Q1, Q2 = Q.chunk(2, dim=-1)\n        K1, K2 = K.chunk(2, dim=-1)\n        V = self.V(x)\n        A1 = Q1 @ K1.transpose(-2, -1) * self.scale\n        A2 = Q2 @ K2.transpose(-2, -1) * self.scale\n        lambda_1 = torch.exp(torch.sum(self.lambda_q1 * self.lambda_k1, dim=-1).float()).type_as(Q1)\n        lambda_2 = torch.exp(torch.sum(self.lambda_q2 * self.lambda_k2, dim=-1).float()).type_as(Q2)\n        lambda_ = lambda_1 - lambda_2 + lambda_init\n        return (F.softmax(A1, dim=-1) - lambda_ * F.softmax(A2, dim=-1)) @ V\n\n\nclass MultiHeadDifferentialAttention(nn.Module):\n    def __init__(self, dim_model: int, head_nums: int, depth: int):\n        super().__init__()\n        self.heads = nn.ModuleList([DifferentialAttention(dim_model, head_nums, depth) for _ in range(head_nums)])\n        self.group_norm = RMSNorm(dim_model)\n        self.output = nn.Linear(2 * dim_model, dim_model, bias=False)\n        self.lambda_init = lambda_init_fn(depth)\n    \n    def forward(self, x):\n        o = torch.cat([self.group_norm(h(x)) for h in self.heads], dim=-1)\n        o = o * (1 - self.lambda_init)\n        return self.output(o)\n\n\nclass RotaryEmbedding(nn.Module):\n    def __init__(self, dim: int):\n        super().__init__()\n        self.dim = dim\n        inv_freq = 1.0 / (10000 ** (torch.arange(0, dim, 2).float() / dim))\n        self.register_buffer(\"inv_freq\", inv_freq)\n\n    def forward(self, max_seq_len: int, *, device: torch.device):\n        seq = torch.arange(max_seq_len, device=device)\n        freqs = torch.einsum(\"i,j->ij\", seq, self.inv_freq)\n        emb = torch.cat((freqs, freqs), dim=-1)\n        return emb.cos(), emb.sin()\n\ndef rotate_half(x):\n    x1, x2 = x.chunk(2, dim=-1)\n    return torch.cat((-x2, x1), dim=-1)\n\ndef apply_rotary_pos_emb(q, k, cos, sin):\n    q_embed = (q * cos) + (rotate_half(q) * sin)\n    k_embed = (k * cos) + (rotate_half(k) * sin)\n    return q_embed, k_embed\n\nclass RMSNorm(nn.Module):\n    def __init__(self, dim: int, eps: float = 1e-6):\n        super().__init__()\n        self.eps = eps\n        self.weight = nn.Parameter(torch.zeros(dim))\n\n    def forward(self, x):\n        return x * torch.rsqrt(x.pow(2).mean(dim=-1, keepdim=True) + self.eps)\n\nclass FeedForward(nn.Module):\n    def __init__(self, dim: int, hidden_dim: int):\n        super().__init__()\n        self.w1 = nn.Linear(dim, hidden_dim, bias=False)\n        self.w2 = nn.Linear(hidden_dim, dim, bias=False)\n        self.w3 = nn.Linear(dim, hidden_dim, bias=False)\n        self.silu = nn.SiLU()\n    \n    def forward(self, x):\n        return self.w2(self.silu(self.w1(x)) * self.w3(x))\n\nclass DifferentialTransformer(nn.Module):\n    def __init__(self, input_dim, dim, depth, heads=8, head_dim=16):\n        super().__init__()\n        self.input_dim = input_dim\n        self.dim = dim\n        \n        # Separate linear layers for each time step\n        self.input_projection = nn.Sequential(\n            nn.Linear(input_dim, dim),\n            nn.LayerNorm(dim),\n            nn.ReLU()\n        )\n        \n        self.layers = nn.ModuleList([\n            MultiHeadDifferentialAttention(dim, heads, depth_idx)\n            for depth_idx in range(depth) #Create a list of transformer layers\n        ])\n        \n        self.ln1 = RMSNorm(dim)\n        self.ln2 = RMSNorm(dim)\n        self.ffn = FeedForward(dim, (dim // 3) * 8)\n        self.output = nn.Linear(dim, 3)  # 3 classes: sell (-1), hold (0), buy (1)\n        \n    def forward(self, x):\n        # x shape: [batch, sequence, features]\n        batch_size, seq_len, _ = x.shape\n        \n        # Project each time step independently\n        x = x.view(-1, self.input_dim)  # [batch * sequence, features]\n        x = self.input_projection(x)     # [batch * sequence, dim]\n        x = x.view(batch_size, seq_len, self.dim)  # [batch, sequence, dim]\n        \n        # Apply transformer layers\n        for attn in self.layers:\n            y = attn(self.ln1(x)) + x\n            x = self.ffn(self.ln2(y)) + y\n        \n        # Take the output corresponding to the last time step\n        return self.output(x[:, -1, :])\n\nclass EnhancedTradingLoss(nn.Module):\n    def __init__(self, class_weights, alpha=0.3, beta=0.3, gamma=0.2):\n        super().__init__()\n        self.class_weights = class_weights\n        self.alpha = alpha  # Weight for directional accuracy\n        self.beta = beta    # Weight for confidence penalty\n        self.gamma = gamma  # Weight for temporal consistency\n        self.base_criterion = nn.CrossEntropyLoss(weight=class_weights, reduction='none')\n        \n    def forward(self, outputs, targets, prev_outputs=None):\n        # Basic classification loss with class weights\n        base_loss = self.base_criterion(outputs, targets)\n        \n        # Softmax probabilities\n        probs = F.softmax(outputs, dim=1)\n        \n        # Directional penalty (penalize opposite predictions more heavily)\n        directional_penalty = torch.zeros_like(base_loss)\n        pred_classes = torch.argmax(probs, dim=1)\n        opposite_mistakes = ((pred_classes == 0) & (targets == 2)) | \\\n                          ((pred_classes == 2) & (targets == 0))\n        directional_penalty[opposite_mistakes] = 1.0\n        \n        # Confidence penalty (penalize uncertain predictions)\n        max_probs, _ = torch.max(probs, dim=1)\n        confidence_penalty = 1.0 - max_probs\n        \n        # Temporal consistency penalty (if previous outputs available)\n        temporal_penalty = torch.zeros_like(base_loss)\n        if prev_outputs is not None:\n            prev_probs = F.softmax(prev_outputs, dim=1)\n            temporal_penalty = F.kl_div(\n                probs.log(), \n                prev_probs, \n                reduction='none'\n            ).sum(dim=1)\n        \n        # Combine all components\n        total_loss = (base_loss + \n                    self.alpha * directional_penalty + \n                    self.beta * confidence_penalty + \n                    self.gamma * temporal_penalty)\n        \n        return total_loss.mean()\n\n\n# Replace this function with your 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 10 minutes of the batch features being provided.\n# Global variables for model and scaler\ntransformer_model = None\nfeature_scaler = None\nsequence_length = 30\nlags_ = None\n\ndef initialize_model():\n    \"\"\"Initialize the transformer model (call this once before predictions start)\"\"\"\n    global transformer_model, feature_scaler\n    \n    # Count actual features from first batch\n    feature_count = len([col for col in test.columns \n                        if col not in ['row_id', 'time_id', 'responder_6']])\n    \n    # Initialize transformer model with correct input dimension\n    transformer_model = DifferentialTransformer(\n        input_dim=feature_count,  # Use actual feature count\n        dim=128,\n        depth=6,\n        heads=8,\n        head_dim=16\n    )\n    \n    transformer_model.eval()  # Set to evaluation mode\n    print(f\"Model initialized with {feature_count} input features\")\n\n# Global variables\ntransformer_model = None\nfeature_scaler = None\nsequence_length = 30\nlags_ = None\n\ndef initialize_model(input_dim):\n    \"\"\"Initialize the transformer model with the pre-trained weights\"\"\"\n    global transformer_model, feature_scaler\n    \n    # Initialize transformer model with the actual feature count\n    transformer_model = DifferentialTransformer(\n        input_dim=input_dim,\n        dim=128,\n        depth=6,\n        heads=8,\n        head_dim=16\n    )\n    \n    # Load pre-trained weights if available\n    try:\n        # Update this path to match your file location\n        checkpoint = torch.load('/kaggle/input/differentialtransformer/pytorch/default/1/transformer_model.pth')  # or wherever your .pth file is located\n        transformer_model.load_state_dict(checkpoint['model_state_dict'])\n        feature_scaler = StandardScaler()\n        feature_scaler.__dict__.update(checkpoint['scaler_state_dict'])\n        print(\"Loaded pre-trained model successfully\")\n    except Exception as e:\n        print(f\"Error loading pre-trained model: {e}\")\n        # Fall back to untrained model if loading fails\n        feature_scaler = StandardScaler()\n    \n    transformer_model.eval()\n    print(f\"Model initialized with input_dim={input_dim}\")\n\ndef predict(test: pl.DataFrame, lags: pl.DataFrame | None) -> pl.DataFrame:\n    try:\n        global lags_, transformer_model, feature_scaler\n        \n        # Store lags if provided\n        if lags is not None:\n            lags_ = lags\n        \n        # Prepare features\n        feature_cols = [col for col in test.columns \n                       if col not in ['row_id', 'time_id', 'responder_6']]\n        features = test.select(feature_cols)\n        feature_count = len(feature_cols)\n        \n        # Initialize model if needed\n        if transformer_model is None:\n            initialize_model(feature_count)\n        \n        # Convert to numpy and clean data efficiently\n        feature_array = features.to_numpy()\n        np.nan_to_num(feature_array, copy=False, nan=0.0)  # In-place operation\n        np.clip(feature_array, -1e9, 1e9, out=feature_array)  # In-place operation\n        \n        # Initialize and fit scaler if needed\n        if feature_scaler is None or not hasattr(feature_scaler, 'mean_'):\n            feature_scaler = StandardScaler()\n            feature_scaler.fit(feature_array)\n        \n        # Normalize features\n        normalized_features = feature_scaler.transform(feature_array)\n        \n        # Process in batches\n        batch_size = 1024  # Adjust based on your memory constraints\n        window_size = min(sequence_length, len(normalized_features))\n        all_predictions = np.zeros(len(normalized_features))\n        \n        for batch_start in range(0, len(normalized_features), batch_size):\n            batch_end = min(batch_start + batch_size, len(normalized_features))\n            batch_sequences = []\n            \n            for i in range(batch_start, batch_end):\n                start_idx = max(0, i - window_size + 1)\n                sequence = normalized_features[start_idx:i+1]\n                \n                if len(sequence) < window_size:\n                    padding = np.zeros((window_size - len(sequence), sequence.shape[1]))\n                    sequence = np.concatenate([padding, sequence])\n                \n                batch_sequences.append(sequence)\n            \n            # Convert batch to tensor\n            batch_tensor = torch.tensor(np.array(batch_sequences), dtype=torch.float32)\n            \n            # Generate predictions for batch\n            with torch.no_grad():\n                outputs = transformer_model(batch_tensor)\n                probs = F.softmax(outputs, dim=1)\n                pred_values = (probs[:, 2] - probs[:, 0]).numpy()\n                np.clip(pred_values, -1, 1, out=pred_values)\n                all_predictions[batch_start:batch_end] = pred_values\n        \n        # Create prediction DataFrame\n        predictions = pl.DataFrame({\n            'row_id': test['row_id'],\n            'responder_6': all_predictions\n        })\n        \n        return predictions\n        \n    except Exception as e:\n        print(f\"\\nError in prediction: {e}\")\n        return pl.DataFrame({\n            'row_id': test['row_id'],\n            'responder_6': [0.0] * len(test)\n        })\n        \n    except Exception as e:\n        print(f\"\\nError in prediction: {e}\")\n        print(\"Detailed error information:\")\n        import traceback\n        print(traceback.format_exc())\n        \n        # Print debug information if available\n        if 'feature_array' in locals():\n            print(f\"\\nFeature array shape: {feature_array.shape}\")\n        if 'sequence_tensor' in locals():\n            print(f\"Sequence tensor shape: {sequence_tensor.shape}\")\n        if 'outputs' in locals():\n            print(f\"Model output shape: {outputs.shape}\")\n        if 'all_predictions' in locals():\n            print(f\"Number of predictions generated: {len(all_predictions)}\")\n        \n        # Return safe default predictions\n        return pl.DataFrame({\n            'row_id': test['row_id'],\n            'responder_6': [0.0] * len(test)\n        })\n\n","metadata":{"_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","papermill":{"duration":0.018344,"end_time":"2024-10-10T11:58:33.59684","exception":false,"start_time":"2024-10-10T11:58:33.578496","status":"completed"},"tags":[],"trusted":true,"execution":{"iopub.status.busy":"2025-01-03T09:02:04.413611Z","iopub.execute_input":"2025-01-03T09:02:04.414068Z","iopub.status.idle":"2025-01-03T09:02:04.461823Z","shell.execute_reply.started":"2025-01-03T09:02:04.414035Z","shell.execute_reply":"2025-01-03T09:02:04.460706Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"When your notebook is run on the hidden test set, inference_server.serve must be called within 15 minutes of the notebook starting or the gateway will throw an error. If you need more than 15 minutes to load your model you can do so during the very first `predict` call, which does not have the usual 10 minute response deadline.","metadata":{"papermill":{"duration":0.002521,"end_time":"2024-10-10T11:58:33.6023","exception":false,"start_time":"2024-10-10T11:58:33.599779","status":"completed"},"tags":[]}},{"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-realtime-marketdata-forecasting/test.parquet',\n            '/kaggle/input/jane-street-realtime-marketdata-forecasting/lags.parquet',\n        )\n    )","metadata":{"_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","execution":{"iopub.status.busy":"2025-01-03T09:02:04.463039Z","iopub.execute_input":"2025-01-03T09:02:04.463370Z","iopub.status.idle":"2025-01-03T09:02:05.329905Z","shell.execute_reply.started":"2025-01-03T09:02:04.463339Z","shell.execute_reply":"2025-01-03T09:02:05.328777Z"},"papermill":{"duration":2.225871,"end_time":"2024-10-10T11:58:35.830964","exception":false,"start_time":"2024-10-10T11:58:33.605093","status":"completed"},"tags":[],"trusted":true},"outputs":[],"execution_count":null}]}