akg-factor-bridge/plan_review.py

355 lines
20 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,
caps: dict | None = None, cap_base: dict | None = None) -> dict | None:
"""一份名单在一个计划日、一个期限上的读数。
excess_all 相对全池等权excess_main 相对主榜等权(同为有券商覆盖的篮子,剥掉大票对小票的贝塔);
excess_cap 相对同市值桶均值(每票减当日同桶全池均值再平均,方案 1.8c 的市值中性口径)。"""
r = ret.reindex([c for c in codes if c in ret.index]).dropna()
if r.empty:
return None
out = {"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,
"excess_cap": None,
"beat": round(float((r > base_all).mean() * 100), 1)}
if caps and cap_base:
adj = [float(v) - cap_base[caps[c]] for c, v in r.items() if caps.get(c) in cap_base]
if adj:
out["excess_cap"] = round(sum(adj) / len(adj), 2)
return out
# ============================================================================
# 主流程
# ============================================================================
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:
# 无 PMS 快照的日子不再退化为桥重算(那不是 PMS 拿到的名单):交付名单当日剔除并注记(台账 010
prod_note = f"{prod_note};无快照,交付名单当日剔除"
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
cap_base: dict = {} # 当日三个市值桶各自的全池均值,供 excess_cap 用
if caps:
for cap in ("小盘", "中盘", "大盘"):
v = ret.reindex([c for c, b in caps.items() if b == cap and c in ret.index]).dropna()
if not v.empty:
cap_base[cap] = float(v.mean())
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, caps, cap_base)
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, caps, cap_base)
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, caps, cap_base)
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, caps, cap_base)
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, caps, cap_base)
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
days, total = g["date"].nunique(), int(n.sum())
# 样本纪律(方案 2.4)机器标注,只有两档:可读=计划日≥20 且样本≥100其余只看方向
grade = "可读" if (days >= 20 and total >= 100) else "方向"
return pd.Series({"days": days, "n": total,
"ret": round(g["ret"].mean(), 2),
"excess_all": round(g["excess_all"].mean(), 2),
"excess_main": round(g["excess_main"].mean(), 2) if g["excess_main"].notna().any() else None,
"excess_cap": round(g["excess_cap"].mean(), 2) if g["excess_cap"].notna().any() else None,
"beat": round(g["beat"].mean(), 1),
"share": round(float(n.max() / n.sum() * 100), 0) if n.sum() else None,
"grade": grade,
"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()
h0 = int(df["h"].min())
daily_c = df[(df["list"] == "候选单") & (df["group"] == "全部") & (df["h"] == h0)] \
[["date", "n", "ret", "excess_all", "excess_cap", "beat", "regime_post"]].sort_values("date")
md = [f"# 候选单复盘 · {days_plan[0]}{stamp}(起点价 {start}", "",
f"计划日 {len(days_plan)} 期;样本纪律:分组样本少于一百只或覆盖计划日少于二十个只看方向;"
f"口径切换的总样本少于六十个计划日只写方向。", "",
"## 一、六份名单 × 期限", "",
"列的读法excess_all 相对全池等权excess_main 相对主榜等权(剥掉大票对小票的贝塔),"
"excess_cap 相对同市值桶(成交额除以换手率三分位,方案 1.8c 口径beat 是跑赢全池比例,"
"全池自己的 beat 只有四成多分布右偏读它时与全池等权那一行相减share 是最大单日样本占比"
"超过三成说明读数由一天主导只看逐日表grade 是样本纪律三档。"
"不带后缀的列按计划日等权(主口径),带 _w 的列按样本加权,只用于与方案第一节手工读数对账。", "",
_md(lists_tbl), "",
"## 二、分组读数", "", _md(groups_tbl), "",
f"## 二之二、候选单逐日(期限 {h0} 日;第一节按日等权的读数就是这张表的平均,看集中度)", "",
_md(daily_c), "",
"## 三、按事后环境分组(未来 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())