#!/usr/bin/env python3 """第1段階向け LightGBM 学習。 ウォークフォワードCV・ハイパラ探索・アンサンブル・リーク検査をオプションで実施する。 """ from __future__ import annotations import argparse import json import os import sys from typing import Any sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) from ai.feature_engineering import ( # noqa: E402 build_training_frame, get_feature_columns, get_target_column_names, validate_training_frame_columns, ) from collectors.db_util import connect # noqa: E402 def _classification_metrics(y_true: list[int], y_prob: list[float], threshold: float = 0.5) -> dict[str, float]: tp = 0 fp = 0 fn = 0 tn = 0 for label, prob in zip(y_true, y_prob): pred = 1 if prob >= threshold else 0 if label == 1 and pred == 1: tp += 1 elif label == 0 and pred == 1: fp += 1 elif label == 1 and pred == 0: fn += 1 else: tn += 1 total = max(1, len(y_true)) precision = tp / max(1, tp + fp) recall = tp / max(1, tp + fn) accuracy = (tp + tn) / total return { "accuracy": round(float(accuracy), 4), "precision": round(float(precision), 4), "recall": round(float(recall), 4), "tp": tp, "fp": fp, "fn": fn, "tn": tn, } def _average_precision_binary(y_true: list[int], y_prob: list[float]) -> float | None: try: from sklearn.metrics import average_precision_score return float(average_precision_score(y_true, y_prob)) except Exception: return None def _eval_fold(model, x_valid: Any, y_valid: list[int]) -> tuple[float, dict[str, float], float | None]: import numpy as np prob = model.predict_proba(x_valid)[:, 1].tolist() m = _classification_metrics(y_valid, prob) pr_auc = _average_precision_binary(y_valid, prob) ll = -float( np.mean( [ y * np.log(max(p, 1e-15)) + (1 - y) * np.log(max(1 - p, 1e-15)) for y, p in zip(y_valid, prob) ], ), ) return ll, m, pr_auc def _walk_forward_indices( df, n_folds: int, min_train_rows: int, ) -> tuple[list[tuple[list[int], list[int]]], dict[str, Any]]: """race_date 昇順。戻り値の辞書に要求折数・実効折数・スキップ理由を含める。""" import numpy as np import pandas as pd df = df.sort_values(["race_date", "race_id", "horse_id"]).reset_index(drop=True) uniq = sorted(df["race_date"].dropna().unique()) if len(uniq) < n_folds + 2: raise ValueError( f"ウォークフォワードに日付が不足です(ユニーク日数={len(uniq)}、要 n_folds+2 程度)", ) chunks = np.array_split(np.array(uniq, dtype=object), n_folds) splits: list[tuple[list[int], list[int]]] = [] skipped: list[dict[str, Any]] = [] for chi, chunk in enumerate(chunks): if chunk.size == 0: skipped.append({"chunk_index": chi, "reason": "empty_date_chunk", "date_sample": []}) continue d0 = pd.Timestamp(chunk.min()) train_mask = pd.to_datetime(df["race_date"]) < d0 valid_mask = df["race_date"].isin(chunk) train_idx = df.index[train_mask].tolist() valid_idx = df.index[valid_mask].tolist() tr_n = len(train_idx) va_n = len(valid_idx) if tr_n < min_train_rows or va_n < 50: ds = sorted(chunk.tolist()) skipped.append( { "chunk_index": chi, "reason": "train_rows_lt_min_or_valid_lt_50", "train_rows": tr_n, "valid_rows": va_n, "date_sample": [str(x) for x in ds[:8]], }, ) continue splits.append((train_idx, valid_idx)) diagnostics: dict[str, Any] = { "requested_cv_folds": int(n_folds), "unique_race_dates": len(uniq), "date_chunks_total": int(len(chunks)), "evaluated_folds": len(splits), "skipped_chunks": skipped, } if not splits: raise ValueError("ウォークフォワード分割が成立しませんでした(min_train_rows を下げるかデータを増やしてください)") return splits, diagnostics def run_walk_forward_cv( df: Any, feature_cols: list[str], target: str, y: list[int], n_folds: int, min_train_rows: int, lgb_params: dict[str, Any], ) -> dict[str, Any]: import statistics import lightgbm as lgb splits, wf_diag = _walk_forward_indices(df, n_folds=n_folds, min_train_rows=min_train_rows) fold_logs: list[dict[str, Any]] = [] loglosses: list[float] = [] pr_aucs: list[float] = [] for fi, (tr_idx, va_idx) in enumerate(splits): x_tr = df.iloc[tr_idx][feature_cols] y_tr = [y[i] for i in tr_idx] x_va = df.iloc[va_idx][feature_cols] y_va = [y[i] for i in va_idx] model = lgb.LGBMClassifier(**lgb_params) model.fit( x_tr, y_tr, eval_set=[(x_va, y_va)], eval_metric="binary_logloss", callbacks=[lgb.early_stopping(30, verbose=False)], ) ll, m, pr_auc = _eval_fold(model, x_va, y_va) loglosses.append(ll) fl: dict[str, Any] = { "fold": fi + 1, "train_rows": len(tr_idx), "valid_rows": len(va_idx), "valid_logloss": round(ll, 5), "valid_metrics": m, "best_iteration": int(model.best_iteration_ or lgb_params.get("n_estimators", 500)), } if pr_auc is not None: fl["valid_pr_auc"] = round(pr_auc, 5) pr_aucs.append(pr_auc) fold_logs.append(fl) mean_ll = float(sum(loglosses) / max(1, len(loglosses))) std_ll = round(statistics.pstdev(loglosses), 5) if len(loglosses) > 1 else 0.0 mean_pr = round(float(sum(pr_aucs) / len(pr_aucs)), 5) if pr_aucs else None out: dict[str, Any] = { "walk_forward_split": wf_diag, "folds": fold_logs, "mean_valid_logloss": round(mean_ll, 5), "std_valid_logloss": std_ll, } if mean_pr is not None: out["mean_valid_pr_auc"] = mean_pr return out def _median_best_iteration(cv_report: dict[str, Any] | None) -> int | None: if not cv_report: return None folds = cv_report.get("folds") or [] its = [int(f["best_iteration"]) for f in folds if f.get("best_iteration") is not None] if not its: return None import statistics return int(statistics.median(its)) def _final_n_estimators_for_deploy( *, cv_report: dict[str, Any] | None, ceiling: int, ignore_cv_iteration: bool, ) -> tuple[int, dict[str, Any]]: """CV の best_iteration 中央値から本番学習の木数を決める(検証と独立した全データ fit 用)。""" base_note: dict[str, Any] = { "mode": "ceiling_only", "n_estimators_used": ceiling, "ceiling": ceiling, } if ignore_cv_iteration or not cv_report: return ceiling, base_note med = _median_best_iteration(cv_report) if med is None: return ceiling, base_note # 全データ再学習では検証よりやや余裕を見る(過小の木数回避) buffered = int(med + max(15, round(med * 0.08))) n_used = min(ceiling, max(50, buffered)) note = { "mode": "from_cv_median_best_iteration", "median_best_iteration": med, "buffered_iteration": buffered, "ceiling": ceiling, "n_estimators_used": n_used, } return n_used, note def _build_lgb_params( *, n_estimators: int, learning_rate: float, num_leaves: int, min_child_samples: int, random_state: int, ) -> dict[str, Any]: return { "objective": "binary", "random_state": random_state, "n_estimators": n_estimators, "learning_rate": learning_rate, "num_leaves": num_leaves, "min_child_samples": min_child_samples, "class_weight": "balanced", } def main() -> None: targets = get_target_column_names() p = argparse.ArgumentParser(description="第1段階向け LightGBM 学習(WF-CV・チューニング・アンサンブル対応)") p.add_argument("--lookback-runs", type=int, default=10) p.add_argument("--limit-rows", type=int, default=None) p.add_argument("--output-csv", type=str, default="ai/models/stage1_training.csv", help="生成する学習CSV") p.add_argument("--target", type=str, default="label_top3", choices=targets) p.add_argument("--train-ratio", type=float, default=0.8, help="単純分割モード時の学習比率(--cv-folds 0 のとき)") p.add_argument("--cv-folds", type=int, default=0, help=">0 でウォークフォワードCV(日付ベース)。0 で単純時系列1分割") p.add_argument("--min-train-rows-cv", type=int, default=5000, help="WF-CV の各フォールドの train 最小行数") p.add_argument("--tune", action="store_true", help="ハイパラを狭いグリッドで探索(WF-CV の平均 logloss 最小)") p.add_argument("--tune-trials", type=int, default=0, help="0=全グリッド、>0 でランダムに試行回数を上限") p.add_argument("--ensemble", type=int, default=1, help="最終モデルを複数本学習し平均アンサンブル(2 以上)") p.add_argument("--skip-leak-check", action="store_true", help="特徴列のリークパターンチェックをスキップ") p.add_argument("--num-leaves", type=int, default=31) p.add_argument("--learning-rate", type=float, default=0.05) p.add_argument("--n-estimators", type=int, default=500) p.add_argument("--min-child-samples", type=int, default=20) p.add_argument("--output-model", type=str, default="ai/models/stage1_lgbm.txt") p.add_argument("--output-meta", type=str, default="ai/models/stage1_lgbm.meta.json") p.add_argument( "--ignore-cv-best-iteration", action="store_true", help="WF-CV があっても最終学習の木数を --n-estimators 固定にする(既定は CV の best_iteration 中央値から決定)", ) args = p.parse_args() if args.lookback_runs <= 0: print("--lookback-runs は 1 以上で指定してください", file=sys.stderr) sys.exit(1) if not (0.5 <= args.train_ratio < 1.0): print("--train-ratio は 0.5 以上 1.0 未満で指定してください", file=sys.stderr) sys.exit(1) if args.ensemble < 1: print("--ensemble は 1 以上で指定してください", file=sys.stderr) sys.exit(1) try: import lightgbm as lgb except ImportError: print("lightgbm が未インストールです。`python3.11 -m pip install lightgbm` を実行してください。", file=sys.stderr) sys.exit(1) conn = connect() try: df = build_training_frame( conn=conn, lookback_runs=args.lookback_runs, limit_rows=args.limit_rows, ) finally: conn.close() if df.empty: print("学習データが0件でした。race_entries / race_results を確認してください。", file=sys.stderr) sys.exit(1) feature_cols = get_feature_columns() if not args.skip_leak_check: validate_training_frame_columns(feature_cols) missing = [c for c in feature_cols if c not in df.columns] if missing: print(f"特徴量列が不足しています: {missing}", file=sys.stderr) sys.exit(1) out_csv = os.path.abspath(args.output_csv) os.makedirs(os.path.dirname(out_csv), exist_ok=True) df.to_csv(out_csv, index=False) df = df.sort_values(["race_date", "race_id", "horse_id"], ascending=True).reset_index(drop=True) target = args.target if target not in df.columns: print(f"ターゲット列 `{target}` がありません。", file=sys.stderr) sys.exit(1) y = df[target].astype(int).tolist() if len(set(y)) < 2: print(f"ターゲット `{target}` が単一クラスのため学習できません。", file=sys.stderr) sys.exit(1) meta_out: dict[str, Any] = { "target": target, "rows_total": int(len(df)), "feature_columns": feature_cols, "lookback_runs": args.lookback_runs, "cv_config": { "cv_folds": args.cv_folds, "tune": args.tune, "ensemble_n": args.ensemble, }, } # ハイパラ探索(狭いグリッド) best_params = { "num_leaves": args.num_leaves, "learning_rate": args.learning_rate, "min_child_samples": args.min_child_samples, "n_estimators": args.n_estimators, } cv_report: dict[str, Any] | None = None if args.tune and args.cv_folds > 0: grid = [] for nl in (23, 31, 47): for lr in (0.035, 0.05, 0.07): for mcs in (15, 25, 40): grid.append({"num_leaves": nl, "learning_rate": lr, "min_child_samples": mcs}) import random if args.tune_trials > 0 and args.tune_trials < len(grid): random.seed(42) grid = random.sample(grid, args.tune_trials) best_ll = float("inf") for g in grid: lp = _build_lgb_params( n_estimators=min(args.n_estimators, 400), learning_rate=g["learning_rate"], num_leaves=g["num_leaves"], min_child_samples=g["min_child_samples"], random_state=42, ) cv_res = run_walk_forward_cv( df, feature_cols, target, y, n_folds=args.cv_folds, min_train_rows=args.min_train_rows_cv, lgb_params=lp, ) mll = cv_res["mean_valid_logloss"] if mll < best_ll: best_ll = mll best_params = {**g, "n_estimators": args.n_estimators} cv_report = cv_res meta_out["tuning"] = {"best_grid_params": best_params, "best_mean_valid_logloss": best_ll} meta_out["walk_forward_cv"] = cv_report print(f"[tune] best mean_valid_logloss={best_ll} params={best_params}") elif args.tune and args.cv_folds <= 0: print("[tune] --cv-folds を 1 以上にしてください(WF-CV と併用)", file=sys.stderr) sys.exit(1) elif args.cv_folds > 0 and not args.tune: lp = _build_lgb_params( n_estimators=best_params["n_estimators"], learning_rate=best_params["learning_rate"], num_leaves=best_params["num_leaves"], min_child_samples=best_params["min_child_samples"], random_state=42, ) cv_report = run_walk_forward_cv( df, feature_cols, target, y, n_folds=args.cv_folds, min_train_rows=args.min_train_rows_cv, lgb_params=lp, ) meta_out["walk_forward_cv"] = cv_report print(f"[cv] mean_valid_logloss={cv_report['mean_valid_logloss']} ± {cv_report.get('std_valid_logloss', 0)}") n_estimators_ceiling = int(best_params["n_estimators"]) n_final, deploy_note = _final_n_estimators_for_deploy( cv_report=cv_report, ceiling=n_estimators_ceiling, ignore_cv_iteration=args.ignore_cv_best_iteration, ) meta_out["final_training"] = { **deploy_note, "n_estimators_ceiling": n_estimators_ceiling, "ignore_cv_best_iteration": bool(args.ignore_cv_best_iteration), } print(f"[deploy] final n_estimators={n_final} ({deploy_note.get('mode')})") # 最終学習: 全データ(運用デプロイ用) x_all = df[feature_cols] out_model_base = os.path.abspath(args.output_model) out_meta = os.path.abspath(args.output_meta) os.makedirs(os.path.dirname(out_model_base), exist_ok=True) meta_dir = os.path.dirname(out_meta) ensemble_entries: list[dict[str, Any]] = [] last_model: Any | None = None for ei in range(args.ensemble): rs = 42 + ei lp = _build_lgb_params( n_estimators=n_final, learning_rate=best_params["learning_rate"], num_leaves=best_params["num_leaves"], min_child_samples=best_params["min_child_samples"], random_state=rs, ) model = lgb.LGBMClassifier(**lp) # 本番用は全行で学習(検証データを eval に混ぜない=リーク防止) model.fit(x_all, y) last_model = model stem, ext = os.path.splitext(out_model_base) path_e = f"{stem}_e{ei + 1}{ext}" if args.ensemble > 1 else out_model_base model.booster_.save_model(path_e) w = 1.0 / args.ensemble rel = os.path.relpath(path_e, start=meta_dir) ensemble_entries.append({"path": rel, "weight": round(w, 4)}) if last_model is not None: meta_out["feature_importance_gain"] = { k: float(v) for k, v in zip(feature_cols, last_model.booster_.feature_importance(importance_type="gain")) } else: meta_out["feature_importance_gain"] = {} meta_out["ensemble"] = {"models": ensemble_entries} meta_out["training_target_note"] = ( "label_win=1着、label_top3=複勝圏、label_top5=5着以内。" "運用の主指標は通常 label_top3 または label_win。" ) with open(out_meta, "w", encoding="utf-8") as fw: json.dump(meta_out, fw, ensure_ascii=False, indent=2) print(f"training rows: {len(df)}") print(f"feature columns: {', '.join(feature_cols)}") print(f"csv: {out_csv}") print(f"meta: {out_meta}") for e in ensemble_entries: print(f"model: {e['path']} (w={e['weight']})") if __name__ == "__main__": main()