diff --git a/plan_review.py b/plan_review.py new file mode 100644 index 0000000..abee21d --- /dev/null +++ b/plan_review.py @@ -0,0 +1,311 @@ +"""候选单复盘(只读):这个系统唯一的评价方式——名单级观察收益,不是回测。 + +## 为什么要有它 + +方案(docs/主观选股改进方案_2026-09-02.md 第 2.4 节)把评价单位定为"候选单"而不是因子: +不建仓、不计成本、不定仓位、不算净值,只看几份名单在后五、十、二十个交易日相对 +全池与主榜的超额和命中率,按环境与证据线分组,每周五出报告。它取代了原对比器 +score_lab(两套读数工具两种口径的坑),也承接了 09-02 手工归因的三组读数: +第一步先"对账"——用 signal_close 口径复现方案 1.8b 的四档读数(强传导 −2.04、 +观察档 +0.52),证明脚本口径与手工一致;第二步再切常规口径 next_close 出周报。 + +本脚本不提供权重参数,只提供分组与期限——防止复盘变成调参或模拟交易。 + +## 六份名单 + + 生产名单 PMS 实际拿到的候选:PMS 计划快照名册(pms_plan_snapshot.roster_json)里档位为 + 强传导、按分数降序前 N(N 取当时 PMS_PLAN_TOP_N,历史值不可还原,用现值并注明); + 快照缺失的日子退为桥按同规则重算并标注 + 候选单 候选卡判决为"候选"(候选卡上线日之前的历史日,按当日快照或重算得出, + "已启动"依赖已动成员视图,视图未建的日子候选为空并标注) + 关注单 判决为"关注" + 环节名单 当日被传导指向的环节的全体成员等权(把"定位对"与"选票对"分开) + 主榜等权 当日全部主榜 + 全池等权 基座行情快照当日全部个股(基准) + +## 口径 + + 起点价 next_close(默认,实盘买得到):T 日出计划,T+1 收盘买,收益 = Σ pct[T+2 .. T+1+h] + signal_close(对账用):收益 = Σ pct[T+1 .. T+h],与方案 1.8 系列的手工口径一致 + 期限 5 / 10 / 20 个交易日 + 指标 均收益、相对全池超额、相对主榜超额、跑赢全池比例(命中率) + 分组 档位、判决、吸筹三态(基座行情快照当日的 accum 状态)、事后环境(全池后 h 日涨跌, + 只作解释)、事前标签(当日快照 regime 段,上线后才有)、市值三分位(成交额除以换手率) + +## 跑法【桥机 155 · ~/project/akg-factor-bridge】 + + docker compose exec -T akg-factor-bridge python plan_review.py --since 2026-07-29 --horizons 5,10,20 + docker compose exec -T akg-factor-bridge python plan_review.py --since 2026-07-29 --start-price signal_close --horizons 5 # 对账 1.8b + +只读:因子表、基座视图与行情快照、PMS 计划快照全部 SELECT;只写 data/review/ 下的报告与明细。 +""" +from __future__ import annotations + +import argparse +import datetime as dt +import json +import os + +import pandas as pd + +import common +import config +import db +import plan + +HORIZONS_DEFAULT = (5, 10, 20) +MAIN_MIN = 150.0 + + +# ============================================================================ +# 数据 +# ============================================================================ +def plan_dates(since: str, until: str | None) -> list[str]: + q = "SELECT DISTINCT trade_date FROM t_factor_akg_gate WHERE trade_date >= %s" + args = [since] + if until: + q += " AND trade_date <= %s" + args.append(until) + df = db.read_mysql("factor", q + " ORDER BY trade_date", tuple(args)) + return [pd.Timestamp(x).date().isoformat() for x in df["trade_date"]] + + +def price_panel(since: str, days_after: int = 30) -> pd.DataFrame: + """基座行情快照:trade_date × 前缀码 -> 日涨幅(百分数)。""" + end = (dt.date.fromisoformat(since) + dt.timedelta(days=200)).isoformat() + df = db.read_pg( + "SELECT trade_date, code, (metrics->>'pct_change')::float AS pct, " + "metrics->'accum'->>'state' AS accum FROM mkt_daily " + "WHERE kind='stock' AND trade_date >= %s AND trade_date <= %s", (since, end)) + df["k"] = df["code"].map(common.to_prefix) + df["trade_date"] = pd.to_datetime(df["trade_date"]).dt.date.astype(str) + return df + + +def cap_bucket(day: str) -> dict[str, str]: + """市值三分位(成交额除以换手率的近似流通市值),同日分桶。读不到返回空。""" + try: + df = db.read_mysql( + "factor", "SELECT symbol, amount, turnoverrate FROM gp_day_data " + "WHERE DATE(`timestamp`) = %s AND turnoverrate > 0 AND amount > 0", (day,)) + except Exception as e: # noqa: BLE001 + print(f" (市值分桶读取失败 {day}: {e!r})") + return {} + if df.empty: + return {} + df["k"] = df["symbol"].astype(str).str.strip().map(common.to_prefix) + df["mv"] = df["amount"] / df["turnoverrate"] + try: + df["cap"] = pd.qcut(df["mv"], 3, labels=["小盘", "中盘", "大盘"]) + except ValueError: + return {} + return dict(zip(df["k"], df["cap"].astype(str))) + + +def pms_roster(day: str) -> tuple[list[str], str]: + """PMS 当日拿到的生产名单:计划快照名册里档位强传导、按分数序前 N。返回 (代码, 注记)。""" + try: + df = db.read_mysql( + "pms", "SELECT roster_json, fetched_at FROM pms_plan_snapshot " + "WHERE plan_date = %s ORDER BY id DESC LIMIT 1", (day,)) + n_df = db.read_mysql( + "pms", "SELECT param_value FROM pms_runtime_param WHERE param_key='PMS_PLAN_TOP_N'") + except Exception as e: # noqa: BLE001 + return [], f"PMS 快照读取失败({e!r}),生产名单退为桥重算" + if df.empty: + return [], "PMS 无当日快照,生产名单退为桥重算" + n = int(n_df.iloc[0, 0]) if not n_df.empty else 30 + roster = json.loads(df.iloc[0]["roster_json"] or "[]") + rows = [r for r in roster if str(r.get("t") or "") == "强传导" and str(r.get("b") or "main") == "main"] + rows.sort(key=lambda r: -(r.get("s") or 0)) + return [common.to_prefix(str(r["c"]).strip()) for r in rows[:n] if r.get("c")], \ + f"PMS 快照名册(N={n} 为现值,历史 N 不可还原)" + + +# ============================================================================ +# 收益 +# ============================================================================ +def forward(pivot: pd.DataFrame, days: list[str], day: str, h: int, start: str) -> pd.Series | None: + """从 day 起按口径取后 h 日累计涨幅(每票)。数据不够返回 None。""" + if day not in days: + return None + i = days.index(day) + lo = i + (2 if start == "next_close" else 1) + hi = lo + h # 切片 [lo, hi) + if hi > len(days): + return None + block = pivot.loc[days[lo:hi]] + return block.sum(min_count=h) + + +def _md(df: pd.DataFrame) -> str: + """自己拼 Markdown 表:桥镜像没装 tabulate,pandas.to_markdown 用不了。""" + if df is None or df.empty: + return "(无数据)" + cols = [str(c) for c in df.columns] + lines = ["| " + " | ".join(cols) + " |", "|" + "---|" * len(cols)] + for _, r in df.iterrows(): + lines.append("| " + " | ".join("" if (isinstance(v, float) and pd.isna(v)) else str(v) + for v in r.tolist()) + " |") + return "\n".join(lines) + + +def summarize(ret: pd.Series, codes: list[str], base_all: float, base_main: float) -> dict | None: + r = ret.reindex([c for c in codes if c in ret.index]).dropna() + if r.empty: + return None + return {"n": int(len(r)), "ret": round(float(r.mean()), 2), + "excess_all": round(float(r.mean() - base_all), 2), + "excess_main": round(float(r.mean() - base_main), 2) if base_main is not None else None, + "beat": round(float((r > base_all).mean() * 100), 1)} + + +# ============================================================================ +# 主流程 +# ============================================================================ +def run(since: str, until: str | None, horizons: tuple, start: str, out_dir: str, + with_cards: bool = True) -> dict: + days_plan = plan_dates(since, until) + if not days_plan: + raise SystemExit("区间内没有档位日。") + px = price_panel(since) + pivot = px.pivot_table(index="trade_date", columns="k", values="pct") + days = sorted(pivot.index.tolist()) + accum_by_day = {d: dict(zip(g["k"], g["accum"])) for d, g in px.groupby("trade_date")} + + rows, notes = [], [] + for day in days_plan: + try: + data = plan.collect(day, top=5000, obs_top=5000, theme_cap=0) + except Exception as e: # noqa: BLE001 + notes.append(f"{day}: 计划重算失败 {e!r}") + continue + full = data.get("_full") or {} + main_rows, obs_rows = full.get("main", []), full.get("observe", []) + main_codes = [r["code"] for r in main_rows] + tier_of = {r["code"]: r.get("tier") for r in main_rows} + verdict_of = {r["code"]: r.get("verdict") for r in main_rows + obs_rows} + seg_codes = sorted({r["code"] for r in main_rows + obs_rows if r.get("evidence")}) + prod, prod_note = pms_roster(day) + if not prod: + prod = [r["code"] for r in main_rows if r.get("tier") == "强传导"][:100] + cands = [r["code"] for r in data.get("candidates") or []] + watch = [r["code"] for r in data.get("watch") or []] + caps = cap_bucket(day) + acc = accum_by_day.get(day, {}) + reg = None + try: + import regime + reg = regime.read_from_snapshot(day) + except Exception: # noqa: BLE001 + reg = None + + for h in horizons: + ret = forward(pivot, days, day, h, start) + if ret is None: + continue + base_all = float(ret.dropna().mean()) + main_ret = ret.reindex([c for c in main_codes if c in ret.index]).dropna() + base_main = float(main_ret.mean()) if not main_ret.empty else None + regime_post = "涨周" if base_all > 0 else "跌周" + regime_pre = (("弱势日" if reg.get("weak_day") else "非弱势日") + if reg and reg.get("weak_day") is not None else "无标签") + lists = { + "生产名单": prod, "候选单": cands, "关注单": watch, "环节名单": seg_codes, + "主榜等权": main_codes, "全池等权": list(ret.dropna().index), + } + for name, codes in lists.items(): + s = summarize(ret, codes, base_all, base_main) + if s: + rows.append({"date": day, "h": h, "list": name, "group": "全部", + "regime_post": regime_post, "regime_pre": regime_pre, **s}) + # 分组:档位、判决、吸筹三态、市值 + for tier in ("强传导", "弱传导", "无传导"): + codes = [c for c, t in tier_of.items() if t == tier] + s = summarize(ret, codes, base_all, base_main) + if s: + rows.append({"date": day, "h": h, "list": "主榜", "group": f"档位={tier}", + "regime_post": regime_post, "regime_pre": regime_pre, **s}) + for v in ("候选", "关注", "仅展示"): + codes = [c for c, vv in verdict_of.items() if vv == v] + s = summarize(ret, codes, base_all, base_main) + if s: + rows.append({"date": day, "h": h, "list": "档位表", "group": f"判决={v}", + "regime_post": regime_post, "regime_pre": regime_pre, **s}) + for st_label, pred in (("明确吸筹", lambda s: str(s).startswith("明确")), + ("潜在吸筹", lambda s: str(s).startswith("潜在")), + ("其他", lambda s: not (str(s).startswith("明确") or str(s).startswith("潜在")))): + codes = [c for c, s in acc.items() if s and pred(s)] + s = summarize(ret, codes, base_all, base_main) + if s: + rows.append({"date": day, "h": h, "list": "全池", "group": f"吸筹={st_label}", + "regime_post": regime_post, "regime_pre": regime_pre, **s}) + if caps: + for cap in ("小盘", "中盘", "大盘"): + codes = [c for c in cands if caps.get(c) == cap] + s = summarize(ret, codes, base_all, base_main) + if s: + rows.append({"date": day, "h": h, "list": "候选单", "group": f"市值={cap}", + "regime_post": regime_post, "regime_pre": regime_pre, **s}) + notes.append(f"{day}: 主榜 {len(main_codes)} 观察 {len(obs_rows)} 候选 {len(cands)} " + f"关注 {len(watch)} 生产 {len(prod)}({prod_note})") + + df = pd.DataFrame(rows) + if df.empty: + raise SystemExit("没有任何可算的期限(数据尾部不足)。") + os.makedirs(out_dir, exist_ok=True) + stamp = days_plan[-1] + csv_path = os.path.join(out_dir, f"复盘明细_{stamp}_{start}.csv") + df.to_csv(csv_path, index=False, encoding="utf-8-sig") + + # 汇总:按名单 × 期限(全部);按分组 × 期限;按事后环境 × 名单(只作解释) + def agg(g: pd.DataFrame) -> pd.Series: + return pd.Series({"days": g["date"].nunique(), "n": int(g["n"].sum()), + "ret": round(g["ret"].mean(), 2), + "excess_all": round(g["excess_all"].mean(), 2), + "beat": round(g["beat"].mean(), 1)}) + lists_tbl = df[df["group"] == "全部"].groupby(["list", "h"]).apply(agg).reset_index() + groups_tbl = df[df["group"] != "全部"].groupby(["list", "group", "h"]).apply(agg).reset_index() + regime_tbl = df[df["group"] == "全部"].groupby(["regime_post", "list", "h"]).apply(agg).reset_index() + pre_tbl = df[(df["group"] == "全部") & (df["regime_pre"] != "无标签")] \ + .groupby(["regime_pre", "list", "h"]).apply(agg).reset_index() + + md = [f"# 候选单复盘 · {days_plan[0]} 至 {stamp}(起点价 {start})", "", + f"计划日 {len(days_plan)} 期;样本纪律:分组样本少于一百只或覆盖计划日少于二十个只看方向;" + f"口径切换的总样本少于六十个计划日只写方向。", "", + "## 一、六份名单 × 期限(超额=相对全池等权,命中=跑赢全池比例)", "", + _md(lists_tbl), "", + "## 二、分组读数", "", _md(groups_tbl), "", + "## 三、按事后环境分组(未来 h 日全池涨跌,只作解释,不作交易前置)", "", + _md(regime_tbl), ""] + if not pre_tbl.empty: + md += ["## 四、按事前线上标签分组(快照 regime 段,上线后才有)", "", + _md(pre_tbl), ""] + else: + md += ["## 四、按事前线上标签分组", "", "(区间内没有带环境标签的快照,本节待环境标签上线后出现。)", ""] + md += ["## 五、逐日注记", ""] + [f"- {n}" for n in notes] + ["", + "## 六、拍板建议", "", "(只列读数与选项,不改任何东西——由每周五人工填写。)", ""] + md_path = os.path.join(out_dir, f"复盘_{stamp}_{start}.md") + with open(md_path, "w", encoding="utf-8") as f: + f.write("\n".join(md)) + print("\n".join(md[:8])) + print(f"\n已写入 {md_path} 与 {csv_path}") + return {"md": md_path, "csv": csv_path, "days": len(days_plan)} + + +def main() -> int: + ap = argparse.ArgumentParser(description="候选单复盘(只读,名单级观察收益,不是回测)") + ap.add_argument("--since", default="2026-07-29") + ap.add_argument("--until") + ap.add_argument("--horizons", default="5,10,20") + ap.add_argument("--start-price", choices=["next_close", "signal_close"], default="next_close", + help="next_close=次日收盘起算(默认,实盘口径);signal_close=信号日收盘起算(对账方案 1.8 系列)") + ap.add_argument("--out", default="data/review") + a = ap.parse_args() + hs = tuple(int(x) for x in a.horizons.split(",") if x.strip()) + run(a.since, a.until, hs, a.start_price, a.out) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main())