2026-07-30 12:44:56 +08:00
|
|
|
|
"""每日选股计划(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}"
|
|
|
|
|
|
|
|
|
|
|
|
|
2026-07-30 12:55:24 +08:00
|
|
|
|
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:
|
2026-07-30 12:44:56 +08:00
|
|
|
|
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)
|
2026-07-30 12:55:24 +08:00
|
|
|
|
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)
|
2026-07-30 12:44:56 +08:00
|
|
|
|
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)}(T−1 口径:"
|
|
|
|
|
|
f"\"谁已经动了\"看的是上个交易日收盘;拍点方案定版前均如此)。")
|
|
|
|
|
|
L.append("")
|
|
|
|
|
|
|
2026-07-30 12:55:24 +08:00
|
|
|
|
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})")
|
2026-07-30 12:44:56 +08:00
|
|
|
|
L.append("")
|
|
|
|
|
|
L.append("| # | 代码 | 名称 | 总分 | 档位 | 传导证据 | 热度 | 预期空间 |")
|
|
|
|
|
|
L.append("|---|------|------|------|------|----------|------|----------|")
|
2026-07-30 12:55:24 +08:00
|
|
|
|
for i, (k, s) in enumerate(main_rows, 1):
|
2026-07-30 12:44:56 +08:00
|
|
|
|
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("")
|
|
|
|
|
|
|
2026-07-30 12:55:24 +08:00
|
|
|
|
obs_rows = _pick(obs, ev, obs_top, theme_cap)
|
|
|
|
|
|
L.append(f"## 观察档 Top {len(obs_rows)}"
|
|
|
|
|
|
f"(无券商预期、但在传导链上——没有估值锚,置信度低{cap_txt})")
|
2026-07-30 12:44:56 +08:00
|
|
|
|
L.append("")
|
|
|
|
|
|
L.append("| # | 代码 | 名称 | 分 | 传导证据 | 热度 |")
|
|
|
|
|
|
L.append("|---|------|------|----|----------|------|")
|
2026-07-30 12:55:24 +08:00
|
|
|
|
for i, (k, s) in enumerate(obs_rows, 1):
|
2026-07-30 12:44:56 +08:00
|
|
|
|
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 + 组内分(还没热、还便宜);"
|
2026-07-30 13:02:43 +08:00
|
|
|
|
"观察档分 = 100 + 0.6z(传导) + 0.4z(−热度)。")
|
2026-07-30 12:44:56 +08:00
|
|
|
|
|
|
|
|
|
|
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
|