{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"name":"python","version":"3.10.14","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":30787,"isInternetEnabled":false,"language":"python","sourceType":"notebook","isGpuEnabled":true}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"markdown","source":"# Mainly for model training\n\nDepending on the size of your training set, you will need an [inference notebook](https://www.kaggle.com/code/regisvargas/inference-jane-street-a-beginner-s-notebook).","metadata":{}},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"import os\nimport gc\nimport numpy as np\nimport pandas as pd\nimport tensorflow as tf\nfrom tensorflow.keras import layers, models\nfrom tensorflow.keras.callbacks import ModelCheckpoint\nfrom sklearn.preprocessing import StandardScaler\nfrom sklearn.impute import SimpleImputer\nfrom tensorflow.keras.callbacks import EarlyStopping\n\n# Set test flag to switch between training and testing phases\nTEST = False  # Change to True for testing, False for training\n\npd.set_option('display.max_rows', None)  # No limit on rows\npd.set_option('display.max_columns', None)  # No limit on columns\n\nscaler = StandardScaler()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-11-29T10:09:01.557041Z","iopub.execute_input":"2024-11-29T10:09:01.557416Z","iopub.status.idle":"2024-11-29T10:09:13.603963Z","shell.execute_reply.started":"2024-11-29T10:09:01.557382Z","shell.execute_reply":"2024-11-29T10:09:13.603302Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"def prepare_data():\n    # Initialize list for chunks\n    samples = []\n\n    # Load data from partition (you can loop through other partitions if needed)\n    for i in [8, 9]:  # Add more partition ids if needed\n        file_path = f\"/kaggle/input/jane-street-real-time-market-data-forecasting/train.parquet/partition_id={i}/part-0.parquet\"\n        chunk = pd.read_parquet(file_path)\n        samples.append(chunk)\n\n    # Concatenate all chunks into a single DataFrame\n    sample_df = pd.concat(samples, ignore_index=False)\n    del samples\n    gc.collect()\n\n    # Fill missing values (forward fill, then zero if needed)\n    sample_df = sample_df.ffill()  # Forward fill missing values\n    sample_df = sample_df.fillna(0)\n\n    # Feature engineering\n    features = sample_df.filter(regex='^feature_')\n    responders = sample_df.filter(regex='^responder_')\n    weights = sample_df['weight']\n    \n    # Drop 'responder_6' and 'responder_3' columns, and add an engineered feature\n    responders_selected = responders.drop(columns=['responder_6', 'responder_3'])\n    responders_selected['arctan_responder_3'] = np.arctan(responders['responder_3'].values)\n\n    # Use only 'features' for X_df (no concatenation with responders)\n    X_df = features  # Directly use 'features' as X_df\n    y_df = responders[['responder_6']]  # Target variable\n\n    # Add 'time_id' and 'symbol_id' columns to the feature set using .loc\n    X_df.loc[:, 'time_id'] = sample_df['time_id']\n    X_df.loc[:, 'symbol_id'] = sample_df['symbol_id']\n\n    # Convert to numpy arrays for model input\n    X = X_df.values  # Input features\n    y = y_df.values  # Target variable\n\n    del X_df, y_df\n    gc.collect()\n\n    return X, y, weights\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-11-29T10:09:13.605422Z","iopub.execute_input":"2024-11-29T10:09:13.605883Z","iopub.status.idle":"2024-11-29T10:09:13.612642Z","shell.execute_reply.started":"2024-11-29T10:09:13.605855Z","shell.execute_reply":"2024-11-29T10:09:13.611808Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"# Prepare data","metadata":{}},{"cell_type":"code","source":"# Custom ordered time-series splitter\ndef ordered_time_series_split(X, n_splits=3):\n    n_samples = len(X)\n    if n_splits < 2:\n        raise ValueError(\"Number of splits must be at least 2.\")\n    fold_size = n_samples // n_splits\n    for i in range(n_splits):\n        test_start = i * fold_size\n        test_end = (i + 1) * fold_size if i < n_splits - 1 else n_samples\n        test_indices = np.arange(test_start, test_end)\n        train_indices = np.concatenate([np.arange(0, test_start), np.arange(test_end, n_samples)])\n        yield train_indices, test_indices","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-11-29T10:09:13.613647Z","iopub.execute_input":"2024-11-29T10:09:13.613896Z","iopub.status.idle":"2024-11-29T10:09:13.636399Z","shell.execute_reply.started":"2024-11-29T10:09:13.613872Z","shell.execute_reply":"2024-11-29T10:09:13.635437Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# Model creation function\ndef create_ae_mlp(num_columns, num_labels, hidden_units, dropout_rates, ls=0, lr=1e-3):\n    inp = tf.keras.layers.Input(shape=(num_columns,))\n    x0 = tf.keras.layers.BatchNormalization()(inp)\n\n    # Encoder\n    encoder = tf.keras.layers.GaussianNoise(dropout_rates[0])(x0)\n    encoder = tf.keras.layers.Dense(hidden_units[0])(encoder)\n    encoder = tf.keras.layers.BatchNormalization()(encoder)\n    encoder = tf.keras.layers.Activation('swish')(encoder)\n\n    # Decoder\n    decoder = tf.keras.layers.Dropout(dropout_rates[1])(encoder)\n    decoder = tf.keras.layers.Dense(num_columns, name='decoder')(decoder)\n\n    # Autoencoder Action Branch\n    x_ae = tf.keras.layers.Dense(hidden_units[1])(decoder)\n    x_ae = tf.keras.layers.BatchNormalization()(x_ae)\n    x_ae = tf.keras.layers.Activation('swish')(x_ae)\n    x_ae = tf.keras.layers.Dropout(dropout_rates[2])(x_ae)\n\n    out_ae = tf.keras.layers.Dense(num_labels, activation='linear', name='ae_res_6')(x_ae)\n\n    # Main Output Branch\n    x = tf.keras.layers.Concatenate()([x0, encoder])\n    x = tf.keras.layers.BatchNormalization()(x)\n    x = tf.keras.layers.Dropout(dropout_rates[3])(x)\n\n    for i in range(2, len(hidden_units)):\n        x = tf.keras.layers.Dense(hidden_units[i])(x)\n        x = tf.keras.layers.BatchNormalization()(x)\n        x = tf.keras.layers.Activation('swish')(x)\n        x = tf.keras.layers.Dropout(dropout_rates[i + 2])(x)\n\n    out = tf.keras.layers.Dense(num_labels, activation='linear', name='res_6')(x)\n\n    # Define and compile the model\n    model = tf.keras.models.Model(inputs=inp, outputs=[decoder, out_ae, out])\n    model.compile(\n        optimizer=tf.keras.optimizers.Adam(learning_rate=lr),\n        loss={\n            'decoder': tf.keras.losses.MeanSquaredError(),\n            'ae_res_6': tf.keras.losses.MeanSquaredError(),\n            'res_6': tf.keras.losses.MeanSquaredError(),\n        },\n        metrics={\n            'decoder': tf.keras.metrics.MeanAbsoluteError(name='MAE'),\n            'ae_res_6': tf.keras.metrics.RootMeanSquaredError(name='RMSE'),\n            'res_6': tf.keras.metrics.RootMeanSquaredError(name='RMSE'),\n        },\n    )\n\n    return model","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-11-29T10:09:13.638385Z","iopub.execute_input":"2024-11-29T10:09:13.638739Z","iopub.status.idle":"2024-11-29T10:09:13.650583Z","shell.execute_reply.started":"2024-11-29T10:09:13.638711Z","shell.execute_reply":"2024-11-29T10:09:13.649769Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"def get_model_checkpoint_callback(fold_num):\n    return tf.keras.callbacks.ModelCheckpoint(\n        f\"best_model_fold_{fold_num}.h5.keras\",  # Path where the best model will be saved\n        monitor='val_loss',  # You can also monitor 'val_accuracy' or another metric\n        mode='min',  # Save the model with the minimum validation loss\n        save_best_only=True,  # Only save the model when the validation loss improves\n        verbose=1  # Print messages when saving the model\n    )","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-11-29T10:09:13.651808Z","iopub.execute_input":"2024-11-29T10:09:13.652425Z","iopub.status.idle":"2024-11-29T10:09:13.661781Z","shell.execute_reply.started":"2024-11-29T10:09:13.652385Z","shell.execute_reply":"2024-11-29T10:09:13.660919Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"params = {\n    'num_columns': None,  # To be set dynamically based on the data\n    'num_labels': None,   # To be set dynamically based on the data\n    'hidden_units': [96, 96, 896, 448, 448, 256], \n    'dropout_rates': [0.03527936123679956, 0.038424974585075086, 0.42409238408801436, 0.10431484318345882, 0.49230389137187497, 0.32024444956111164, 0.2716856145683449, 0.4379233941604448], \n}","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-11-29T10:09:13.662911Z","iopub.execute_input":"2024-11-29T10:09:13.663282Z","iopub.status.idle":"2024-11-29T10:09:13.669573Z","shell.execute_reply.started":"2024-11-29T10:09:13.663223Z","shell.execute_reply":"2024-11-29T10:09:13.668770Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"gc.collect()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-11-29T10:09:13.670542Z","iopub.execute_input":"2024-11-29T10:09:13.670774Z","iopub.status.idle":"2024-11-29T10:09:13.860700Z","shell.execute_reply.started":"2024-11-29T10:09:13.670751Z","shell.execute_reply":"2024-11-29T10:09:13.859816Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"def evaluate_model_on_fold(fold_num, X, y, weights, model_paths):\n    train_idx, test_idx = list(ordered_time_series_split(X, n_splits=5))[fold_num]\n    X_train, X_test = X[train_idx], X[test_idx]\n    y_train, y_test = y[train_idx], y[test_idx]\n    sample_weights_test = weights.iloc[test_idx].values\n    \n    # Load the best model for this fold\n    model = tf.keras.models.load_model(model_paths[fold_num])\n    \n    # Predict the target variable (responder_6)\n    y_pred = model.predict(X_test)[2]\n    \n    # Calculate the weighted R² score\n    r2 = calculate_weighted_r2(y_test, y_pred, sample_weights_test)\n    \n    print(f\"Results for Fold {fold_num}: R² = {r2}\")\n    return r2","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-11-29T10:09:13.861850Z","iopub.execute_input":"2024-11-29T10:09:13.862135Z","iopub.status.idle":"2024-11-29T10:09:13.868611Z","shell.execute_reply.started":"2024-11-29T10:09:13.862092Z","shell.execute_reply":"2024-11-29T10:09:13.867669Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"def train_and_save_model(X, y, weights, params):\n    for fold_num, (train_idx, test_idx) in enumerate(ordered_time_series_split(X, n_splits=3)):\n        if fold_num !=9:\n            X_train, X_test = X[train_idx], X[test_idx]\n            y_train, y_test = y[train_idx], y[test_idx]\n            sample_weights_train = weights.iloc[train_idx].values\n            sample_weights_test = weights.iloc[test_idx].values\n            \n            # Model Parameters\n            num_columns = params['num_columns']\n            num_labels = params['num_labels']\n            hidden_units = params['hidden_units']\n            dropout_rates = params['dropout_rates']\n            \n            # Create and train the model\n            model = create_ae_mlp(num_columns, num_labels, hidden_units, dropout_rates)\n            checkpoint_callback = get_model_checkpoint_callback(fold_num)\n            \n            early_stopping_callback = EarlyStopping(\n                monitor='val_loss',  # You can change this to 'val_accuracy' or other metrics if needed\n                patience=7,  # Number of epochs with no improvement after which training will stop\n                min_delta=0.001,  # Minimum change to qualify as an improvement\n                restore_best_weights=True,  # Restore the best model weights after early stopping\n                verbose=1\n            )\n            \n            # Train the model\n            history = model.fit(\n                X_train,\n                {\"decoder\": X_train, \"ae_res_6\": y_train, \"res_6\": y_train},\n                validation_data=(X_test, {\"decoder\": X_test, \"ae_res_6\": y_test, \"res_6\": y_test}),\n                epochs=73,\n                batch_size=4096,\n                sample_weight=sample_weights_train,\n                callbacks=[checkpoint_callback,early_stopping_callback],\n                verbose=1,\n            )\n            \n            # Optionally clear session and delete model to free memory\n            del model  # Delete the model object from memory\n            tf.keras.backend.clear_session()  ","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-11-29T10:09:13.869938Z","iopub.execute_input":"2024-11-29T10:09:13.870202Z","iopub.status.idle":"2024-11-29T10:09:13.879720Z","shell.execute_reply.started":"2024-11-29T10:09:13.870176Z","shell.execute_reply":"2024-11-29T10:09:13.878980Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"if not TEST:\n    # Train the models\n    X, y, weights = prepare_data()\n    params['num_columns'] = X.shape[1]\n    params['num_labels'] = y.shape[1]\n    train_and_save_model(X, y, weights, params)\n\nelse:\n    # Testing code (evaluate previously saved models)\n    model_paths = [f\"best_model_fold_{fold_num}.h5.keras\" for fold_num in range(5)]\n    r2_scores = []\n    for fold_num in range(5):\n        r2 = evaluate_model_on_fold(fold_num, X, y, weights, model_paths)\n        r2_scores.append(r2)\n    \n    # Compute average R² score across all folds\n    average_r2 = np.mean(r2_scores)\n    print(f\"Average R² score across all folds: {average_r2}\")\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-11-29T10:09:13.881915Z","iopub.execute_input":"2024-11-29T10:09:13.882403Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# def preprocess_for_prediction(test: pd.DataFrame, lags: pd.DataFrame | None = None) -> np.ndarray:\n#     \"\"\"Preprocess test data to match the feature engineering done during training.\"\"\"\n#     global lags_\n\n#     # Convert test to pandas DataFrame if it is not already a pandas DataFrame\n#     if not isinstance(test, pd.DataFrame):\n#         test = test.to_pandas()\n\n#     # Update lags if provided (if applicable to your model)\n#     if lags is not None:\n#         lags_ = lags\n\n#     # Fill missing values (same as during training)\n#     test.fillna(method='ffill', inplace=True)\n#     test.fillna(0, inplace=True)\n\n#     # Feature engineering: Extract feature columns\n#     feat_cols = [col for col in test.columns if col.startswith('feature_')]\n#     features = test[feat_cols].values\n\n#     # Process the responders (similar to the training phase)\n#     responders = test.filter(regex='^responder_')\n    \n#     # responders_selected = responders.drop(columns=['responder_6', 'responder_3'])\n#     # responders_selected['arctan_responder_3'] = np.arctan(responders['responder_3'].values)\n\n#     # Combine features and responders\n#     # X_df = np.hstack([features])\n#     del responders  # Clear intermediate variables\n#     gc.collect()\n\n#     # Additional columns to features: time_id and symbol_id\n#     time_id = test['time_id'].values.reshape(-1, 1)  # Ensure it is (n_samples, 1)\n#     symbol_id = test['symbol_id'].values.reshape(-1, 1)  # Ensure it is (n_samples, 1)\n\n\n#     X = np.hstack([features, time_id, symbol_id])\n#     del  time_id, symbol_id  # Cleanup\n#     gc.collect()\n#     print(\"Final X shape:\", X.shape)\n\n#     return X\n","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# def predict(test: pd.DataFrame, lags: pd.DataFrame | None = None) -> pd.DataFrame:\n#     \"\"\"Generate predictions using the pre-trained model.\"\"\"\n#     # Preprocess the test data to prepare features (matching the training preprocessing)\n#     X = preprocess_for_prediction(test, lags)\n\n#     # Load the pre-trained model\n#     model_path = \"/kaggle/input/janestreet-demo/keras/default/1/best_model_fold_0.h5.keras\"\n#     model = tf.keras.models.load_model(model_path)\n\n#     # Generate predictions\n#     pred = model.predict(X)[1].reshape(-1)  # Assuming responder_6 is the first output\n   \n#     # print(pred)\n#     # Prepare the predictions DataFrame\n#     predictions = pd.DataFrame({\n#         'row_id': test['row_id'],\n#         'responder_6': pred\n#     })\n\n#     # Validate output\n#     assert predictions.columns.tolist() == ['row_id', 'responder_6'], \"Output DataFrame structure mismatch\"\n#     assert len(predictions) == len(test), \"Mismatch in number of rows between input and output\"\n\n#     return predictions","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# 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":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null}]}