{"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":10277864,"sourceType":"datasetVersion","datasetId":6359705},{"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\nimport numpy as np\nimport pandas as pd\nimport polars as pls\nimport torch\nimport torch.nn as nn\nimport torch.nn.functional as F\n\nimport pytorch_lightning as pl\nfrom pytorch_lightning.callbacks import EarlyStopping, ModelCheckpoint\n\nfrom sklearn.model_selection import KFold\nfrom sklearn.linear_model import Ridge\nfrom sklearn.preprocessing import StandardScaler\nfrom sklearn.pipeline import Pipeline\nfrom sklearn.metrics import r2_score\n\nfrom xgboost import XGBRegressor\nfrom catboost import CatBoostRegressor\nimport lightgbm as lgb\n\nfrom torch.utils.data import Dataset, DataLoader\n\n# コンペ指定のインフェレンスサーバー\nimport kaggle_evaluation.jane_street_inference_server\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    n_folds = 5\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    # モデルの保存先\n    model_dir_nn = \"nn_folds\"\n    model_dir_gbdt = \"gbdt_folds\"\n    model_path_stacking = \"stacking_model.pkl\"\n\n\n#########################\n# 3) データ読み込み\n#########################\nvalid_pl = pls.scan_parquet(\"/kaggle/input/js24-preprocessing-create-lags/validation.parquet/\").collect()\ndf = valid_pl.to_pandas()\nprint(\"df.shape =\", df.shape)\n\nX_all = df[CONFIG.feature_cols].fillna(method=\"ffill\").fillna(0)\ny_all = df[CONFIG.target_col].values\nw_all = df[\"weight\"].values\ndel valid_pl, df\ngc.collect()\n\n\n#########################\n# 4) 評価指標 (Weighted R^2)\n#########################\ndef weighted_r2_score(y_true, y_pred, weight):\n    \"\"\"\n    Weighted R^2 を計算\n    \"\"\"\n    numerator = np.average((y_pred - y_true)**2, weights=weight)\n    denominator = np.average(y_true**2, weights=weight) + 1e-38\n    return 1.0 - numerator / denominator\n\n\n#########################\n# 5) NNモデル定義\n#########################\nclass NNModel(pl.LightningModule):\n    def __init__(\n        self,\n        input_dim,\n        hidden_dims=[512, 256],\n        dropouts=[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, hdim in enumerate(hidden_dims):\n            layers.append(nn.BatchNorm1d(in_dim))\n            layers.append(nn.SiLU())\n            layers.append(nn.Dropout(dropouts[i]))\n            layers.append(nn.Linear(in_dim, hdim))\n            in_dim = hdim\n        \n        # 出力層 + Tanh => [-1,1], ×5 => [-5,5]\n        layers.append(nn.Linear(in_dim, 1))\n        layers.append(nn.Tanh())\n        self.model = nn.Sequential(*layers)\n\n        self.lr = lr\n        self.weight_decay = weight_decay\n        self.validation_step_outputs = []\n\n    def forward(self, x):\n        return 5.0 * self.model(x).squeeze(-1)\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)\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)\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        # validationデータに対する R^2 を算出\n        y_pred_all = []\n        y_true_all = []\n        w_all = []\n        for (yp, yt, wt) in self.validation_step_outputs:\n            y_pred_all.append(yp.cpu().numpy())\n            y_true_all.append(yt.cpu().numpy())\n            w_all.append(wt.cpu().numpy())\n        y_pred_all = np.concatenate(y_pred_all, axis=0).squeeze()\n        y_true_all = np.concatenate(y_true_all, axis=0).squeeze()\n        w_all = np.concatenate(w_all, axis=0).squeeze()\n\n        val_r2 = weighted_r2_score(y_true_all, y_pred_all, w_all)\n        self.log(\"val_r2\", 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(\n            optimizer, mode=\"min\", factor=0.5, patience=5, verbose=True\n        )\n        return {\n            \"optimizer\": optimizer,\n            \"lr_scheduler\": {\n                \"scheduler\": scheduler,\n                \"monitor\": \"val_loss\",\n            }\n        }\n\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[idx]),\n            torch.FloatTensor([self.y[idx]]),\n            torch.FloatTensor([self.w[idx]])\n        )\n\nclass MarketDataModule(pl.LightningDataModule):\n    def __init__(self, X, y, w, X_val, y_val, w_val, batch_size=4096):\n        super().__init__()\n        self.X = X\n        self.y = y\n        self.w = w\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 train_dataloader(self):\n        ds = MarketDataset(self.X, self.y, self.w)\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# 6) K-Fold で GBDT & NN を学習\n#########################\nos.makedirs(CONFIG.model_dir_nn, exist_ok=True)\nos.makedirs(CONFIG.model_dir_gbdt, exist_ok=True)\n\nkf = KFold(n_splits=CONFIG.n_folds, shuffle=True, random_state=CONFIG.seed)\n\noof_preds_xgb = np.zeros(len(X_all), dtype=float)\noof_preds_cat = np.zeros(len(X_all), dtype=float)\noof_preds_lgb = np.zeros(len(X_all), dtype=float)\noof_preds_nn  = np.zeros(len(X_all), dtype=float)\n\ndevice = \"cuda\" if torch.cuda.is_available() else \"cpu\"\n\nfor fold, (tr_idx, va_idx) in enumerate(kf.split(X_all)):\n    print(f\"========== FOLD {fold} ==========\")\n    X_tr = X_all.iloc[tr_idx].values\n    y_tr = y_all[tr_idx]\n    w_tr = w_all[tr_idx]\n\n    X_va = X_all.iloc[va_idx].values\n    y_va = y_all[va_idx]\n    w_va = w_all[va_idx]\n\n    # ========== (A) XGBoost ==========\n    xgb_model = XGBRegressor(\n        random_state=CONFIG.seed,\n        n_estimators=400,\n        max_depth=7,\n        learning_rate=0.02,\n        subsample=0.8,\n        colsample_bytree=0.8,\n        tree_method=\"gpu_hist\",  # GPU\n    )\n    xgb_model.fit(\n        X_tr, y_tr,\n        sample_weight=w_tr,\n        eval_set=[(X_va, y_va)],\n        eval_metric=\"rmse\",\n        sample_weight_eval_set=[w_va],\n        early_stopping_rounds=50,\n        verbose=False\n    )\n    preds_va_xgb = xgb_model.predict(X_va)\n    oof_preds_xgb[va_idx] = preds_va_xgb\n    with open(f\"{CONFIG.model_dir_gbdt}/xgb_fold{fold}.pkl\", \"wb\") as fp:\n        pickle.dump(xgb_model, fp)\n\n    # ========== (B) CatBoost ==========\n    cat_model = CatBoostRegressor(\n        iterations=400,\n        depth=7,\n        learning_rate=0.02,\n        random_seed=CONFIG.seed,\n        task_type=\"GPU\",\n        verbose=False\n    )\n    cat_model.fit(\n        X_tr, y_tr, sample_weight=w_tr,\n        eval_set=(X_va, y_va),\n        use_best_model=True\n    )\n    preds_va_cat = cat_model.predict(X_va)\n    oof_preds_cat[va_idx] = preds_va_cat\n    with open(f\"{CONFIG.model_dir_gbdt}/cat_fold{fold}.pkl\", \"wb\") as fp:\n        pickle.dump(cat_model, fp)\n\n    # ========== (C) LightGBM ==========\n    lgb_model = lgb.LGBMRegressor(\n        n_estimators=1000,\n        learning_rate=0.02,\n        max_depth=7,\n        random_state=CONFIG.seed,\n        subsample=0.8,\n        colsample_bytree=0.8,\n        device=\"gpu\",\n        verbosity=-1\n    )\n    lgb_model.fit(\n        X_tr, y_tr,\n        sample_weight=w_tr,\n        eval_set=[(X_va, y_va)],\n        eval_sample_weight=[w_va],\n        eval_metric=\"rmse\",\n        callbacks=[\n            lgb.early_stopping(50, verbose=False),\n        ],\n        # verbose=... => NG\n    )\n    preds_va_lgb = lgb_model.predict(X_va)\n    oof_preds_lgb[va_idx] = preds_va_lgb\n    with open(f\"{CONFIG.model_dir_gbdt}/lgb_fold{fold}.pkl\", \"wb\") as fp:\n        pickle.dump(lgb_model, fp)\n\n    # ========== (D) NN ==========\n    scaler = StandardScaler()\n    X_tr_nn = scaler.fit_transform(X_tr)\n    X_va_nn = scaler.transform(X_va)\n\n    dm = MarketDataModule(X_tr_nn, y_tr, w_tr, X_va_nn, y_va, w_va, batch_size=4096)\n    nn_model = NNModel(input_dim=X_tr_nn.shape[1])\n    \n    trainer = pl.Trainer(\n        max_epochs=10,\n        accelerator=\"gpu\" if torch.cuda.is_available() else \"cpu\",\n        precision=16,\n        callbacks=[\n            EarlyStopping(monitor=\"val_loss\", mode=\"min\", patience=3),\n            ModelCheckpoint(\n                dirpath=CONFIG.model_dir_nn,\n                filename=f\"nn_fold{fold}\",\n                monitor=\"val_loss\",\n                mode=\"min\",\n                save_top_k=1\n            )\n        ]\n    )\n    trainer.fit(nn_model, dm)\n    best_ckpt = trainer.checkpoint_callback.best_model_path\n    best_nn_model = NNModel.load_from_checkpoint(best_ckpt).eval().to(device)\n    X_va_torch = torch.FloatTensor(X_va_nn).to(device)\n    with torch.no_grad():\n        preds_va_nn = best_nn_model(X_va_torch).cpu().numpy()\n    oof_preds_nn[va_idx] = preds_va_nn.flatten()\n\n# ---------- OOF スコア ----------\ndef print_oof_score(name, preds):\n    r2 = weighted_r2_score(y_all, preds, w_all)\n    print(f\"  {name} = {r2:.6f}\")\n\nprint(\"OOF Weighted R^2:\")\nprint_oof_score(\"XGB\", oof_preds_xgb)\nprint_oof_score(\"CAT\", oof_preds_cat)\nprint_oof_score(\"LGB\", oof_preds_lgb)\nprint_oof_score(\"NN \", oof_preds_nn)\n\n\n#########################\n# 7) Stacking の学習\n#########################\nstack_X = np.column_stack([oof_preds_xgb, oof_preds_cat, oof_preds_lgb, oof_preds_nn])\nstack_y = y_all\nstack_w = w_all\n\n# Pipeline でRidge + StandardScaler\nstack_model = Pipeline([\n    (\"scaler\", StandardScaler()),\n    (\"ridge\", Ridge(alpha=1.0)),\n])\n\n# Pipeline.fit() は sample_weight=... を受け付けないので \"ridge__sample_weight=...\" とする\nstack_model.fit(\n    stack_X, stack_y,\n    ridge__sample_weight=stack_w\n)\n\nstack_oof_pred = stack_model.predict(stack_X)\nr2_stack = weighted_r2_score(stack_y, stack_oof_pred, stack_w)\nprint(f\"OOF Weighted R^2 (Stacking) = {r2_stack:.6f}\")\n\nwith open(CONFIG.model_path_stacking, \"wb\") as fp:\n    pickle.dump(stack_model, fp)\n\n\n#########################\n# 8) 推論用 関数\n#########################\nlags_: pls.DataFrame | None = None\n\ndef predict(test: pls.DataFrame, lags: pls.DataFrame | None) -> pls.DataFrame | pd.DataFrame:\n    global lags_\n    if lags is not None:\n        lags_ = lags\n\n    # 出力データフレーム\n    preds_df = test.select(\"row_id\", pls.lit(0.0).alias(CONFIG.target_col))\n\n    # ラグ結合\n    lags_last = lags_.clone().group_by([\"date_id\",\"symbol_id\"], maintain_order=True).last()\n    test_merged = test.join(lags_last, on=[\"date_id\",\"symbol_id\"], how=\"left\")\n\n    # 特徴量をNumpy化\n    X_test = test_merged[CONFIG.feature_cols].to_pandas().fillna(method=\"ffill\").fillna(0).values\n\n    # Foldごとに予測 → 平均\n    # (A) XGB\n    xgb_preds_fold = []\n    for fold in range(CONFIG.n_folds):\n        with open(f\"{CONFIG.model_dir_gbdt}/xgb_fold{fold}.pkl\",\"rb\") as fp:\n            xgb_model = pickle.load(fp)\n        xgb_preds_fold.append(xgb_model.predict(X_test))\n    xgb_preds = np.mean(xgb_preds_fold, axis=0)\n\n    # (B) Cat\n    cat_preds_fold = []\n    for fold in range(CONFIG.n_folds):\n        with open(f\"{CONFIG.model_dir_gbdt}/cat_fold{fold}.pkl\",\"rb\") as fp:\n            cat_model = pickle.load(fp)\n        cat_preds_fold.append(cat_model.predict(X_test))\n    cat_preds = np.mean(cat_preds_fold, axis=0)\n\n    # (C) LGB\n    lgb_preds_fold = []\n    for fold in range(CONFIG.n_folds):\n        with open(f\"{CONFIG.model_dir_gbdt}/lgb_fold{fold}.pkl\",\"rb\") as fp:\n            lgb_model = pickle.load(fp)\n        lgb_preds_fold.append(lgb_model.predict(X_test))\n    lgb_preds = np.mean(lgb_preds_fold, axis=0)\n\n    # (D) NN\n    nn_preds_fold = []\n    device = \"cuda\" if torch.cuda.is_available() else \"cpu\"\n    for fold in range(CONFIG.n_folds):\n        ckpt_path = f\"{CONFIG.model_dir_nn}/nn_fold{fold}.ckpt\"\n        nn_model = NNModel.load_from_checkpoint(ckpt_path)\n        nn_model.eval().to(device)\n        X_test_torch = torch.FloatTensor(X_test).to(device)\n        with torch.no_grad():\n            p_nn = nn_model(X_test_torch).cpu().numpy().flatten()\n        nn_preds_fold.append(p_nn)\n    nn_preds = np.mean(nn_preds_fold, axis=0)\n\n    # Stackingモデルで最終予測\n    stack_features = np.column_stack([xgb_preds, cat_preds, lgb_preds, nn_preds])\n    with open(CONFIG.model_path_stacking, \"rb\") as fp:\n        stack_model = pickle.load(fp)\n    final_preds = stack_model.predict(stack_features)\n\n    # clip to [-5,5]\n    final_preds = np.clip(final_preds, a_min=-5, a_max=5)\n\n    # 結果\n    preds_df = test_merged.select(\"row_id\").with_columns(\n        pls.Series(name=CONFIG.target_col, values=final_preds, dtype=pls.Float64)\n    )\n    return preds_df\n\n\n#########################\n# 9) インフェレンスサーバーを起動\n#########################\ninference_server = kaggle_evaluation.jane_street_inference_server.JSInferenceServer(predict)\n\nif os.getenv(\"KAGGLE_IS_COMPETITION_RERUN\"):\n    # 提出環境 (再実行) では serve()\n    inference_server.serve()\nelse:\n    # ローカル実行 (Notebook) では run_local_gateway()\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":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","trusted":true,"execution":{"iopub.status.busy":"2024-12-30T08:25:06.122082Z","iopub.execute_input":"2024-12-30T08:25:06.122316Z","iopub.status.idle":"2024-12-30T08:48:37.785038Z","shell.execute_reply.started":"2024-12-30T08:25:06.122295Z","shell.execute_reply":"2024-12-30T08:48:37.784204Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null}]}