akg-factor-bridge/plan_review.py

314 lines
16 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""候选单复盘(只读):这个系统唯一的评价方式——名单级观察收益,不是回测。
## 为什么要有它
方案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里档位为
强传导、按分数降序前 NN 取当时 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 表:桥镜像没装 tabulatepandas.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:
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())