{"metadata":{"kaggle":{"accelerator":"none","dataSources":[{"sourceId":84493,"databundleVersionId":9871156,"sourceType":"competition"},{"sourceId":9795629,"sourceType":"datasetVersion","datasetId":5996723}],"dockerImageVersionId":30786,"isInternetEnabled":false,"language":"python","sourceType":"notebook","isGpuEnabled":false},"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"},"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.\nModeling\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.\nExecution\n\nThe notebook 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","metadata":{"execution":{"iopub.status.busy":"2024-11-03T23:55:00.452992Z","iopub.execute_input":"2024-11-03T23:55:00.453997Z","iopub.status.idle":"2024-11-03T23:55:00.461014Z","shell.execute_reply.started":"2024-11-03T23:55:00.453931Z","shell.execute_reply":"2024-11-03T23:55:00.459570Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"## 1. Configuration","metadata":{}},{"cell_type":"code","source":"class CFG:\n    is_training = False","metadata":{"execution":{"iopub.status.busy":"2024-11-03T23:55:00.463386Z","iopub.execute_input":"2024-11-03T23:55:00.465142Z","iopub.status.idle":"2024-11-03T23:55:00.498357Z","shell.execute_reply.started":"2024-11-03T23:55:00.464983Z","shell.execute_reply":"2024-11-03T23:55:00.496729Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"","metadata":{},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"## 2. Load traning and validation data","metadata":{}},{"cell_type":"code","source":"feature_columns = [f\"feature_{i:02d}\" for i in range(79)]\ntarget_column = \"responder_6\"","metadata":{"execution":{"iopub.status.busy":"2024-11-03T23:55:00.501578Z","iopub.execute_input":"2024-11-03T23:55:00.502571Z","iopub.status.idle":"2024-11-03T23:55:00.518368Z","shell.execute_reply.started":"2024-11-03T23:55:00.502431Z","shell.execute_reply":"2024-11-03T23:55:00.516173Z"},"trusted":true},"execution_count":null,"outputs":[]},{"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    n_samples = 0\n    feature_sum = None\n    feature_sq_sum = 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[feature_columns]\n        # Initialize sums if first partition\n        if feature_sum is None:\n            feature_sum = X_part.sum()\n            feature_sq_sum = (X_part ** 2).sum()\n            n_samples = X_part.shape[0]\n            X_train = X_part\n            y_train = df[target_column]\n        else:\n            feature_sum += X_part.sum()\n            feature_sq_sum += (X_part ** 2).sum()\n            n_samples += X_part.shape[0]\n            X_train = pd.concat([X_train, X_part])\n            y_train = pd.concat([y_train, df[target_column]])\n    # Calculate mean and standard deviation\n    mean_values = feature_sum / n_samples\n    variance = (feature_sq_sum / n_samples) - (mean_values ** 2)\n    std_values = np.sqrt(variance)\n    print(\"Mean values:\\n\", mean_values)\n    print(\"\\nStandard deviation values:\\n\", std_values)\n    mean_values = mean_values.astype(np.float32)\n    std_values = std_values.astype(np.float32)\n    for j, val in enumerate(std_values):\n        if pd.isna(val):\n            std_values[j] = X_train[feature_columns[j]].std()\n    with open(\"mean.pkl\", \"wb\") as f:\n        pickle.dump(mean_values, f)\n    with open(\"std.pkl\", \"wb\") as f:\n        pickle.dump(std_values, f)\nelse:\n    with open(\"/kaggle/input/jane-street-rmf-keras-model/mean.pkl\", \"rb\") as f:\n        mean_values = pickle.load(f) \n    with open(\"/kaggle/input/jane-street-rmf-keras-model/std.pkl\", \"rb\") as f:\n        std_values = pickle.load(f)","metadata":{"execution":{"iopub.status.busy":"2024-11-03T23:55:00.519990Z","iopub.execute_input":"2024-11-03T23:55:00.520477Z","iopub.status.idle":"2024-11-03T23:55:00.549093Z","shell.execute_reply.started":"2024-11-03T23:55:00.520422Z","shell.execute_reply":"2024-11-03T23:55:00.547594Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"### Load Validation data","metadata":{}},{"cell_type":"code","source":"X_val = None\ny_val = 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    if X_val is None:\n        X_val = df[feature_columns]\n        y_val = df[target_column]\n    else:\n        X_val = pd.concat([X_val, df[feature_columns]])\n        y_val = pd.concat([y_val, df[target_column]])\nprint(X_val.shape, y_val.shape)","metadata":{"execution":{"iopub.status.busy":"2024-11-03T23:55:00.554691Z","iopub.execute_input":"2024-11-03T23:55:00.555142Z","iopub.status.idle":"2024-11-03T23:55:25.490408Z","shell.execute_reply.started":"2024-11-03T23:55:00.555097Z","shell.execute_reply":"2024-11-03T23:55:25.489107Z"},"trusted":true},"execution_count":null,"outputs":[]},{"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) - 1)\n            current_indices = indices[start_index: end_index]\n            features = X.iloc[start_index: end_index].replace(np.NAN, -1).values\n            label = y.iloc[start_index: end_index]\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(feature_columns),), dtype=tf.float32),  # Adjust shape to number of features\n        tf.TensorSpec(shape=(None, 1), dtype=tf.float32)\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-03T23:55:25.491894Z","iopub.execute_input":"2024-11-03T23:55:25.492263Z","iopub.status.idle":"2024-11-03T23:55:25.503063Z","shell.execute_reply.started":"2024-11-03T23:55:25.492226Z","shell.execute_reply":"2024-11-03T23:55:25.501660Z"},"trusted":true},"execution_count":null,"outputs":[]},{"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-03T23:55:25.504566Z","iopub.execute_input":"2024-11-03T23:55:25.504954Z","iopub.status.idle":"2024-11-03T23:55:25.532106Z","shell.execute_reply.started":"2024-11-03T23:55:25.504911Z","shell.execute_reply":"2024-11-03T23:55:25.530878Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"valid_ds = create_dataset(X_val, y_val, shuffle=False)","metadata":{"execution":{"iopub.status.busy":"2024-11-03T23:55:25.535626Z","iopub.execute_input":"2024-11-03T23:55:25.536776Z","iopub.status.idle":"2024-11-03T23:55:25.632457Z","shell.execute_reply.started":"2024-11-03T23:55:25.536717Z","shell.execute_reply":"2024-11-03T23:55:25.631139Z"},"trusted":true},"execution_count":null,"outputs":[]},{"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-03T23:55:25.634001Z","iopub.execute_input":"2024-11-03T23:55:25.634415Z","iopub.status.idle":"2024-11-03T23:55:25.755254Z","shell.execute_reply.started":"2024-11-03T23:55:25.634374Z","shell.execute_reply":"2024-11-03T23:55:25.753926Z"},"trusted":true},"execution_count":null,"outputs":[]},{"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        \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        mean_true = tf.reduce_mean(y_true)\n        self.total_sum_squares.assign_add(tf.reduce_sum(tf.square(y_true - mean_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-03T23:55:25.756444Z","iopub.execute_input":"2024-11-03T23:55:25.756828Z","iopub.status.idle":"2024-11-03T23:55:25.846290Z","shell.execute_reply.started":"2024-11-03T23:55:25.756788Z","shell.execute_reply":"2024-11-03T23:55:25.845031Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"def get_model():\n    inputs = tf.keras.Input(shape=(79, ), dtype=tf.float32)\n    #mean_val = tf.constant(mean_values)\n    #std_val = tf.constant(std_values)\n    #x = (inputs - mean_val) / std_val\n    x = tf.keras.layers.Dense(128, activation=\"swish\")(inputs)\n    x = tf.keras.layers.Dense(64, activation=\"swish\")(x)\n    x = tf.keras.layers.Dense(32, activation=\"swish\")(x)\n    outputs = 5 * tf.keras.layers.Dense(1, activation=\"tanh\")(x)\n    model = tf.keras.Model(inputs=inputs, outputs=outputs)\n    optimizer = tf.keras.optimizers.Adam(3e-4)\n    model.compile(loss=\"huber\", optimizer=optimizer, metrics=[R2Metric()])\n    return model","metadata":{"execution":{"iopub.status.busy":"2024-11-03T23:55:25.847698Z","iopub.execute_input":"2024-11-03T23:55:25.848097Z","iopub.status.idle":"2024-11-03T23:55:25.856287Z","shell.execute_reply.started":"2024-11-03T23:55:25.848028Z","shell.execute_reply":"2024-11-03T23:55:25.854797Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"model = get_model()\nmodel.summary()","metadata":{"execution":{"iopub.status.busy":"2024-11-03T23:55:25.858021Z","iopub.execute_input":"2024-11-03T23:55:25.858601Z","iopub.status.idle":"2024-11-03T23:55:25.980795Z","shell.execute_reply.started":"2024-11-03T23:55:25.858533Z","shell.execute_reply":"2024-11-03T23:55:25.979633Z"},"trusted":true},"execution_count":null,"outputs":[]},{"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=1\n        )\n    ]\n    # Train the model with callbacks\n    history = model.fit(\n        train_ds,\n        epochs=10,\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-03T23:55:25.982380Z","iopub.execute_input":"2024-11-03T23:55:25.983367Z","iopub.status.idle":"2024-11-03T23:55:26.166789Z","shell.execute_reply.started":"2024-11-03T23:55:25.983311Z","shell.execute_reply":"2024-11-03T23:55:26.165803Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"## 6. Model Evaluation","metadata":{}},{"cell_type":"code","source":"val_loss, val_r2 = model.evaluate(valid_ds)\nprint(f\"Validation R2:{val_r2:.4f}\")","metadata":{"execution":{"iopub.status.busy":"2024-11-03T23:55:26.168900Z","iopub.execute_input":"2024-11-03T23:55:26.169302Z","iopub.status.idle":"2024-11-03T23:55:51.861857Z","shell.execute_reply.started":"2024-11-03T23:55:26.169260Z","shell.execute_reply":"2024-11-03T23:55:51.860550Z"},"trusted":true},"execution_count":null,"outputs":[]},{"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(feature_columns).fill_null(-1).to_numpy()\n    \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\n    print(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-03T23:55:51.863562Z","iopub.execute_input":"2024-11-03T23:55:51.864030Z","iopub.status.idle":"2024-11-03T23:55:51.873565Z","shell.execute_reply.started":"2024-11-03T23:55:51.863978Z","shell.execute_reply":"2024-11-03T23:55:51.872267Z"},"trusted":true},"execution_count":null,"outputs":[]},{"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.00196,"end_time":"2024-10-10T13:05:45.853279","exception":false,"start_time":"2024-10-10T13:05:45.851319","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-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-03T23:55:51.876606Z","iopub.execute_input":"2024-11-03T23:55:51.877073Z","iopub.status.idle":"2024-11-03T23:55:52.448395Z","shell.execute_reply.started":"2024-11-03T23:55:51.876995Z","shell.execute_reply":"2024-11-03T23:55:52.447019Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"","metadata":{},"execution_count":null,"outputs":[]}]}