{"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":"none","dataSources":[{"sourceId":84493,"databundleVersionId":9871156,"sourceType":"competition"},{"sourceId":9859479,"sourceType":"datasetVersion","datasetId":6050832},{"sourceId":9871648,"sourceType":"datasetVersion","datasetId":6059588},{"sourceId":10285006,"sourceType":"datasetVersion","datasetId":6060014},{"sourceId":161606,"sourceType":"modelInstanceVersion","modelInstanceId":137419,"modelId":160126},{"sourceId":162898,"sourceType":"modelInstanceVersion","modelInstanceId":138527,"modelId":161177},{"sourceId":163785,"sourceType":"modelInstanceVersion","modelInstanceId":139311,"modelId":161940},{"sourceId":163933,"sourceType":"modelInstanceVersion","modelInstanceId":139440,"modelId":162071},{"sourceId":166855,"sourceType":"modelInstanceVersion","modelInstanceId":141974,"modelId":164550},{"sourceId":208287,"sourceType":"modelInstanceVersion","modelInstanceId":177575,"modelId":199880},{"sourceId":208464,"sourceType":"modelInstanceVersion","isSourceIdPinned":true,"modelInstanceId":177724,"modelId":200030}],"dockerImageVersionId":30786,"isInternetEnabled":false,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"from xgboost import XGBRegressor\nimport numpy as np\nimport joblib\nimport json\nimport os\n\nimport pandas as pd\nimport polars as pl\nfrom copy import deepcopy\nimport kaggle_evaluation.jane_street_inference_server\ndef r2_xgb(y_true, y_pred, sample_weight):\n    r2 = 1 - np.average((y_pred - y_true) ** 2, weights=sample_weight) / (np.average((y_true) ** 2, weights=sample_weight) + 1e-38)\n    return -r2","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-31T08:11:06.717311Z","iopub.execute_input":"2024-12-31T08:11:06.717688Z","iopub.status.idle":"2024-12-31T08:11:06.727864Z","shell.execute_reply.started":"2024-12-31T08:11:06.717662Z","shell.execute_reply":"2024-12-31T08:11:06.724549Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"class ModelGroup:\n    def __init__(self):\n        self.models = []\n\n    def add_model(self, model):\n        \"\"\"Add a trained model to the group.\"\"\"\n        self.models.append(model)\n\n    def predict(self, test_data):\n        \"\"\"Make predictions using all models in the group and return the average.\"\"\"\n        preds = []\n        for model in self.models:\n            \n            pred = model.predict(test_data)\n            \n            preds.append(pred)\n        \n        # Average the predictions from all models\n        avg_pred = np.mean(preds, axis=0)\n        return avg_pred\n    \n    @classmethod\n    def load(cls, file_path):\n        \"\"\"Load a model group from a file.\"\"\"\n        model_group = joblib.load(file_path)\n        return model_group\n\nmodel_group = ModelGroup.load('/kaggle/input/xgb_4/scikitlearn/default/1/xgb_model_group_4.pkl')","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-31T08:11:06.733563Z","iopub.execute_input":"2024-12-31T08:11:06.734444Z","iopub.status.idle":"2024-12-31T08:11:06.820254Z","shell.execute_reply.started":"2024-12-31T08:11:06.734411Z","shell.execute_reply":"2024-12-31T08:11:06.818777Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"import numpy as np\n\n\nclass data_processor:\n    def __init__(self, df, func, kwargs) -> None:\n        self.df = df\n        self.func = func\n        self.kwargs = kwargs\n\n    def check_consistency(self):\n        raise NotImplementedError\n\n    def run(self):\n        df = self.func(self.df, **self.kwargs)\n        # if feature_name in self.df.columns:\n        #     Y = input(f\"{feature_name} already exists. Do you want to overwrite it? (y/n)\")\n        #     if Y == 'n':\n        #         return self.df, feature_name\n        #     return self.df, feature_name\n\n        return df\n\n\ndef format_pl(x):\n    x = x.with_columns([pl.col(column).cast(pl.Float64) for column in x.columns])\n    return x\n\n\ndef lagging_last_value(df):\n\n    lag_cols_original = [\"date_id\", \"symbol_id\"] + [\n        f\"responder_{idx}\" for idx in range(9)\n    ]\n    lag_cols_rename = {f\"responder_{idx}\": f\"responder_{idx}_lag_1\" for idx in range(9)}\n\n    lags = df.select(pl.col(lag_cols_original))\n    lags = lags.rename(lag_cols_rename)\n    lags = lags.with_columns(\n        date_id=pl.col(\"date_id\") + 1,  # lagged by 1 day\n    )\n    lags = lags.group_by(\n        [\"date_id\", \"symbol_id\"], maintain_order=True\n    ).last()  # pick up last record of previous date\n    df = df.join(lags, on=[\"date_id\", \"symbol_id\"], how=\"left\")\n    return df\n\n\ndef rolling_mean(df, column, window):\n    # Calculate the rolling mean for each group (symbol_id) using the `rolling_mean` method\n    rolling_mean_column = f\"{column}_rolling_mean_{window}\"\n    df = df.sort([\"symbol_id\", \"date_id\", \"time_id\"])\n    df = df.with_columns(\n        pl.col(column)\n        .rolling_mean(window_size=window)\n        .over(\"symbol_id\")\n        .alias(rolling_mean_column)\n    )\n    return df\n\n\ndef rolling_zscore(df, column, window):\n    zscore_column = \"responder_6_zscore\"\n    df = df.sort([\"symbol_id\", \"date_id\", \"time_id\"])\n    col = column\n\n    df = df.with_columns(\n        (\n            (\n                pl.col(col)\n                - pl.col(col).rolling_mean(window_size=window).over(\"symbol_id\")\n            )\n            / pl.col(col).rolling_std(window_size=window).over(\"symbol_id\")\n        ).alias(f\"{col}_zscore\")\n    )\n    return df","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-31T08:11:06.821888Z","iopub.execute_input":"2024-12-31T08:11:06.822235Z","iopub.status.idle":"2024-12-31T08:11:06.845379Z","shell.execute_reply.started":"2024-12-31T08:11:06.822203Z","shell.execute_reply":"2024-12-31T08:11:06.843700Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"feature_list = ['symbol_id', 'time_id']\nfeature_list += [f\"feature_{idx:02d}\" for idx in range(79)]\nfeature_list += [f'responder_{i}_lag_1' for i in range(9)]\nfeature_list += [f\"feature_{idx:02d}_zscore\" for idx in range(79) if idx != 9 and idx != 10 and idx != 11]\nlags_ : pl.DataFrame | None = None\ncnt = 0\nmean_ : pl.DataFrame | None = None\nstd_ : pl.DataFrame | None = None\n\ndef format_pl(x):\n    x = x.with_columns([\n        pl.col(column).cast(pl.Float64) for column in x.columns\n    ])\n    return x\n\ntest_history = pl.read_parquet('/kaggle/input/d/yechaochen/test-history/test_history.parquet')\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 1 minute 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    global test_history\n    global cnt\n    global mean_\n    global std_\n    \n    test = format_pl(test)\n    if lags is not None:\n        lags = format_pl(lags)\n        lags = lags.group_by([\"date_id\", \"symbol_id\"], maintain_order=True).last() # pick up last record of previous date\n        test = test.join(lags, on=[\"date_id\", \"symbol_id\"], how=\"left\").drop('time_id_right')\n        if lags_ is None:\n            lags_ = lags\n        else:\n            lags_ = pl.concat([lags_, lags])\n        lags_ = lags.group_by([\"date_id\", \"symbol_id\"], maintain_order=True).last()\n        lags_ = format_pl(lags_)\n    else:\n        test = test.join(lags_, on=[\"date_id\", \"symbol_id\"], how=\"left\").drop('time_id_right')\n    \n    test_history = test_history[test.columns]\n    test_history = format_pl(test_history)\n    test_history = pl.concat([test_history, test])\n    # test_history = test_history.unique(['symbol_id', 'date_id', 'time_id'])\n    \n    test_history = (\n        test_history\n        .sort(by=[\"symbol_id\", \"date_id\", \"time_id\"])\n        .with_row_index(name=\"index\")  # 使用 with_row_index() 替代 with_row_count()\n        .with_columns(\n            pl.col(\"index\").over(\"symbol_id\").alias(\"rank\")\n        )\n        .filter(pl.col(\"rank\") > (pl.col(\"rank\").max().over(\"symbol_id\")-4000.0-pl.lit(test.shape[0]).cast(pl.Float64)))\n        .drop([\"rank\", \"index\"])\n    )\n    test_history = format_pl(test_history)\n    X = deepcopy(test_history)\n    cols = []\n    cols += [f\"feature_{idx:02d}\" for idx in range(79) if idx != 9 and idx != 10 and idx != 11]\n    if cnt % 4000 == 0:\n        mean_ = test_history.select([pl.col(col).mean().alias(col) for col in cols])\n        std_ = test_history.select([pl.col(col).std().alias(col) for col in cols])\n    cnt+=1\n    mean_dict = mean_.to_dict(as_series=False)\n    std_dict = std_.to_dict(as_series=False)\n    X = format_pl(X)\n    \n    real_test = test.select(['date_id', 'time_id', 'symbol_id']).join(X, on=['date_id', 'time_id', 'symbol_id'], how='left')\n    normalized_columns = [\n        ((pl.col(col) - mean_dict[col][0]) / std_dict[col][0]).alias(f\"{col}_zscore\")\n        for col in cols\n    ]\n    normalized_df = test.select(normalized_columns)\n    real_test = pl.concat([test, normalized_df], how=\"horizontal\")\n    \n    \n    real_test = real_test[feature_list].to_numpy()\n    real_test = np.nan_to_num(real_test, nan=0, posinf=1000, neginf=-1000)\n    pred = model_group.predict(real_test)\n    \n    # Replace this section with your own predictions\n    predictions = test.select(\n        'row_id'\n    ).with_columns(\n        pl.Series(\n            name   = 'responder_6', \n            values = np.clip(pred, a_min = -5, a_max = 5),\n            dtype  = pl.Float64,\n        )\n    )\n    predictions = predictions.with_columns(pl.Series(\"responder_6\", pred))\n    predictions = format_pl(predictions)\n\n    if isinstance(predictions, pl.DataFrame):\n        assert predictions.columns == ['row_id', 'responder_6']\n    elif isinstance(predictions, pd.DataFrame):\n        assert (predictions.columns == ['row_id', 'responder_6']).all()\n    else:\n        raise TypeError('The predict function must return a DataFrame')\n    # Confirm has as many rows as the test data.\n    assert len(predictions) == len(test)\n    \n    return predictions\n\ninference_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":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-31T08:11:06.860773Z","iopub.execute_input":"2024-12-31T08:11:06.861152Z","iopub.status.idle":"2024-12-31T08:11:07.686833Z","shell.execute_reply.started":"2024-12-31T08:11:06.861126Z","shell.execute_reply":"2024-12-31T08:11:07.684826Z"}},"outputs":[],"execution_count":null}]}