{"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":"markdown","source":"# 🧠 Jane Street Market Prediction – Local Inference Simulation\n\nThis notebook demonstrates how to locally simulate the **inference loop** for the [Jane Street Market Prediction](https://www.kaggle.com/competitions/jane-street-market-prediction) competition, using a custom `predict()` function that processes streaming `test.parquet` and `lags.parquet` data.\n\nThe goal is to:\n- Emulate Kaggle's step-by-step test-time environment\n- Feed in daily market data and lag features\n- Make predictions using a pre-trained ensemble model\n- Generate a submission-ready CSV file with `responder_6` values\n\n### 💡 Key Components\n- `test.parquet`: Streaming market data for multiple dates and timesteps\n- `lags.parquet`: Feature-lagged data to simulate prior knowledge\n- `predict()`: Main prediction function that processes, updates and ensembles model outputs\n- `submission.parquet`: Final output file in the required format for submission\n\nThis notebook is useful for:\n- Debugging the inference pipeline locally\n- Validating model behavior across time\n- Measuring model performance before final submission\n\n🔁 Let's begin with loading the data and running the simulation step-by-step.\n\n-----------------------------","metadata":{"papermill":{"duration":0.004199,"end_time":"2025-01-12T23:04:33.984314","exception":false,"start_time":"2025-01-12T23:04:33.980115","status":"completed"},"tags":[]}},{"cell_type":"markdown","source":"### Import Libraries","metadata":{}},{"cell_type":"code","source":"#  Install the custom Jane Street package (provided by competition organizers)\nimport os\nwhl_path = \"/kaggle/input/janestreet2025-code/janestreet-0.1-py3-none-any.whl\"\nos.system(f\"pip install --force-reinstall {whl_path}\")\n\n#  Import core utilities\nimport time   # for timing\nimport copy   # for deep copying objects\n\n# Import essential libraries for data & ML\nimport numpy as np\nimport pandas as pd\nimport polars as pl       # optional, faster dataframe library\nimport torch              # for deep learning models\n\n# Kaggle's evaluation module for inference (do not modify)\nfrom kaggle_evaluation import jane_street_inference_server\n\n# Import pipeline modules from the installed Jane Street package\nfrom janestreet.pipeline import FullPipeline, PipelineCV   # training & cross-validation pipelines\nfrom janestreet.data_processor import DataProcessor        # handles preprocessing\nfrom janestreet.config import PATH_DATA                    # config path to input data\n","metadata":{"execution":{"iopub.status.busy":"2025-06-17T21:13:04.450316Z","iopub.execute_input":"2025-06-17T21:13:04.450620Z","iopub.status.idle":"2025-06-17T21:13:50.362295Z","shell.execute_reply.started":"2025-06-17T21:13:04.450582Z","shell.execute_reply":"2025-06-17T21:13:50.361602Z"},"papermill":{"duration":17.717942,"end_time":"2025-01-12T23:04:51.706046","exception":false,"start_time":"2025-01-12T23:04:33.988104","status":"completed"},"tags":[],"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## SEt Configurations","metadata":{"papermill":{"duration":0.002354,"end_time":"2025-01-12T23:04:51.711232","exception":false,"start_time":"2025-01-12T23:04:51.708878","status":"completed"},"tags":[]}},{"cell_type":"code","source":"# Run Configuration\nRUN_NAME = \"full\"  # Run ID or tag to track experiments\n\n#  Model Ensemble Configuration\nMODEL_NAMES = [\n    \"gru_2.0_700\", \n    \"gru_2.1_700\", \n    \"gru_2.2_700\", \n    \"gru_3.0_700\", \n    \"gru_3.1_700\", \n    \"gru_3.2_700\"\n]\n\n# Equal weight distribution across all models (normalized)\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\n# Log current ensemble setup\nprint(\"Using ensemble of models:\", \"_\".join(MODEL_NAMES))\nprint(\"Normalized model weights:\", WEIGHTS)\n","metadata":{"execution":{"iopub.status.busy":"2025-06-17T21:13:50.363860Z","iopub.execute_input":"2025-06-17T21:13:50.364129Z","iopub.status.idle":"2025-06-17T21:13:50.369982Z","shell.execute_reply.started":"2025-06-17T21:13:50.364104Z","shell.execute_reply":"2025-06-17T21:13:50.369134Z"},"papermill":{"duration":0.009996,"end_time":"2025-01-12T23:04:51.723697","exception":false,"start_time":"2025-01-12T23:04:51.713701","status":"completed"},"tags":[],"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"# Load models","metadata":{"papermill":{"duration":0.00225,"end_time":"2025-01-12T23:04:51.728408","exception":false,"start_time":"2025-01-12T23:04:51.726158","status":"completed"},"tags":[]}},{"cell_type":"code","source":"# -------------------------------------------------------------------------------------\n# 🔧 Step 1: Load Data Processor\n# -------------------------------------------------------------------------------------\nfrom janestreet.data_processor import DataProcessor\n\nprint(\"🔧 Loading data processor using:\", MODEL_NAMES[0])\ndata_processor = DataProcessor(MODEL_NAMES[0]).load()\nprint(\"✅ Data processor loaded successfully!\\n\")\n\n# -------------------------------------------------------------------------------------\n# 🔁 Step 2: Initialize Pipelines for Each Model\n# -------------------------------------------------------------------------------------\npipelines = {}\n\nfor model_name in MODEL_NAMES:\n    print(f\"\\n🚀 Setting up pipeline for: {model_name}\")\n    \n    pipeline = FullPipeline(\n        None,                   # Required positional arg — model will be loaded via joblib\n        run_name=RUN_NAME,\n        name=model_name,\n        load_model=True,\n        features=None,          # Features will be auto-loaded\n        save_to_disc=False\n    )\n    \n    pipeline.fit(verbose=True)\n    pipelines[model_name] = pipeline\n    \n    print(\"-\" * 80)\n    print(f\"📌 Model Name: {model_name}\")\n    print(f\"🛠️  Model Params: {pipeline.model.get_params()}\")\n    print(f\"🧬 Number of Features: {len(pipeline.features)}\")\n    print(f\"🎯 Number of Response Targets: {pipeline.model.model.num_resp}\")\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-06-17T21:13:50.371058Z","iopub.execute_input":"2025-06-17T21:13:50.371292Z","iopub.status.idle":"2025-06-17T21:13:54.879421Z","shell.execute_reply.started":"2025-06-17T21:13:50.371271Z","shell.execute_reply":"2025-06-17T21:13:54.878432Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"# Load tail of train data for rolling features calculation","metadata":{"papermill":{"duration":0.002733,"end_time":"2025-01-12T23:04:57.515147","exception":false,"start_time":"2025-01-12T23:04:57.512414","status":"completed"},"tags":[]}},{"cell_type":"code","source":"# -------------------------------------------------------------------------------------\n# 🔢 Step 1: Define constants and ID columns\n# -------------------------------------------------------------------------------------\nMAX_DATE = 1698\nCOLS_ID = ['row_id', 'date_id', 'time_id', 'symbol_id', 'weight', 'is_scored']\n\n# -------------------------------------------------------------------------------------\n# 📦 Step 2: Load last 10 days of data from training set using Polars (Lazy mode)\n# -------------------------------------------------------------------------------------\ndf_raw = pl.scan_parquet(f\"{PATH_DATA}/train.parquet\")  # Lazy load to save memory\ndf_raw = df_raw.filter(pl.col(\"date_id\") >= MAX_DATE - 10)  # Filter for last 10 days only\ndf_raw = df_raw.collect()  # Collect into memory (eager)\n\n# -------------------------------------------------------------------------------------\n# 🛠️ Step 3: Preprocess columns to match expected structure for inference\n# -------------------------------------------------------------------------------------\ndf_raw = df_raw.with_columns(\n    pl.lit(-1).cast(pl.Int64).alias(\"row_id\"),            # Dummy row_id for compatibility\n    pl.lit(True).alias(\"is_scored\"),                      # Mark all rows as scored\n    (pl.col(\"date_id\") - MAX_DATE - 1).alias(\"date_id\")   # Shift date for evaluation window\n)\n\n# -------------------------------------------------------------------------------------\n# 🎯 Step 4: Select only ID columns + feature columns required by the model\n# -------------------------------------------------------------------------------------\ndf_raw = df_raw.select(COLS_ID + data_processor.COLS_FEATURES_INIT)\n\n# -------------------------------------------------------------------------------------\n# 📊 Step 5: Sort data properly before feeding to models\n# -------------------------------------------------------------------------------------\ndf_raw = (\n    df_raw\n    .filter(pl.col(\"date_id\") >= -5)  # Keep most recent 5 \"virtual\" days\n    .sort(['date_id', 'time_id', 'symbol_id'])  # Ensure time consistency\n)\n\n# -------------------------------------------------------------------------------------\n# 🧠 Step 6: Prepare for inference\n# -------------------------------------------------------------------------------------\nhidden_states = [None] * len(pipelines)  # One hidden state per model\ndfs = []  # To store model outputs / predictions\n","metadata":{"execution":{"iopub.status.busy":"2025-06-17T21:14:17.439470Z","iopub.execute_input":"2025-06-17T21:14:17.440280Z","iopub.status.idle":"2025-06-17T21:14:18.875648Z","shell.execute_reply.started":"2025-06-17T21:14:17.440245Z","shell.execute_reply":"2025-06-17T21:14:18.874936Z"},"papermill":{"duration":1.728541,"end_time":"2025-01-12T23:04:59.246573","exception":false,"start_time":"2025-01-12T23:04:57.518032","status":"completed"},"tags":[],"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## Inference Pipeline (Jane Street)\n### Setup Constants and Time Tracking","metadata":{}},{"cell_type":"code","source":"DEBUG = False  # Set to True for internal debug logs\n\nCNT_DATES = 9  # Total test days\nCNT_DATES_NOT_SCORED = 4  # First few days are unscored\n\n# Track execution time\ntime_start = time.time()\ntime_start_not_scored = time.time()\ntime_est = 0\ntime_est_not_scored = 0\ncnt_dates = 0  # To track daily iterations\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-06-17T21:14:22.670367Z","iopub.execute_input":"2025-06-17T21:14:22.670674Z","iopub.status.idle":"2025-06-17T21:14:22.675268Z","shell.execute_reply.started":"2025-06-17T21:14:22.670649Z","shell.execute_reply":"2025-06-17T21:14:22.674344Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### 🧪 Define predict() function used during inference","metadata":{}},{"cell_type":"code","source":"def predict(test: pl.DataFrame, lags: pl.DataFrame | None) -> pl.DataFrame:\n    \"\"\"Main prediction logic for each test timestep.\"\"\"\n\n    global df_raw, hidden_states, pipeline, dfs\n    global time_est, time_start, time_start_not_scored, time_est_not_scored, cnt_dates\n\n    # Extract current timestamp\n    date_id = test[\"date_id\"][0]\n    time_id = test[\"time_id\"][0]\n    is_scored = test[\"is_scored\"][0]\n\n    # Time tracking for debug\n    if DEBUG:\n        if not is_scored:\n            time_est_not_scored = time.time() - time_start_not_scored\n        else:\n            time_est = time.time() - time_start\n        if time_id == 0:\n            print(\"-\" * 100)\n            if date_id == 1:\n                time_start_not_scored = time.time()\n            if date_id == CNT_DATES_NOT_SCORED:\n                time_start = time.time()\n\n    # 🧹 Reset states and update with lag features\n    if time_id == 0:\n        cnt_dates += 1\n        hidden_states = [None] * len(pipelines)\n        lags = lags.with_columns(\n            pl.col(\"responder_6_lag_1\").alias(\"responder_6\"),\n            pl.lit(date_id - 1).cast(pl.Int16).alias(\"date_id\")\n        ).select([\"date_id\", \"time_id\", \"symbol_id\", \"responder_6\"])\n\n        # 🧩 Join lag data to prior features (only after day 0)\n        if cnt_dates > 1:\n            df = pl.concat(dfs)\n            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\n    # 🧪 Append new test row to raw data\n    test = test.select(df_raw.columns)  # Match input format\n    df_raw = pl.concat([df_raw, test], how=\"vertical_relaxed\")\n    df_raw = df_raw.select(test.columns)\n\n    # ⏳ Only keep recent N_ROLL samples per symbol\n    df_raw = df_raw.group_by(\"symbol_id\").tail(N_ROLL)\n\n    # 🔧 Preprocess using data_processor\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    df_cur = df_cur.with_columns(pl.lit(None).alias(\"responder_6\"))\n\n    # 🔁 Update model weights using full history\n    if (time_id == 0) & (cnt_dates > 1):\n        if len(df) > 968:  # Only update if there's enough data\n            for i, (name, pipeline) in enumerate(pipelines.items()):\n                pipeline.update(df)\n\n    # 📊 Ensemble Predictions\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        pred = np.average(preds, axis=0, weights=WEIGHTS)\n\n        df_cur = df_cur.with_columns(pl.Series(\"responder_6\", pred))\n        df_cur = test.select([\"date_id\", \"time_id\", \"symbol_id\"]).join(df_cur, on=[\"date_id\", \"time_id\", \"symbol_id\"], how=\"left\")\n        predictions = df_cur.select([\"row_id\", \"responder_6\"])\n    else:\n        predictions = test.select(\"row_id\", pl.lit(0.0).alias(\"responder_6\"))\n\n    # 🐞 Debug output\n    if DEBUG and time_id % 100 == 0:\n        n_nans = sum(sum(predictions.fill_nan(None).null_count().to_numpy()))\n        print(\n            f\"{date_id} {time_id:3.0f} (is_scored {is_scored}): \"\n            f\"time {time.time()-start_time:.4f}, # NaNs: {n_nans}\"\n        )\n    elif (time_id == 0) and (date_id == 0):\n        print(predictions)\n\n    return predictions\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-06-17T21:14:24.052586Z","iopub.execute_input":"2025-06-17T21:14:24.052879Z","iopub.status.idle":"2025-06-17T21:14:24.065572Z","shell.execute_reply.started":"2025-06-17T21:14:24.052852Z","shell.execute_reply":"2025-06-17T21:14:24.064841Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### 📡 Initialize Inference Server (local or competition)","metadata":{}},{"cell_type":"code","source":"import polars as pl\nimport os\n\nprint(os.listdir(PATH_DATA))\n\ntest_df = pl.read_parquet(f\"{PATH_DATA}/test.parquet\")\ntest_df = test_df.with_columns(pl.lit(True).alias(\"is_scored\"))\n\nlags_df = pl.read_parquet(f\"{PATH_DATA}/lags.parquet\")\n\ndisplay(test_df.head())\ndisplay(lags_df.head())\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-06-17T21:14:48.939952Z","iopub.execute_input":"2025-06-17T21:14:48.940748Z","iopub.status.idle":"2025-06-17T21:14:48.993086Z","shell.execute_reply.started":"2025-06-17T21:14:48.940695Z","shell.execute_reply":"2025-06-17T21:14:48.992362Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# make sure column types align\nlags_df = lags_df.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\nall_predictions = []\n\n# Unique time steps\ntime_steps = test_df.select([\"date_id\", \"time_id\"]).unique()\n\n# Iterate over each timestep\nfor row in time_steps.iter_rows():\n    date_id, time_id = row\n\n    # Filter test and lags for current 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    # 🟢 Force scoring ON\n    test_batch = test_batch.with_columns(pl.lit(True).alias(\"is_scored\"))\n\n    # Predict\n    preds = predict(test_batch, lags_batch)\n\n    # Collect results\n    all_predictions.append(preds)\n\n# Concatenate all predictions\nfinal_predictions = pl.concat(all_predictions)\n\n# Save submission format\nfinal_predictions.write_csv(\"submission.parquet\")\n","metadata":{"papermill":{"duration":0.002894,"end_time":"2025-01-12T23:05:00.570260","exception":false,"start_time":"2025-01-12T23:05:00.567366","status":"completed"},"tags":[],"trusted":true,"execution":{"iopub.status.busy":"2025-06-17T21:14:55.562526Z","iopub.execute_input":"2025-06-17T21:14:55.562966Z","iopub.status.idle":"2025-06-17T21:14:56.133670Z","shell.execute_reply.started":"2025-06-17T21:14:55.562925Z","shell.execute_reply":"2025-06-17T21:14:56.132805Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"# ✅ Conclusion\n\nIn this notebook, we built and validated a **local inference simulation** for the Jane Street Market Prediction challenge. This setup closely mimics the real-time evaluation logic used by Kaggle, allowing us to test and improve our model with confidence.\n\n### 📌 Key Highlights:\n- Processed `test.parquet` and `lags.parquet` files to simulate live data stream.\n- Successfully integrated the `predict()` function with temporal logic and lag updates.\n- Produced valid predictions for `responder_6` and prepared a ready-to-submit file.\n\n### 🚀 Next Steps:\n- Enhance the ensemble model performance using better validation strategies.\n- Experiment with different lag features and rolling window sizes.\n- Implement runtime optimizations for faster inference.\n\n---\n\n# 🤝 Let’s Connect!\n\nIf you found this project helpful, feel free to connect and collaborate:\n\n- 💼 [LinkedIn](https://www.linkedin.com/in/sheemamasood/)  \n- 💻 [GitHub](https://github.com/SheemaMasood381)\n-   \n- 📬 Drop a message for project feedback or collaboration ideas!\n\nTogether, let's keep building impactful and intelligent solutions! 💡🚀\n\n--------\n","metadata":{}},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null}]}