{"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":9910436,"sourceType":"datasetVersion","datasetId":5996723},{"sourceId":207221536,"sourceType":"kernelVersion"}],"dockerImageVersionId":30786,"isInternetEnabled":false,"language":"python","sourceType":"notebook","isGpuEnabled":false},"papermill":{"default_parameters":{},"duration":4.669361,"end_time":"2024-10-10T13:05:46.686069","environment_variables":{},"exception":null,"input_path":"__notebook__.ipynb","output_path":"__notebook__.ipynb","parameters":{},"start_time":"2024-10-10T13:05:42.016708","version":"2.6.0"}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"markdown","source":"# Jane Street Real-Time Market Data Forecasting with Tensorflow\n\n## Overview\nThis notebook is designed for a real-time market data forecasting task using TensorFlow. The primary goal is to predict a target variable (responder_6) based on 79 features. This notebook provides a comprehensive workflow for real-time market data forecasting using TensorFlow. It includes data loading, model training, validation, and a submission pipeline for inference. The model is a simple feedforward neural network, and the notebook uses custom metrics and callbacks to enhance training and evaluation.\n\n## Key Components\n1. **Configuration**: A configuration class CFG is defined to control whether the model is in training mode or not.\n\n2. **Load traning and validation data**\n\n    **Feature Columns**: The notebook defines 79 feature columns named feature_00 to feature_78.\n\n    **Target Column**: The target column is named responder_6.\n\n    **Training Data**: The training data is loaded from multiple Parquet files partitioned by partition_id. The data is concatenated into X_train (features) and y_train (target).\n\n    **Validation Data**: Similar to training data, validation data is loaded from Parquet files and concatenated into X_val and y_val.\n\n\n3. **Generators**: Two generator functions (train_generator and validation_generator) are created to yield batches of data for training and validation.\n4. **TensorFlow Dataset**\n: TensorFlow datasets (train_ds and valid_ds) are created from these generators.\n\n5. **Custom Metric**: A custom R² metric (R2Metric) is defined to evaluate the model's performance.\n\n6. **Model Architecture**: A simple feedforward neural network is defined with three hidden layers and a final output layer. The model is compiled with Mean Squared Error (MSE) loss and the custom R² metric.\n\n7. **Training**: The model is trained with callbacks for checkpointing, early stopping, and terminating on NaN loss.\n\n\n8. **Predict Function**: A predict function is defined to make predictions on test data. The function handles null values, makes predictions using the trained model, and formats the output as required.\n9. **Inference Server**: The notebook sets up an inference server using the JSInferenceServer class, which handles the prediction requests. And checks if it is running in a competition rerun environment and either serves the inference server or runs a local gateway for testing.","metadata":{}},{"cell_type":"code","source":"import os\nimport pandas as pd\nimport polars as pl\nimport kaggle_evaluation.jane_street_inference_server\nimport tensorflow as tf\nimport numpy as np\nimport pickle\nimport gc","metadata":{"execution":{"iopub.status.busy":"2024-11-14T15:52:09.019223Z","iopub.execute_input":"2024-11-14T15:52:09.019726Z","iopub.status.idle":"2024-11-14T15:52:09.026384Z","shell.execute_reply.started":"2024-11-14T15:52:09.019679Z","shell.execute_reply":"2024-11-14T15:52:09.024722Z"},"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## 1. Configuration","metadata":{}},{"cell_type":"code","source":"class CFG:\n    \n    is_training = False\n    \n    feature_columns =  [f\"feature_{i:02d}\" for i in range(79)]\n    \n    target_column = \"responder_6\"","metadata":{"execution":{"iopub.status.busy":"2024-11-14T15:46:53.995883Z","iopub.execute_input":"2024-11-14T15:46:53.996356Z","iopub.status.idle":"2024-11-14T15:46:54.004573Z","shell.execute_reply.started":"2024-11-14T15:46:53.996312Z","shell.execute_reply":"2024-11-14T15:46:54.003266Z"},"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## 2. Load training and validation data","metadata":{}},{"cell_type":"markdown","source":"### Load statics values","metadata":{}},{"cell_type":"code","source":"with open(\"/kaggle/input/jane-street-rmf-computing-feature-statistics/max.pkl\", \"rb\") as f:\n    max_values = pickle.load(f)\nwith open(\"/kaggle/input/jane-street-rmf-computing-feature-statistics/min.pkl\", \"rb\") as f:\n    min_values = pickle.load(f)\nwith open(\"/kaggle/input/jane-street-rmf-computing-feature-statistics/mean.pkl\", \"rb\") as f:\n    mean_values = pickle.load(f)\nwith open(\"/kaggle/input/jane-street-rmf-computing-feature-statistics/mean_pandas.pkl\", \"rb\") as f:\n    mean_pandas = pickle.load(f)\nmax_values = max_values.astype(np.float16)\nmin_values = min_values.astype(np.float16)\nmax_min_diff = max_values - min_values\nmean_values = mean_values.astype(np.float16)\nmean_pandas = mean_pandas.astype(np.float16)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-11-14T15:50:59.046390Z","iopub.execute_input":"2024-11-14T15:50:59.046911Z","iopub.status.idle":"2024-11-14T15:50:59.060846Z","shell.execute_reply.started":"2024-11-14T15:50:59.046867Z","shell.execute_reply":"2024-11-14T15:50:59.059719Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### Load training data","metadata":{}},{"cell_type":"code","source":"if CFG.is_training:\n    X_train = None\n    y_train = None\n    for i in range(8):\n        print(\"=\" * 30)\n        print(f\"Partition {i}\")\n        print(\"=\" * 30)\n        df = pd.read_parquet(f\"/kaggle/input/jane-street-real-time-market-data-forecasting/train.parquet/partition_id={i}/part-0.parquet\")\n        X_part = df[CFG.feature_columns].astype(np.float16)\n        X_part.fillna(mean_pandas, inplace=True)\n        X_part = (X_part - min_values) / max_min_diff\n        targets = df[CFG.target_column]\n        if X_train is None:\n            X_train = X_part\n            y_train = targets\n        else:\n            X_train = pd.concat([X_train, X_part])\n            y_train = pd.concat([y_train, targets])\n        del X_part\n        del targets\n        gc.collect()","metadata":{"execution":{"iopub.status.busy":"2024-11-14T15:52:22.364600Z","iopub.execute_input":"2024-11-14T15:52:22.365958Z","iopub.status.idle":"2024-11-14T15:55:21.920954Z","shell.execute_reply.started":"2024-11-14T15:52:22.365900Z","shell.execute_reply":"2024-11-14T15:55:21.919305Z"},"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### Load Validation data","metadata":{}},{"cell_type":"code","source":"X_val = None\ny_val = None\nval_weights = None\nfor i in range(8, 10):\n    df = pd.read_parquet(f\"/kaggle/input/jane-street-real-time-market-data-forecasting/train.parquet/partition_id={i}/part-0.parquet\")\n    X_part = df[CFG.feature_columns].astype(np.float16)\n    X_part.fillna(mean_pandas, inplace=True)\n    X_part = (X_part - min_values) / max_min_diff\n    targets = df[CFG.target_column]\n    weights = df[\"weight\"].astype(np.float16)\n    if X_val is None:\n        X_val = X_part\n        y_val = targets\n        val_weights = weights\n    else:\n        X_val = pd.concat([X_val, X_part])\n        y_val = pd.concat([y_val, targets])\n        val_weights = pd.concat([val_weights, weights])\n    del X_part\n    del weights\n    del targets\n    del df\n    gc.collect()","metadata":{"execution":{"iopub.status.busy":"2024-11-14T15:56:28.924255Z","iopub.execute_input":"2024-11-14T15:56:28.924704Z","iopub.status.idle":"2024-11-14T15:57:29.792486Z","shell.execute_reply.started":"2024-11-14T15:56:28.924664Z","shell.execute_reply":"2024-11-14T15:57:29.790995Z"},"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## 3. Create Training and Validation dataset","metadata":{}},{"cell_type":"code","source":"def get_generator(X, y, shuffle=True, batch_size=4096):\n    def generator():\n        # Create a shuffled index for random access\n        if shuffle:\n            indices = np.random.permutation(len(X))\n        else:\n            indices = np.arange(len(X))\n        num_batch = len(indices) // batch_size + (1 if len(indices) % batch_size > 0 else 0)\n        for i in range(num_batch):\n            start_index = i * batch_size\n            end_index = min((i + 1) * batch_size, len(X))\n            current_indices = indices[start_index: end_index]\n            features = X.iloc[current_indices].values\n            label = y.iloc[current_indices]\n            yield features, np.array(label).reshape(-1, 1)\n    return generator\n        \ndef create_dataset(X, y, shuffle=True, batch_size=4096):\n    output_signature = (\n        tf.TensorSpec(shape=(None, len(CFG.feature_columns),), dtype=tf.float16),  # Adjust shape to number of features\n        tf.TensorSpec(shape=(None, 1), dtype=tf.float16)\n    )\n    # Create a TensorFlow Dataset from the generator\n    ds = tf.data.Dataset.from_generator(\n        get_generator(X, y, shuffle=shuffle, batch_size=batch_size),\n        output_signature=output_signature\n    )\n    return ds","metadata":{"execution":{"iopub.status.busy":"2024-11-14T15:58:08.825034Z","iopub.execute_input":"2024-11-14T15:58:08.825622Z","iopub.status.idle":"2024-11-14T15:58:08.842342Z","shell.execute_reply.started":"2024-11-14T15:58:08.825573Z","shell.execute_reply":"2024-11-14T15:58:08.840509Z"},"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"if CFG.is_training:\n    train_ds = create_dataset(X_train, y_train, shuffle=True)","metadata":{"execution":{"iopub.status.busy":"2024-11-14T15:58:11.589433Z","iopub.execute_input":"2024-11-14T15:58:11.589872Z","iopub.status.idle":"2024-11-14T15:58:11.746741Z","shell.execute_reply.started":"2024-11-14T15:58:11.589831Z","shell.execute_reply":"2024-11-14T15:58:11.745539Z"},"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"valid_ds = create_dataset(X_val, y_val, shuffle=False)","metadata":{"execution":{"iopub.status.busy":"2024-11-14T15:58:14.167259Z","iopub.execute_input":"2024-11-14T15:58:14.167719Z","iopub.status.idle":"2024-11-14T15:58:14.209171Z","shell.execute_reply.started":"2024-11-14T15:58:14.167676Z","shell.execute_reply":"2024-11-14T15:58:14.207588Z"},"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"Check data format","metadata":{}},{"cell_type":"code","source":"for batch in valid_ds:\n    print(batch)\n    break","metadata":{"execution":{"iopub.status.busy":"2024-11-14T15:58:17.145139Z","iopub.execute_input":"2024-11-14T15:58:17.145577Z","iopub.status.idle":"2024-11-14T15:58:17.315235Z","shell.execute_reply.started":"2024-11-14T15:58:17.145537Z","shell.execute_reply":"2024-11-14T15:58:17.313890Z"},"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## 4. Modeling","metadata":{}},{"cell_type":"code","source":"class R2Metric(tf.keras.metrics.Metric):\n    def __init__(self, name='r2', **kwargs):\n        super(R2Metric, self).__init__(name=name, **kwargs)\n        self.squared_residuals_sum = self.add_weight(name='squared_residuals_sum', initializer='zeros')\n        self.total_sum_squares = self.add_weight(name='total_sum_squares', initializer='zeros')\n        self.count = self.add_weight(name='count', initializer='zeros')\n\n    def update_state(self, y_true, y_pred, sample_weight=None):\n        # Flatten tensors to ensure shape compatibility\n        y_true = tf.reshape(y_true, (-1,))\n        y_pred = tf.reshape(y_pred, (-1,))\n        y_true = tf.cast(y_true, dtype=tf.float32)\n        y_pred = tf.cast(y_pred, dtype=tf.float32)\n        # Residual Sum of Squares\n        residuals = y_true - y_pred\n        self.squared_residuals_sum.assign_add(tf.reduce_sum(tf.square(residuals)))\n        \n        # Total Sum of Squares\n        self.total_sum_squares.assign_add(tf.reduce_sum(tf.square(y_true)))\n        \n        # Increment count\n        self.count.assign_add(tf.cast(tf.size(y_true), tf.float32))\n\n    def result(self):\n        # Calculate R2: 1 - (SS_res / SS_tot)\n        return 1 - (self.squared_residuals_sum / (self.total_sum_squares + tf.keras.backend.epsilon()))\n\n    def reset_states(self):\n        # Reset all variables at the beginning of each epoch\n        self.squared_residuals_sum.assign(0.0)\n        self.total_sum_squares.assign(0.0)\n        self.count.assign(0.0)\n","metadata":{"execution":{"iopub.status.busy":"2024-11-14T16:03:34.827610Z","iopub.execute_input":"2024-11-14T16:03:34.828117Z","iopub.status.idle":"2024-11-14T16:03:34.843279Z","shell.execute_reply.started":"2024-11-14T16:03:34.828072Z","shell.execute_reply":"2024-11-14T16:03:34.842052Z"},"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"def get_model():\n    inputs = tf.keras.Input(shape=(79, ), dtype=tf.float16)\n    x = tf.keras.layers.Dense(512, activation=\"swish\")(inputs)\n    x = tf.keras.layers.Dense(256, activation=\"swish\")(x)\n    x = tf.keras.layers.Dense(128, activation=\"swish\")(x)\n    x = tf.keras.layers.Dense(32, activation=\"swish\", kernel_regularizer=\"l2\")(x)\n    outputs = tf.keras.layers.Dense(1)(x)\n    model = tf.keras.Model(inputs=inputs, outputs=outputs)\n    optimizer = tf.keras.optimizers.Adam(5e-5)\n    model.compile(loss=\"huber\", optimizer=optimizer, metrics=[R2Metric()])\n    return model","metadata":{"execution":{"iopub.status.busy":"2024-11-14T16:03:36.968809Z","iopub.execute_input":"2024-11-14T16:03:36.969291Z","iopub.status.idle":"2024-11-14T16:03:36.979515Z","shell.execute_reply.started":"2024-11-14T16:03:36.969245Z","shell.execute_reply":"2024-11-14T16:03:36.978120Z"},"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"model = get_model()\nmodel.summary()\ntf.keras.utils.plot_model(model, show_shapes=True, show_dtype=True)","metadata":{"execution":{"iopub.status.busy":"2024-11-14T16:03:38.968369Z","iopub.execute_input":"2024-11-14T16:03:38.968873Z","iopub.status.idle":"2024-11-14T16:03:39.338737Z","shell.execute_reply.started":"2024-11-14T16:03:38.968822Z","shell.execute_reply":"2024-11-14T16:03:39.337387Z"},"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## 5. Training or loading model","metadata":{}},{"cell_type":"code","source":"if CFG.is_training:\n    model = get_model()\n    callbacks = [\n        tf.keras.callbacks.ModelCheckpoint(\n            filepath='model.keras',  # Path where model is saved\n            monitor='val_r2',  # Monitor validation loss\n            save_best_only=True,  # Save only the best model\n            save_weights_only=False,  # Save the entire model, not just weights\n            mode='max',  # Minimize the monitored metric\n            verbose=1\n        ),\n        tf.keras.callbacks.TerminateOnNaN(),  # Stop training if NaN loss is encountered\n        tf.keras.callbacks.EarlyStopping(\n            monitor='val_loss',  # Monitor validation loss\n            patience=5,  # Number of epochs with no improvement after which training will stop\n            mode='min',  # Minimize the monitored metric\n            verbose=2\n        )\n    ]\n    # Train the model with callbacks\n    history = model.fit(\n        train_ds,\n        epochs=30,\n        validation_data=valid_ds,\n        callbacks=callbacks,\n        verbose=1\n    )\nelse:\n    model = tf.keras.models.load_model(\"/kaggle/input/jane-street-rmf-keras-model/model.keras\", custom_objects={\n        \"R2Metric\": R2Metric\n    })\nmodel.summary()","metadata":{"execution":{"iopub.status.busy":"2024-11-14T16:03:43.032076Z","iopub.execute_input":"2024-11-14T16:03:43.033299Z"},"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## 6. Model Evaluation","metadata":{}},{"cell_type":"code","source":"def calculate_r2(y_true, y_pred, weights):\n    \"\"\"\n    Calculate the sample weighted zero-mean R-squared score (R2).\n\n    Parameters:\n    - y_true (pd.Series or np.array): Ground truth values.\n    - y_pred (pd.Series or np.array): Predicted values.\n    - weights (pd.Series or np.array): Sample weights.\n\n    Returns:\n    - float: R2 score.\n    \"\"\"\n    numerator = np.sum(weights * (y_true - y_pred) ** 2)\n    denominator = np.sum(weights * (y_true ** 2))\n    r2_score = 1 - (numerator / denominator)\n    return r2_score","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-11-07T17:09:26.396823Z","iopub.execute_input":"2024-11-07T17:09:26.397703Z","iopub.status.idle":"2024-11-07T17:09:26.404725Z","shell.execute_reply.started":"2024-11-07T17:09:26.397647Z","shell.execute_reply":"2024-11-07T17:09:26.403368Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"y_pred = model.predict(valid_ds, verbose=0).reshape(-1)\nsample_weight_val = val_weights.values.reshape(-1)\nr2 = calculate_r2(y_val, y_pred, sample_weight_val)\nprint(f\"Validation R2:{r2:.4f}\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-11-07T17:07:48.121013Z","iopub.execute_input":"2024-11-07T17:07:48.122468Z","iopub.status.idle":"2024-11-07T17:07:48.127858Z","shell.execute_reply.started":"2024-11-07T17:07:48.122412Z","shell.execute_reply":"2024-11-07T17:07:48.126489Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## 7. Create submission pipeline","metadata":{}},{"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.002051,"end_time":"2024-10-10T13:05:45.83073","exception":false,"start_time":"2024-10-10T13:05:45.828679","status":"completed"},"tags":[]}},{"cell_type":"code","source":"lags_ : pl.DataFrame | None = None\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.\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_\n    if lags is not None:\n        lags_ = lags\n    # 1. Select the required feature columns and convert to numpy array for Keras\n    X_test = test.select(CFG.feature_columns).to_numpy()\n    X_test = np.where(np.isnan(X_test), mean_values, X_test)\n    X_test = (X_test - min_values) / max_min_diff\n    # 2. Make predictions using the Keras model\n    y_pred = model.predict(X_test, batch_size=4096)\n    \n    # 3. Prepare the DataFrame for output\n    predictions = test.select('row_id').with_columns(\n        pl.Series(\"responder_6\", y_pred.flatten())\n    )\n    # The predict function must return a DataFrame\n    assert isinstance(predictions, pl.DataFrame | pd.DataFrame)\n    # with columns 'row_id', 'responer_6'\n    assert predictions.columns == ['row_id', 'responder_6']\n    # and as many rows as the test data.\n    assert len(predictions) == len(test)\n\n    return predictions","metadata":{"_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","papermill":{"duration":0.015917,"end_time":"2024-10-10T13:05:45.848958","exception":false,"start_time":"2024-10-10T13:05:45.833041","status":"completed"},"tags":[],"execution":{"iopub.status.busy":"2024-11-02T11:24:56.448353Z","iopub.execute_input":"2024-11-02T11:24:56.448955Z","iopub.status.idle":"2024-11-02T11:24:56.458510Z","shell.execute_reply.started":"2024-11-02T11:24:56.448905Z","shell.execute_reply":"2024-11-02T11:24:56.457142Z"},"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"if CFG.is_training == False:\n\n    inference_server = kaggle_evaluation.jane_street_inference_server.JSInferenceServer(predict)\n\n    if os.getenv('KAGGLE_IS_COMPETITION_RERUN'):\n        inference_server.serve()\n    else:\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":{"_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","papermill":{"duration":0.308219,"end_time":"2024-10-10T13:05:46.163573","exception":false,"start_time":"2024-10-10T13:05:45.855354","status":"completed"},"tags":[],"execution":{"iopub.status.busy":"2024-11-02T11:24:59.078687Z","iopub.execute_input":"2024-11-02T11:24:59.079183Z","iopub.status.idle":"2024-11-02T11:24:59.211405Z","shell.execute_reply.started":"2024-11-02T11:24:59.079113Z","shell.execute_reply":"2024-11-02T11:24:59.209899Z"},"trusted":true},"outputs":[],"execution_count":null}]}