"""候选单复盘(只读):这个系统唯一的评价方式——名单级观察收益,不是回测。 ## 为什么要有它 方案(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] obs_codes = [r["code"] for r in obs_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, "观察档等权": obs_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: # 主口径按日等权:每个计划日先算名单均值,再跨日平均——一份名单一天一票; # 带 _w 的三列按样本加权(大名单的日子权重大),与方案 1.8b/1.8c 手工读数同口径,只作对账。 n = g["n"].astype(float) w = n / n.sum() if n.sum() else n return pd.Series({"days": g["date"].nunique(), "n": int(n.sum()), "ret": round(g["ret"].mean(), 2), "excess_all": round(g["excess_all"].mean(), 2), "beat": round(g["beat"].mean(), 1), "ret_w": round(float((g["ret"] * w).sum()), 2), "excess_w": round(float((g["excess_all"] * w).sum()), 2), "beat_w": round(float((g["beat"] * w).sum()), 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"口径切换的总样本少于六十个计划日只写方向。", "", "## 一、六份名单 × 期限(超额=相对全池等权,命中=跑赢全池比例;" "不带后缀的列按计划日等权,带 _w 的列按样本加权,后者只用于与方案第一节手工读数对账)", "", _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())