akg-factor-bridge/plan.py

223 lines
9.4 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.

"""每日选股计划R4 首版):把 akg_gate / akg_score 变成一份人能读的榜单。
数据全部来自已落库的表,不重算:
平台因子表 t_factor_akg_score / _gate / _upside / _heat —— 当日截面
基座只读视图 v_factor_transmission —— 传导证据(主题、源数、已动比例)
基座 industry_pools —— 股票名称
产出:终端打印 + data/plan/plan_<日期>.md。升降档一节对比前一交易日的档位表
体现"数据到达本身是信号"(首次覆盖 / 新进传导链即升档)。
"""
import json
import os
import pandas as pd
import common
import db
# 分数编码(与 factors.build_score 一致):主榜 = 200 + 传导档位×20 + 组内分,
# 观察档 = 100 + 组内分,组内分 clip ±9.9。150 落在两带中间的空档上,用作分界。
_MAIN_MIN = 150.0
def _factor(table: str, ds: str) -> pd.Series:
df = db.read_mysql(
"factor",
f"SELECT stock_code, factor_value FROM {table} WHERE trade_date = %s", (ds,))
if df.empty:
return pd.Series(dtype=float)
return df.set_index("stock_code")["factor_value"].astype(float)
def _latest_date(table: str, upto: str | None = None):
if upto:
df = db.read_mysql(
"factor", f"SELECT MAX(trade_date) d FROM {table} "
f"WHERE trade_date <= %s", (upto,))
else:
df = db.read_mysql("factor", f"SELECT MAX(trade_date) d FROM {table}")
v = None if df.empty else df.iloc[0, 0]
return None if v is None or pd.isna(v) else pd.Timestamp(v).date().isoformat()
def _prev_date(table: str, before: str):
df = db.read_mysql(
"factor", f"SELECT MAX(trade_date) d FROM {table} "
f"WHERE trade_date < %s", (before,))
v = None if df.empty else df.iloc[0, 0]
return None if v is None or pd.isna(v) else pd.Timestamp(v).date().isoformat()
def _names() -> dict:
out = {}
pools = db.read_pg("SELECT members FROM industry_pools")
for _, r in pools.iterrows():
ms = r["members"]
if isinstance(ms, str):
ms = json.loads(ms)
for m in ms or []:
ts, name = (m or {}).get("ts_code"), (m or {}).get("name")
if ts and name:
out.setdefault(common.to_prefix(ts), str(name))
return out
def _evidence(ds: str):
"""每股最强一条传导证据:主题、源数、已动比例;另返回涉及的行情快照日。"""
tr = db.read_pg(
"SELECT ts_code, target, n_sources, moved_ratio, mkt_trade_date "
"FROM v_factor_transmission WHERE scan_date = %s", (ds,))
if tr.empty:
return {}, set()
tr["k"] = tr["ts_code"].map(common.to_prefix)
tr["n_sources"] = pd.to_numeric(tr["n_sources"], errors="coerce").fillna(0)
tr["moved_ratio"] = pd.to_numeric(tr["moved_ratio"], errors="coerce").fillna(0)
tr["strength"] = tr["n_sources"] * (1.0 - tr["moved_ratio"])
tr = tr.sort_values("strength", ascending=False).drop_duplicates("k")
ev = {r.k: (str(r.target), int(r.n_sources), float(r.moved_ratio))
for r in tr.itertuples()}
days = {str(x) for x in tr["mkt_trade_date"].dropna().unique()}
return ev, days
def _tier_label(score: float) -> str:
return {0: "无传导", 1: "弱传导", 2: "强传导"}.get(
int((score - 190.0) // 20), "?")
def _pct(v) -> str:
return "" if v is None or pd.isna(v) else f"{v:+.0%}"
def _num(v) -> str:
return "" if v is None or pd.isna(v) else f"{v:.2f}"
def _pick(ranked: pd.Series, ev: dict, top: int, theme_cap: int):
"""按分数从高到低取 top 条;每个传导主题最多 theme_cap 条0=不设限)。
为什么设限:传导目标是环节级,同环节全体成员共享同一条证据,不设限时
榜单会被三四个环节刷屏——20 个名额实际只是 4 注。限额后变成
"多条线索 × 每条线索取最冷最便宜的几只",被挤掉的仍在完整档位表里。"""
out, cnt = [], {}
for k, s in ranked.items():
e = ev.get(k)
theme = e[0] if e else "(无传导)"
if theme_cap and cnt.get(theme, 0) >= theme_cap:
continue
cnt[theme] = cnt.get(theme, 0) + 1
out.append((k, s))
if len(out) >= top:
break
return out
def generate(date: str | None = None, top: int = 20, obs_top: int = 10,
theme_cap: int = 5) -> str:
ds = date or _latest_date("t_factor_akg_score")
if not ds:
raise SystemExit("t_factor_akg_score 还没有数据——先 build akg_score。")
score = _factor("t_factor_akg_score", ds)
gate = _factor("t_factor_akg_gate", ds)
if score.empty or gate.empty:
raise SystemExit(f"{ds} 缺 akg_score / akg_gate——先 build 该日再出计划。")
upside = _factor("t_factor_akg_upside", ds)
if upside.empty:
# 晚间 18:40 建当日 upside 表时行情源(~19:50 发布)还没到,当日表常为空。
# 这里现算兜底此刻库里已有晚到的当日价as-of 口径不变consensus<=当日)。
import factors
df_up = factors.build_upside(ds, ds)
if df_up is not None and not df_up.empty:
x = df_up.copy()
x["k"] = x["stock_code"].map(common.to_prefix)
upside = x.groupby("k")["factor_value"].max().astype(float)
hd = _latest_date("t_factor_akg_heat", ds)
heat = _factor("t_factor_akg_heat", hd) if hd else pd.Series(dtype=float)
names = _names()
ev, mkt_days = _evidence(ds)
main = score[score >= _MAIN_MIN].sort_values(ascending=False)
obs = score[score < _MAIN_MIN].sort_values(ascending=False)
L = [f"# 每日选股计划 · {ds}", ""]
L.append(f"主榜 {len(main)} 只 / 观察档 {len(obs)} 只 / 全池档位覆盖 {len(gate)} 只。")
stale = sorted(d for d in mkt_days if d != ds)
if stale:
L.append(f"⚠️ 本日传导用的行情快照 = {''.join(stale)}T1 口径:"
f"\"谁已经动了\"看的是上个交易日收盘;拍点方案定版前均如此)。")
L.append("")
cap_txt = f",每主题限额 {theme_cap}" if theme_cap else ""
main_rows = _pick(main, ev, top, theme_cap)
L.append(f"## 主榜 Top {len(main_rows)}(有券商预期、目标价不低于现价{cap_txt}")
L.append("")
L.append("| # | 代码 | 名称 | 总分 | 档位 | 传导证据 | 热度 | 预期空间 |")
L.append("|---|------|------|------|------|----------|------|----------|")
for i, (k, s) in enumerate(main_rows, 1):
e = ev.get(k)
etxt = f"{e[0]}{e[1]} 源,已动 {e[2]:.0%}" if e else ""
L.append(f"| {i} | {k} | {names.get(k, '')} | {s:.1f} | {_tier_label(s)} "
f"| {etxt} | {_num(heat.get(k))} | {_pct(upside.get(k))} |")
L.append("")
obs_rows = _pick(obs, ev, obs_top, theme_cap)
L.append(f"## 观察档 Top {len(obs_rows)}"
f"(无券商预期、但在传导链上——没有估值锚,置信度低{cap_txt}")
L.append("")
L.append("| # | 代码 | 名称 | 分 | 传导证据 | 热度 |")
L.append("|---|------|------|----|----------|------|")
for i, (k, s) in enumerate(obs_rows, 1):
e = ev.get(k)
etxt = f"{e[0]}{e[1]} 源,已动 {e[2]:.0%}" if e else ""
L.append(f"| {i} | {k} | {names.get(k, '')} | {s:.1f} "
f"| {etxt} | {_num(heat.get(k))} |")
L.append("")
L.append("## 今日升降档")
L.append("")
prev_ds = _prev_date("t_factor_akg_gate", ds)
if not prev_ds:
L.append("(没有更早的档位表可比,升降档从下一个交易日开始。)")
else:
prev = _factor("t_factor_akg_gate", prev_ds)
both = pd.concat([prev.rename("prev"), gate.rename("cur")], axis=1)
both = both.fillna(-1.0) # -1 = 当日不在面板
up = both[both["cur"] > both["prev"]]
down = both[both["cur"] < both["prev"]]
lab = {-1.0: "池外", 0.0: "不采纳", 1.0: "观察档", 2.0: "主榜"}
L.append(f"对比 {prev_ds}:升档 {len(up)} 只,降档 {len(down)} 只。"
f"升档=拿到新锚(首次覆盖 / 新进传导链),本身就是值得看的信号。")
def _rows(d: pd.DataFrame, cap: int = 15):
lines = []
for k, r in d.iterrows():
lines.append(f"- {k} {names.get(k, '')}"
f"{lab.get(r['prev'], '?')}{lab.get(r['cur'], '?')}")
if len(lines) >= cap:
lines.append(f"- ……共 {len(d)} 只,其余见档位表")
break
return lines
if not up.empty:
L.append("")
L.append("**升档**")
L += _rows(up.sort_values("cur", ascending=False))
if not down.empty:
L.append("")
L.append("**降档**")
L += _rows(down.sort_values("prev", ascending=False))
L.append("")
L.append("---")
L.append("口径:主榜分 = 200 + 传导档位×20 + 组内分(还没热、还便宜);"
"观察档分 = 100 + 0.6z(传导) + 0.4z(−热度)。")
text = "\n".join(L)
os.makedirs("data/plan", exist_ok=True)
out = f"data/plan/plan_{ds}.md"
with open(out, "w", encoding="utf-8") as f:
f.write(text + "\n")
print(text)
print(f"\n已写入 {out}")
return out