{"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":"nvidiaTeslaT4","dataSources":[{"sourceId":84493,"databundleVersionId":9871156,"sourceType":"competition"}],"dockerImageVersionId":30822,"isInternetEnabled":false,"language":"python","sourceType":"notebook","isGpuEnabled":true}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"import numpy as np\nimport pandas as pd\nimport tensorflow as tf\nfrom tensorflow.keras.models import Sequential\nfrom tensorflow.keras.layers import Dense, Dropout, BatchNormalization\nfrom tensorflow.keras.callbacks import EarlyStopping\nfrom sklearn.model_selection import train_test_split\nfrom sklearn.preprocessing import StandardScaler\nfrom datasets import Dataset\nfrom datasets import load_dataset  # For handling large datasets\nimport os\nimport kaggle_evaluation.jane_street_inference_server\nimport polars as pl\nfrom pyarrow.dataset import dataset\nfrom sklearn.metrics import r2_score\n\nimport numpy as np\nimport pandas as pd\nimport tensorflow as tf\nfrom tensorflow.keras.models import Sequential\nfrom tensorflow.keras.layers import Dense, Dropout, BatchNormalization\nfrom tensorflow.keras.callbacks import EarlyStopping, ReduceLROnPlateau\nfrom sklearn.model_selection import train_test_split\nfrom sklearn.preprocessing import StandardScaler\nfrom sklearn.metrics import r2_score\nimport os\nimport polars as pl\nfrom pyarrow.dataset import dataset\n\n# Define paths\ntrain_folder = '/kaggle/input/jane-street-real-time-market-data-forecasting/train.parquet'\npartitions = [f for f in os.listdir(train_folder) if os.path.isdir(os.path.join(train_folder, f))]\nlags_path = '/kaggle/input/jane-street-real-time-market-data-forecasting/lags.parquet'\n\n# Load lagged features in chunks\nprint(\"Loading lagged features...\")\nlags_ds = dataset(lags_path, format=\"parquet\")\nlags = lags_ds.to_table().to_pandas()\nlags['date_id'] = lags['date_id'].astype('int32')\nprint(\"Lagged features loaded.\")\n\n# Define the model\ndef build_model(input_dim):\n    model = Sequential([\n        Dense(128, activation='relu', input_dim=input_dim),\n        BatchNormalization(),\n        Dropout(0.3),\n        Dense(64, activation='relu'),\n        BatchNormalization(),\n        Dropout(0.3),\n        Dense(1, activation='linear')\n    ])\n    model.compile(optimizer='adam', loss='mean_squared_error', metrics=['mae'])\n    return model\n\n# Callbacks\nearly_stopping = EarlyStopping(monitor='val_loss', patience=5, restore_best_weights=True)\nreduce_lr = ReduceLROnPlateau(monitor='val_loss', factor=0.5, patience=3, verbose=1)\n\nscaler = StandardScaler()\n\n# Process each partition in chunks\nr2_scores = []\nfor i, partition in enumerate(partitions):\n    partition_path = os.path.join(train_folder, partition, 'part-0.parquet')\n    print(f\"Processing {partition_path}...\")\n\n    # Load partition data in chunks\n    temp_data = pd.read_parquet(partition_path)\n    temp_data = temp_data[temp_data['weight'] > 0]  # Filter based on weight\n    temp_data.fillna(0, inplace=True)\n\n    features = [col for col in temp_data.columns if 'feature' in col]\n\n    # Merge lag features\n    print(\"Merging lag features...\")\n    temp_data = temp_data.merge(\n        lags,\n        on=['date_id', 'time_id', 'symbol_id'],\n        how='left'\n    ).fillna(0)\n\n    # Scale features\n    print(\"Scaling features...\")\n    all_features = features + [col for col in lags.columns if col.startswith('responder_') and 'lag' in col]\n    temp_data[all_features] = scaler.fit_transform(temp_data[all_features])\n\n    # Split data\n    X = temp_data[all_features]\n    y = temp_data['responder_6']\n    del temp_data  # Free memory\n    X_train, X_val, y_train, y_val = train_test_split(X, y, test_size=0.2, random_state=42)\n    del X, y  # Free memory\n\n    # Build and train the model\n    model = build_model(input_dim=X_train.shape[1])\n    print(f\"Training on partition {partition}...\")\n    model.fit(\n        X_train, y_train,\n        validation_data=(X_val, y_val),\n        epochs=7,  # Reduced epochs\n        batch_size=2048,  # Reduced batch size\n        callbacks=[early_stopping, reduce_lr],\n        verbose=1\n    )\n    del X_train, y_train  # Free memory\n\n    # Evaluate\n    y_pred = model.predict(X_val)\n    r2 = r2_score(y_val, y_pred)\n    r2_scores.append(r2)\n    print(f\"R² Score for partition {partition}: {r2:.4f}\")\n    del X_val, y_val, y_pred  # Free memory\n\n# Final evaluation\nprint(\"Average R² Score across partitions:\", np.mean(r2_scores))","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-29T10:05:00.609070Z","iopub.execute_input":"2024-12-29T10:05:00.609313Z","iopub.status.idle":"2024-12-29T10:26:32.280053Z","shell.execute_reply.started":"2024-12-29T10:05:00.609292Z","shell.execute_reply":"2024-12-29T10:26:32.279084Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"import os\nimport kaggle_evaluation.jane_street_inference_server\nimport polars as pl\nfrom pyarrow.dataset import dataset\nfrom datasets import Dataset\nfrom datasets import load_dataset  # For handling large datasets\nimport kaggle_evaluation.jane_street_inference_server\n# Prediction Function\ndef predict(test: pl.DataFrame, lags: pl.DataFrame | None) -> pl.DataFrame | pd.DataFrame:\n    global lags_\n    if lags is not None:\n        lags_ = lags\n\n    test_df = test.to_pandas()\n    test_df.fillna(0, inplace=True)\n\n    features = [col for col in test_df.columns if 'feature' in col]\n    if lags_ is not None:\n        lags_df = lags_.to_pandas()\n        test_df = test_df.merge(\n            lags_df,\n            on=['date_id', 'time_id', 'symbol_id'],\n            how='left'\n        ).fillna(0)\n\n    all_features = features + [col for col in lags_.columns if col.startswith('responder_') and 'lag' in col]\n    missing_features = [f for f in all_features if f not in test_df.columns]\n    for feature in missing_features:\n        test_df[feature] = 0\n\n    test_features = test_df[all_features].astype(float)\n    test_features = scaler.transform(test_features)\n\n    print(\"Generating predictions...\")\n    predictions = model.predict(test_features)\n\n    predictions_df = test.select(\n        'row_id',\n        pl.lit(predictions.flatten()).alias('responder_6')\n    )\n\n    if isinstance(predictions_df, pl.DataFrame):\n        assert predictions_df.columns == ['row_id', 'responder_6']\n    elif isinstance(predictions_df, pd.DataFrame):\n        assert (predictions_df.columns == ['row_id', 'responder_6']).all()\n    else:\n        raise TypeError('The predict function must return a DataFrame')\n\n    assert len(predictions_df) == len(test)\n\n    print(\"Predictions generated successfully.\")\n    # After generating predictions in the inference process\n    # Convert predictions to Pandas DataFrame if needed\n    if isinstance(predictions_df, pl.DataFrame):\n        predictions_df = predictions_df.to_pandas()\n    \n    # Save predictions to submission.csv\n    submission_file = '/kaggle/working/submission.csv'\n    predictions_df.to_csv(submission_file, index=False)\n    print(f\"Submission file saved to {submission_file}.\")\n\n    return predictions_df\n\nlags_ = None\n\nprint(\"Initializing inference server...\")\ninference_server = kaggle_evaluation.jane_street_inference_server.JSInferenceServer(predict)\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    )\nprint(\"Inference server initialized.\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-29T10:33:26.272296Z","iopub.execute_input":"2024-12-29T10:33:26.272661Z","iopub.status.idle":"2024-12-29T10:33:26.378901Z","shell.execute_reply.started":"2024-12-29T10:33:26.272628Z","shell.execute_reply":"2024-12-29T10:33:26.377955Z"}},"outputs":[],"execution_count":null}]}