{"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"},{"sourceId":9756372,"sourceType":"datasetVersion","datasetId":5973876},{"sourceId":9801075,"sourceType":"datasetVersion","datasetId":6006872},{"sourceId":9806342,"sourceType":"datasetVersion","datasetId":6010899},{"sourceId":203900450,"sourceType":"kernelVersion"}],"dockerImageVersionId":30823,"isInternetEnabled":false,"language":"python","sourceType":"notebook","isGpuEnabled":true}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"#########################\n# 1) 必要なライブラリ読み込み\n#########################\nimport os\nimport gc\nimport pickle\n\nimport numpy as np\nimport pandas as pd\n\n# Polars は pls としてインポート (Lightning の pl と衝突しない)\nimport polars as pls\n\n# PyTorch Lightning は pl としてインポート\nimport pytorch_lightning as pl\nfrom pytorch_lightning import Trainer, LightningModule, LightningDataModule\nfrom pytorch_lightning.callbacks import EarlyStopping, ModelCheckpoint\n\n# PyTorch\nimport torch\nimport torch.nn as nn\nimport torch.nn.functional as F\n\n# Sklearn\nfrom sklearn.metrics import r2_score\nfrom sklearn.model_selection import train_test_split\n\n# XGBoost & CatBoost\nfrom xgboost import XGBRegressor\nfrom catboost import CatBoostRegressor\n\nimport warnings\nwarnings.filterwarnings('ignore')\npd.options.display.max_columns = None\n\n\n#########################\n# 2) CONFIGクラス (設定管理)\n#########################\nclass CONFIG:\n    seed = 42\n    target_col = \"responder_6\"\n    # 学習で使用する特徴量\n    feature_cols = [f\"feature_{idx:02d}\" for idx in range(79)] + [f\"responder_{idx}_lag_1\" for idx in range(9)]\n    \n    # 学習済みモデルを保存するパス (Kaggle なら \"/kaggle/working\" なども可)\n    model_save_path_xgb = \"trained_xgb.pkl\"\n    model_save_path_cat = \"trained_cat.pkl\"\n    model_save_folder_nn = \"trained_nn\"\n\n\n#########################\n# 3) データ読み込み (例)\n#########################\n# 本来は train.parquet と validation.parquet を分けるなどしますが、\n# 例として \"/kaggle/input/js24-preprocessing-create-lags/validation.parquet/\" を全読み込みして、\n# その中で train/val に分割しています。\n\n# Polars で parquet 読み込み\nvalid_pl = pls.scan_parquet(\"/kaggle/input/js24-preprocessing-create-lags/validation.parquet/\").collect()\ndf = valid_pl.to_pandas()\n\n# フィーチャー & ターゲット & weight\nX_all = df[CONFIG.feature_cols].fillna(method=\"ffill\").fillna(0)\ny_all = df[CONFIG.target_col]\nw_all = df[\"weight\"]\n\nprint(\"X_valid.shape =\", X_all.shape)\nprint(\"y_valid.shape =\", y_all.shape)\n\n# 例: hold-out で train / validation に分割 (本来は時系列を意識した分割推奨)\nX_train, X_val, y_train, y_val, w_train, w_val = train_test_split(\n    X_all, y_all, w_all, test_size=0.2, random_state=CONFIG.seed\n)\nprint(\"Train size:\", X_train.shape)\nprint(\"Validation size:\", X_val.shape)\n\ndel df, valid_pl, X_all, y_all, w_all\ngc.collect()\n\n\n#########################\n# 4) XGBoost モデル学習\n#########################\nxgb_model = XGBRegressor(\n    random_state=CONFIG.seed,\n    n_estimators=1000,\n    max_depth=8,\n    learning_rate=0.01,\n    subsample=0.8,\n    colsample_bytree=0.8,\n    tree_method='gpu_hist',  # GPUが使える場合\n    reg_alpha=2.0,\n    reg_lambda=2.0,\n)\n\nxgb_model.fit(\n    X_train, y_train,\n    sample_weight=w_train,\n    eval_set=[(X_val, y_val)],\n    eval_metric='rmse',\n    sample_weight_eval_set=[w_val],\n    early_stopping_rounds=50,\n    verbose=100\n)\n\n\n#########################\n# 5) CatBoost モデル学習\n#########################\ncat_model = CatBoostRegressor(\n    iterations=1000,\n    depth=8,\n    learning_rate=0.01,\n    random_seed=CONFIG.seed,\n    task_type=\"GPU\",  # GPUが使える場合\n    verbose=100,\n)\ncat_model.fit(\n    X_train, y_train,\n    sample_weight=w_train,\n    eval_set=(X_val, y_val),\n    use_best_model=True,\n)\n\n\n#########################\n# 6) PyTorch Lightning 用の NN モデル定義\n#########################\ndef r2_val(y_true, y_pred, sample_weight):\n    \"\"\"\n    Weighted R^2 スコアを計算\n    \"\"\"\n    y_true = np.array(y_true, dtype=float)\n    y_pred = np.array(y_pred, dtype=float)\n    w = np.array(sample_weight, dtype=float)\n\n    # np.average() を利用するためには形状が揃っている必要があります\n    numerator = np.average((y_pred - y_true) ** 2, weights=w)\n    denominator = np.average(y_true ** 2, weights=w) + 1e-38\n    r2 = 1.0 - numerator / denominator\n    return r2\n\n\nclass NNModel(LightningModule):\n    def __init__(\n        self, \n        input_dim=len(CONFIG.feature_cols),\n        hidden_dims=[512, 512, 256],\n        dropouts=[0.2, 0.2, 0.2],\n        lr=1e-3,\n        weight_decay=1e-6\n    ):\n        super().__init__()\n        self.save_hyperparameters()\n\n        layers = []\n        in_dim = input_dim\n        for i, hidden_dim in enumerate(hidden_dims):\n            layers.append(nn.BatchNorm1d(in_dim))\n            if i > 0:\n                layers.append(nn.SiLU())\n            if i < len(dropouts):\n                layers.append(nn.Dropout(dropouts[i]))\n            layers.append(nn.Linear(in_dim, hidden_dim))\n            in_dim = hidden_dim\n        layers.append(nn.Linear(in_dim, 1))\n        layers.append(nn.Tanh())\n\n        self.model = nn.Sequential(*layers)\n        self.lr = lr\n        self.weight_decay = weight_decay\n        self.validation_step_outputs = []\n\n    def forward(self, x):\n        return 5 * self.model(x).squeeze(-1)  # Tanh -> [-1,1] → [-5,5] でスケーリング\n\n    def training_step(self, batch, batch_idx):\n        x, y, w = batch\n        y_hat = self(x)\n        loss = F.mse_loss(y_hat, y, reduction='none') * w\n        loss = loss.mean()\n        self.log('train_loss', loss, on_step=False, on_epoch=True, batch_size=x.size(0))\n        return loss\n\n    def validation_step(self, batch, batch_idx):\n        x, y, w = batch\n        y_hat = self(x)\n        loss = F.mse_loss(y_hat, y, reduction='none') * w\n        loss = loss.mean()\n        self.log('val_loss', loss, on_step=False, on_epoch=True, batch_size=x.size(0))\n        self.validation_step_outputs.append((y_hat.detach(), y.detach(), w.detach()))\n        return loss\n\n    def on_validation_epoch_end(self):\n        # 各ミニバッチで蓄積した (y_hat, y, w) をまとめて R^2 を計算\n        y_pred_list = []\n        y_list = []\n        w_list = []\n        for y_hat, y, w in self.validation_step_outputs:\n            y_pred_list.append(y_hat.cpu().numpy())\n            y_list.append(y.cpu().numpy())\n            w_list.append(w.cpu().numpy())\n\n        # 連結したうえで (N,) の形に揃える\n        y_all = np.concatenate(y_list, axis=0).squeeze()\n        y_pred_all = np.concatenate(y_pred_list, axis=0).squeeze()\n        w_all = np.concatenate(w_list, axis=0).squeeze()\n\n        val_r2 = r2_val(y_all, y_pred_all, w_all)\n        self.log(\"val_r_square\", val_r2, prog_bar=True)\n\n        self.validation_step_outputs.clear()\n\n    def configure_optimizers(self):\n        optimizer = torch.optim.Adam(self.parameters(), lr=self.lr, weight_decay=self.weight_decay)\n        scheduler = torch.optim.lr_scheduler.ReduceLROnPlateau(optimizer, mode='min', factor=0.5, patience=5, verbose=True)\n        return {\n            'optimizer': optimizer,\n            'lr_scheduler': {\n                'scheduler': scheduler,\n                'monitor': 'val_loss',\n            }\n        }\n\n\n#########################\n# 7) Dataset & DataModule\n#########################\nfrom torch.utils.data import Dataset, DataLoader\n\nclass MarketDataset(Dataset):\n    def __init__(self, X, y, w):\n        self.X = X\n        self.y = y\n        self.w = w\n\n    def __len__(self):\n        return len(self.X)\n\n    def __getitem__(self, idx):\n        return (\n            torch.FloatTensor(self.X.iloc[idx].values),\n            torch.FloatTensor([self.y.iloc[idx]]),\n            torch.FloatTensor([self.w.iloc[idx]])\n        )\n\n# ★LightningDataModule を継承★\nclass MarketLightningDataModule(LightningDataModule):\n    def __init__(self, X_train, y_train, w_train, X_val, y_val, w_val, batch_size=4096):\n        super().__init__()\n        self.X_train = X_train\n        self.y_train = y_train\n        self.w_train = w_train\n        self.X_val = X_val\n        self.y_val = y_val\n        self.w_val = w_val\n        self.batch_size = batch_size\n\n    def setup(self, stage=None):\n        # 必要であれば、train/val/test の再分割や前処理などを行う\n        pass\n\n    def train_dataloader(self):\n        ds = MarketDataset(self.X_train, self.y_train, self.w_train)\n        return DataLoader(ds, batch_size=self.batch_size, shuffle=True, num_workers=2)\n\n    def val_dataloader(self):\n        ds = MarketDataset(self.X_val, self.y_val, self.w_val)\n        return DataLoader(ds, batch_size=self.batch_size, shuffle=False, num_workers=2)\n\n\n#########################\n# 8) NN モデルの学習\n#########################\nos.makedirs(CONFIG.model_save_folder_nn, exist_ok=True)\n\nnn_model = NNModel(\n    input_dim=len(CONFIG.feature_cols),\n    hidden_dims=[512, 512, 256],\n    dropouts=[0.2, 0.2, 0.2],\n    lr=1e-3,\n    weight_decay=1e-6\n)\n\ndm = MarketLightningDataModule(X_train, y_train, w_train, X_val, y_val, w_val, batch_size=4096)\n\ntrainer = pl.Trainer(\n    max_epochs=20,\n    precision=16,\n    accelerator=\"gpu\" if torch.cuda.is_available() else \"cpu\",\n    callbacks=[\n        EarlyStopping(monitor=\"val_loss\", mode=\"min\", patience=5),\n        ModelCheckpoint(\n            dirpath=CONFIG.model_save_folder_nn,\n            filename=\"nn_model\",\n            monitor=\"val_loss\",\n            mode=\"min\",\n            save_top_k=1\n        )\n    ]\n)\n\ntrainer.fit(nn_model, datamodule=dm)\n\n# チェックポイントのパス\nckpt_path = trainer.checkpoint_callback.best_model_path\nprint(\"Best NN checkpoint path:\", ckpt_path)\n\n# メモリを開放\ndel X_train, X_val, y_train, y_val, w_train, w_val, dm\ngc.collect()\n\n\n#########################\n# 9) Validation でのスコア確認\n#########################\n# 例: 再度 validation.parquet を全読み込みして評価\ndf_val_pl = pls.scan_parquet(\"/kaggle/input/js24-preprocessing-create-lags/validation.parquet/\").collect()\ndf_val = df_val_pl.to_pandas()\nX_val_2 = df_val[CONFIG.feature_cols].fillna(method=\"ffill\").fillna(0)\ny_val_2 = df_val[CONFIG.target_col]\nw_val_2 = df_val[\"weight\"]\ndel df_val_pl, df_val\ngc.collect()\n\n# XGBoost 予測\ny_pred_xgb = xgb_model.predict(X_val_2)\n# CatBoost 予測\ny_pred_cat = cat_model.predict(X_val_2)\n\n# NN 予測\ndevice = \"cuda\" if torch.cuda.is_available() else \"cpu\"\nbest_nn = NNModel.load_from_checkpoint(ckpt_path).to(device)\nbest_nn.eval()\n\nX_val_torch = torch.FloatTensor(X_val_2.values).to(device)\nwith torch.no_grad():\n    y_pred_nn = best_nn(X_val_torch).cpu().numpy()\n\n# 単体スコア\nr2_xgb = r2_score(y_val_2, y_pred_xgb, sample_weight=w_val_2)\nr2_cat = r2_score(y_val_2, y_pred_cat, sample_weight=w_val_2)\nr2_nn  = r2_score(y_val_2, y_pred_nn,  sample_weight=w_val_2)\nprint(f\"[Validation R^2] XGB={r2_xgb:.6f}, CAT={r2_cat:.6f}, NN={r2_nn:.6f}\")\n\n# アンサンブル例 (XGB:0.5, NN:0.3, CAT:0.2)\nensemble_ratio = [0.5, 0.3, 0.2]\ny_pred_ens = ensemble_ratio[0]*y_pred_xgb + ensemble_ratio[1]*y_pred_nn + ensemble_ratio[2]*y_pred_cat\nr2_ens = r2_score(y_val_2, y_pred_ens, sample_weight=w_val_2)\nprint(f\"Ensemble R^2 = {r2_ens:.6f}\")\n\n\n#########################\n# 10) 学習済みモデルの保存\n#########################\nwith open(CONFIG.model_save_path_xgb, \"wb\") as fp:\n    pickle.dump(xgb_model, fp)\n\nwith open(CONFIG.model_save_path_cat, \"wb\") as fp:\n    pickle.dump(cat_model, fp)\n\nprint(\"All models saved!\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-27T11:08:22.756805Z","iopub.execute_input":"2024-12-27T11:08:22.757088Z","iopub.status.idle":"2024-12-27T11:36:50.626266Z","shell.execute_reply.started":"2024-12-27T11:08:22.757064Z","shell.execute_reply":"2024-12-27T11:36:50.625279Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"import os\nimport gc\nimport numpy as np\nimport pandas as pd\nimport polars as pl\nimport pickle\n\nimport torch\nimport torch.nn as nn\nfrom tqdm.auto import tqdm\n\nimport kaggle_evaluation.jane_street_inference_server\n\n# ===================================================\n# 1) CONFIGクラス & トレーニング時と同じ設定を復元\n# ===================================================\nclass CONFIG:\n    seed = 42\n    target_col = \"responder_6\"\n    feature_cols = [f\"feature_{idx:02d}\" for idx in range(79)] + [f\"responder_{idx}_lag_1\" for idx in range(9)]\n    \n    # 学習済みモデルを読み込むパス (Training Part で保存したパスと合わせる)\n    model_save_path_xgb = \"trained_xgb.pkl\"\n    model_save_path_cat = \"trained_cat.pkl\"\n    model_save_folder_nn = \"trained_nn\"\n\nxgb_model = None\ncat_model = None\nnn_model = None\n\n# ===================================================\n# 2) まず 学習済みモデル を読み込む\n# ===================================================\nwith open(CONFIG.model_save_path_xgb, \"rb\") as fp:\n    xgb_model = pickle.load(fp)\n\nwith open(CONFIG.model_save_path_cat, \"rb\") as fp:\n    cat_model = pickle.load(fp)\n\n# NN の読み込み (単一モデルの例)\n# K-Fold で5モデル等を保存している場合は、複数読み込んで平均する形に変更してください。\ncheckpoint_path = f\"{CONFIG.model_save_folder_nn}/nn_model.ckpt\"  # ModelCheckpoint で保存されたファイル名\n# ただし実際のファイル名が異なる場合は調整してください (例: \"nn_model-v0.ckpt\" 等)\nclass NNModel(nn.Module):\n    def __init__(self):\n        super().__init__()\n        # 学習時と同じネットワーク定義\n        # 簡易的に書いていますが、本来は学習スクリプトからクラスごとコピー推奨です。\n        self.seq = nn.Sequential(\n            nn.BatchNorm1d(len(CONFIG.feature_cols)),\n            nn.SiLU(),\n            nn.Dropout(0.2),\n            nn.Linear(len(CONFIG.feature_cols), 512),\n            nn.BatchNorm1d(512),\n            nn.SiLU(),\n            nn.Dropout(0.2),\n            nn.Linear(512, 512),\n            nn.BatchNorm1d(512),\n            nn.SiLU(),\n            nn.Dropout(0.2),\n            nn.Linear(512, 256),\n            nn.BatchNorm1d(256),\n            nn.SiLU(),\n            nn.Dropout(0.2),\n            nn.Linear(256, 1),\n            nn.Tanh(),\n        )\n\n    def forward(self, x):\n        return 5 * self.seq(x).squeeze(-1)\n\n# 推論用にインスタンス化\nnn_model = NNModel()\n# state_dict をロード\nnn_state_dict = torch.load(checkpoint_path, map_location=\"cpu\")[\"state_dict\"]\n# Pytorch-Lightning 形式のキーを nn.Module 形式に変換 (Lightning の場合キーに \"model.\" などが付くことがあるため)\nnew_state_dict = {}\nfor k, v in nn_state_dict.items():\n    # k が \"model.\" や \"seq.\" などで始まるなら取り除き\n    if k.startswith(\"model.\"):\n        new_k = k.replace(\"model.\", \"\")\n    else:\n        new_k = k\n    new_state_dict[new_k] = v\n\nnn_model.load_state_dict(new_state_dict, strict=False)\nnn_model.eval()\nnn_model = nn_model.to(\"cuda\" if torch.cuda.is_available() else \"cpu\")\n\n# ===================================================\n# 3) 予測用関数 predict() の定義\n# ===================================================\nlags_ : pl.DataFrame | None = None\n\ndef predict(test: pl.DataFrame, lags: pl.DataFrame | None) -> pl.DataFrame | pd.DataFrame:\n    \"\"\"\n    コンペ形式に則った推論関数: \n     - 入力: pl.DataFrame の test, lags\n     - 出力: 予測列 'responder_6' を含む DataFrame\n    \"\"\"\n    global lags_\n\n    # lags の更新 (1日目の場合にのみ非 None が入る想定)\n    if lags is not None:\n        lags_ = lags\n    \n    # 出力フォーマット初期化\n    predictions = test.select(\n        'row_id',\n        pl.lit(0.0).alias('responder_6'),\n    )\n\n    # lags テーブルから group_by([\"date_id\", \"symbol_id\"]) で最後の行を引き継ぐ (＝ラグ特徴量)\n    lags_df = lags_.clone().group_by([\"date_id\", \"symbol_id\"], maintain_order=True).last()\n    test_merged = test.join(lags_df, on=[\"date_id\", \"symbol_id\"], how=\"left\")\n\n    # 必要な特徴量を pandas に変換して欠損値を補完\n    X_test = test_merged[CONFIG.feature_cols].to_pandas()\n    X_test = X_test.fillna(method='ffill').fillna(0)\n\n    # 3つのモデル (XGB, CatBoost, NN) で予測\n    # -------------------------------------------------\n    # XGB\n    preds_xgb = xgb_model.predict(X_test)\n\n    # CatBoost\n    preds_cat = cat_model.predict(X_test)\n\n    # NN\n    X_torch = torch.FloatTensor(X_test.values).to(\"cuda\" if torch.cuda.is_available() else \"cpu\")\n    preds_nn = np.zeros(X_test.shape[0], dtype=np.float32)\n    with torch.no_grad():\n        y_hat_nn = nn_model(X_torch).cpu().numpy()\n        preds_nn += y_hat_nn\n\n    # Weighted Ensemble (例: XGB 0.5, NN 0.3, CAT 0.2)\n    ensemble_ratio = [0.5, 0.3, 0.2]\n    preds = (ensemble_ratio[0] * preds_xgb\n             + ensemble_ratio[1] * preds_nn\n             + ensemble_ratio[2] * preds_cat)\n\n    # 最終出力を clamping\n    preds = np.clip(preds, a_min=-5, a_max=5)\n\n    # 出力を polars DataFrame に反映\n    predictions = test_merged.select('row_id').with_columns(\n        pl.Series(name='responder_6', values=preds, dtype=pl.Float64)\n    )\n\n    # 出力データ検証\n    assert isinstance(predictions, (pl.DataFrame, pd.DataFrame)), \\\n        \"predictions must be a pl.DataFrame or pd.DataFrame.\"\n    assert list(predictions.columns) == ['row_id', 'responder_6'], \\\n        \"Output columns must be ['row_id', 'responder_6']\"\n    assert len(predictions) == len(test_merged), \\\n        \"The number of prediction rows must match test rows.\"\n\n    return predictions","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-27T11:36:50.630794Z","iopub.execute_input":"2024-12-27T11:36:50.631009Z","iopub.status.idle":"2024-12-27T11:36:50.824489Z","shell.execute_reply.started":"2024-12-27T11:36:50.630989Z","shell.execute_reply":"2024-12-27T11:36:50.823772Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"inference_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-27T11:40:53.076725Z","iopub.execute_input":"2024-12-27T11:40:53.077033Z","iopub.status.idle":"2024-12-27T11:40:53.258217Z","shell.execute_reply.started":"2024-12-27T11:40:53.077012Z","shell.execute_reply":"2024-12-27T11:40:53.257453Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null}]}