{"metadata":{"kernelspec":{"display_name":"Python 3","language":"python","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":11305158,"sourceType":"competition"},{"sourceId":10460392,"sourceType":"datasetVersion","datasetId":6468154},{"sourceId":10460393,"sourceType":"datasetVersion","datasetId":6468162},{"sourceId":215844687,"sourceType":"kernelVersion"}],"dockerImageVersionId":30787,"isInternetEnabled":false,"language":"python","sourceType":"notebook","isGpuEnabled":true},"papermill":{"default_parameters":{},"duration":32.337374,"end_time":"2025-01-12T23:05:03.060096","environment_variables":{},"exception":null,"input_path":"__notebook__.ipynb","output_path":"__notebook__.ipynb","parameters":{},"start_time":"2025-01-12T23:04:30.722722","version":"2.6.0"}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"# ==============================================================================\n# 📦 Step 1: Import Libraries and Setup Environment\n# ==============================================================================\nimport os\nimport time\nimport copy\nimport numpy as np\nimport pandas as pd\nimport polars as pl\nimport torch\n\n# Install and import the custom Jane Street package\n# Note: This path is specific to the Kaggle environment.\n# You may need to adjust it if running locally.\nwhl_path = \"/kaggle/input/janestreet2025-code/janestreet-0.1-py3-none-any.whl\"\nif os.path.exists(whl_path):\n    os.system(f\"pip install --no-deps --force-reinstall {whl_path}\")\n\n# It's good practice to place imports after potential installations\nfrom kaggle_evaluation import jane_street_inference_server\nfrom janestreet.pipeline import FullPipeline\nfrom janestreet.data_processor import DataProcessor\nfrom janestreet.config import PATH_DATA  # This will point to /kaggle/input/\n\nprint(\"✅ Libraries imported and environment set up.\")\n\n# ==============================================================================\n# ⚙️ Step 2: Set Configurations\n# ==============================================================================\n# Run Configuration\nRUN_NAME = \"full\"\n\n# Model Ensemble Configuration\nMODEL_NAMES = [\n    \"gru_2.0_700\", \"gru_2.1_700\", \"gru_2.2_700\",\n    \"gru_3.0_700\", \"gru_3.1_700\", \"gru_3.2_700\"\n]\n\n# Equal weight distribution across all models\nmodel_count = len(MODEL_NAMES)\nWEIGHTS = np.array([1.0] * model_count)\nWEIGHTS /= WEIGHTS.sum()  # Normalize weights\n\n# Rolling average window size for smoothing predictions\nN_ROLL = 1000\n\nprint(f\"🚀 Using ensemble of {model_count} models: {'_'.join(MODEL_NAMES)}\")\nprint(f\"⚖️ Normalized model weights: {WEIGHTS}\")\n\n# ==============================================================================\n# 🧠 Step 3: Load Data Processor and Models\n# ==============================================================================\nprint(\"\\n🔧 Loading data processor...\")\n# Initialize the data processor with the feature set from the first model\ndata_processor = DataProcessor(MODEL_NAMES[0]).load()\nprint(\"✅ Data processor loaded successfully!\")\n\npipelines = {}\nprint(\"\\n🔥 Initializing and loading models...\")\nfor model_name in MODEL_NAMES:\n    pipeline = FullPipeline(\n        None,\n        run_name=RUN_NAME,\n        name=model_name,\n        load_model=True,\n        features=None,\n        save_to_disc=False\n    )\n    pipeline.fit(verbose=False)  # Use verbose=False for cleaner output\n    pipelines[model_name] = pipeline\n    print(f\"  -> Model '{model_name}' loaded with {len(pipeline.features)} features.\")\n\nprint(\"✅ All models loaded into pipelines.\")\n\n# ==============================================================================\n# 📊 Step 4: Load Tail of Training Data for Rolling Features\n# ==============================================================================\nMAX_DATE = 1698\nCOLS_ID = ['row_id', 'date_id', 'time_id', 'symbol_id', 'weight', 'is_scored']\n\n# Use Polars for efficient loading and filtering\ndf_raw = (\n    pl.scan_parquet(f\"{PATH_DATA}/train.parquet\")\n    .filter(pl.col(\"date_id\") >= MAX_DATE - 10)\n    .collect()\n    .with_columns(\n        pl.lit(-1).cast(pl.Int64).alias(\"row_id\"),\n        pl.lit(True).alias(\"is_scored\"),\n        (pl.col(\"date_id\") - MAX_DATE - 1).alias(\"date_id\")\n    )\n    .select(COLS_ID + data_processor.COLS_FEATURES_INIT)\n    .filter(pl.col(\"date_id\") >= -5)\n    .sort(['date_id', 'time_id', 'symbol_id'])\n)\n\nprint(f\"\\n📈 Loaded initial training data tail with shape: {df_raw.shape}\")\n\n# Initialize states for inference\nhidden_states = [None] * len(pipelines)\ndfs = []  # To store data for model updates\n\n# ==============================================================================\n# 🧪 Step 5: Define the Inference `predict()` Function\n# ==============================================================================\nDEBUG = False\ncnt_dates = 0\n\ndef predict(test: pl.DataFrame, lags: pl.DataFrame | None) -> pl.DataFrame:\n    global df_raw, dfs, hidden_states, cnt_dates\n\n    # Get metadata from the test batch\n    date_id = test[\"date_id\"][0]\n    time_id = test[\"time_id\"][0]\n    is_scored = test[\"is_scored\"][0]\n    start_time = time.time()\n\n    # On a new day, update models with the previous day's data\n    if time_id == 0 and date_id > 0:\n        cnt_dates += 1\n        # Prepare lag data from previous predictions\n        df = pl.concat(dfs)\n        df = df.join(lags, on=[\"date_id\", \"time_id\", \"symbol_id\"], how=\"left\")\n        df = df.sort([\"date_id\", \"time_id\", \"symbol_id\"])\n        dfs = []  # Reset for the new day\n\n        # Update models if enough data is available\n        if len(df) > 968:\n            for i, (name, pipeline) in enumerate(pipelines.items()):\n                pipeline.update(df)\n\n    # Append new test data to the historical buffer\n    test_subset = test.select(df_raw.columns)\n    df_raw = pl.concat([df_raw, test_subset], how=\"vertical_relaxed\")\n    # Keep only the most recent N_ROLL samples per symbol to manage memory\n    df_raw = df_raw.group_by(\"symbol_id\").tail(N_ROLL)\n\n    # Preprocess the current data slice for prediction\n    df_cur = data_processor.process_test_data(df_raw, fast=True, date_id=date_id, time_id=time_id, symbols=test[\"symbol_id\"])\n    df_cur = df_cur.sort([\"symbol_id\"])\n    dfs.append(df_cur)\n\n    # Generate predictions if the row is scored\n    if is_scored:\n        preds = []\n        for i, (name, pipeline) in enumerate(pipelines.items()):\n            pred, hidden_states[i] = pipeline.predict(df_cur, hidden=hidden_states[i], n_times=1)\n            preds.append(pred)\n\n        # Ensemble the predictions using the defined weights\n        final_pred = np.average(preds, axis=0, weights=WEIGHTS)\n\n        # Create the prediction frame with row_id and responder_6\n        df_pred = test.select(\"row_id\").with_columns(pl.Series(\"responder_6\", final_pred))\n    else:\n        # For unscored rows, create a dummy prediction of 0.0\n        df_pred = test.select(\"row_id\").with_columns(pl.lit(0.0).alias(\"responder_6\"))\n\n    if DEBUG and time_id % 100 == 0:\n        print(f\"{date_id} {time_id:3.0f}: Processed in {time.time()-start_time:.4f}s\")\n\n    return df_pred\n\nprint(\"✅ `predict` function defined.\")\n\n# ==============================================================================\n# 🏃 Step 6: Run Local Inference Loop and Generate Submission File\n# ==============================================================================\nprint(\"\\n🏁 Starting local inference simulation...\")\n\n# Load the full test and lags data for simulation\ntest_df = pl.read_parquet(f\"{PATH_DATA}/test.parquet\")\nlags_df = pl.read_parquet(f\"{PATH_DATA}/lags.parquet\").with_columns(\n    pl.col(\"date_id\").cast(pl.Int16),\n    pl.col(\"time_id\").cast(pl.Int16),\n    pl.col(\"symbol_id\").cast(pl.Int8)\n)\n\n# Storage for predictions from each step\nall_predictions = []\n\n# Get unique time steps to iterate over\ntime_steps = test_df.select([\"date_id\", \"time_id\"]).unique().sort([\"date_id\", \"time_id\"])\n\n# Iterate over each time step to simulate the live environment\nfor date_id, time_id in time_steps.iter_rows():\n    # Filter the data for the current time step\n    test_batch = test_df.filter((pl.col(\"date_id\") == date_id) & (pl.col(\"time_id\") == time_id))\n    lags_batch = lags_df.filter((pl.col(\"date_id\") == date_id) & (pl.col(\"time_id\") == time_id))\n\n    # --- In a real submission, `is_scored` comes from the environment. ---\n    # Here we simulate it. For this example, we'll score everything.\n    test_batch = test_batch.with_columns(pl.lit(True).alias(\"is_scored\"))\n\n    # Get predictions for the current batch\n    preds = predict(test_batch, lags_batch)\n\n    # Store the results\n    all_predictions.append(preds)\n\n# Concatenate all predictions into a single DataFrame\nfinal_predictions = pl.concat(all_predictions)\n\n# Save the final predictions to submission.csv\nfinal_predictions.write_csv(\"submission.csv\")\n\nprint(\"\\n\\n✅ --- Inference Complete! --- ✅\")\nprint(f\"Submission file 'submission.csv' created successfully with {len(final_predictions)} rows.\")\nprint(\"\\nFinal Submission Head:\")\nprint(final_predictions.head())\n","metadata":{"trusted":true},"outputs":[],"execution_count":null}]}