From b813e42f5254a4866cf127bbb347f5101405411f Mon Sep 17 00:00:00 2001 From: zlt Date: Thu, 30 Jul 2026 12:01:11 +0800 Subject: [PATCH] =?UTF-8?q?G2+R3:=20=E8=B5=9B=E9=81=93=E6=98=A0=E5=B0=84?= =?UTF-8?q?=E6=A8=A1=E5=9D=97=20+=20akg=5Fgate/akg=5Fscore=20=E5=90=88?= =?UTF-8?q?=E6=88=90=E5=9B=A0=E5=AD=90=EF=BC=8807-30=20=E4=B8=89=E6=8B=8D?= =?UTF-8?q?=E8=90=BD=E5=9C=B0=EF=BC=89=20=E4=B8=A4=E9=94=9A=E4=B8=89?= =?UTF-8?q?=E6=A1=A3=20+=20=E4=B8=A4=E6=AE=B5=E5=BC=8F=EF=BC=9BC=20?= =?UTF-8?q?=E9=97=B8=E9=BB=98=E8=AE=A4=E5=85=B3=E5=BE=85=20yml=20=E8=BD=AC?= =?UTF-8?q?=E6=AD=A3=EF=BC=9Brequirements=20=E5=8A=A0=20PyYAML=20=E9=9C=80?= =?UTF-8?q?=20rebuild=E3=80=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- config.py | 7 ++ docs/G1实测与待办_2026-07-26.md | 34 +++++++- factors.py | 141 ++++++++++++++++++++++++++++++++ requirements.txt | 1 + run.py | 22 ++++- tracks.py | 106 ++++++++++++++++++++++++ 6 files changed, 308 insertions(+), 3 deletions(-) create mode 100644 tracks.py diff --git a/config.py b/config.py index 29e2f7d..7cd3567 100644 --- a/config.py +++ b/config.py @@ -107,3 +107,10 @@ FROZEN_ROOT = os.environ.get("FROZEN_ROOT", "/app/data/frozen") # 07-28 语义更新:基座二批后 topic_context cap=1000、scan 大主题闸=600、quiet 全量落库。 # 本值对齐基座"大主题闸":members_total 超过它 = 基座制度回退,桥侧告警(factors.py)。 UPSTREAM_MEMBER_CAP = int(os.environ.get("UPSTREAM_MEMBER_CAP", "600")) + +# --- 赛道门槛 C(G2:config/frontier_tracks.yml + tracks.py)----------------- +# yml 是唯一事实源(设计 §3.2)。C 闸默认关:映射还是 v0.1 草案,先跑 +# `python run.py tracks` 做覆盖体检、清单转正后再置 1——届时主榜(gate=2) +# 再交赛道成员,观察档的产业链锚也并入赛道成员。 +TRACKS_YML = os.environ.get("TRACKS_YML", "config/frontier_tracks.yml") +ENABLE_TRACK_GATE = os.environ.get("ENABLE_TRACK_GATE", "0") == "1" diff --git a/docs/G1实测与待办_2026-07-26.md b/docs/G1实测与待办_2026-07-26.md index 54c8a94..4105996 100644 --- a/docs/G1实测与待办_2026-07-26.md +++ b/docs/G1实测与待办_2026-07-26.md @@ -228,4 +228,36 @@ probe 三张读数提醒已钉 **07-30 周四 19:00**。 **传导 live 完好起点 = 07-17**(完好日仅 17/20/22/23 四天); 快照完整性建议纳入周检/巡检(非二批)。 ④ 今晚 17:30 盯行数:≈2391=健康且池源同步;≈1675=sync 用滞后池源(新问题); - 再挂=马上查日志。07-24 缺口的 sync 日志仍待查。 \ No newline at end of file + 再挂=马上查日志。07-24 缺口的 sync 日志仍待查。 + +## 七、07-30 拍板与交付记录 + +1. **快照 T−1 结构性错位定案**:源侧每日 ~19:50 才发布当日行情,两拍 17:30/17:55 + 永远只拿到前一日 → 传导天天用昨日 movers(07-28 扫描用 07-27 快照、07-29 用 + 07-28,台账 mkt_trade_date 如实打标)。07-27 事故是这个常态的极端版;"缺日"是 + 滞后不是丢失(某日行情在次日傍晚自动补上)。**拍点四案待拍**:A 整链后移 ~20:30 / + B 加第三拍+当晚重扫 / C 正式接受 T−1 语义 / D 整链挪到次日盘前。倾向 D(盘前出 + 选股计划,晚间披露的公告也能进当日事件)或 A;07-30 起记录源侧到点,与三批-5 + 采信规则同题拍板。 +2. **9-4 主决策已拍:两锚三档**。业绩锚=券商覆盖(upside 可算);产业链锚=图谱 + 传导链(赛道清单转正后并入赛道成员)。主榜=有覆盖且 upside≥0(C 闸待转正); + 观察档=无覆盖但有产业链锚,算分用传导+热度(暂拟 0.6z_T+0.4z_H),明确标低置信; + 两锚皆无=不采纳。升降档本身是信号(首次覆盖/新进传导链),进日报。 +3. **三批-3 已拍:主榜组内排序=两段式**。传导票按 log 传导分的组内中位数分强弱两档, + 档内按"还没热、还便宜"(0.6z(−热度)+0.4z(upside))排。依据(07-29 截面):传导 + 权重 0.5 时 top50 有 48 只传导票=现行加权事实上就是传导优先;两种排法 top20 仅 + 重合 14/20=组内排序就是榜单本身。三项相关 |ρ|≤0.063 正交复现;9-3 的 q 暂留 0 + (门槛②对有覆盖股只剔 33/994,S2 再议)。 +4. **probe 修复判收**:corr 节 scipy 崩点(pandas 的 Series.corr(spearman) 会 + import scipy,DataFrame.corr 不会)已改为复用整表矩阵,07-30 实机全节跑通。 +5. **07-30 交付(待实机验证)**: + ① 基座 transmission.py 传导主题黑名单(高新技术企业/小型微利企业/联营企业/无/ + 报告分部——联营企业 298 名成员不触 600 闸,必须显式拦;源与目标两侧都挡); + ② 桥 config/frontier_tracks.yml v0.1(6 confirmed + 4 candidate、61 主题映射、 + 排除清单、TODO 决策点待圈);③ 桥 tracks.py + `run.py tracks`(赛道覆盖体检 + + 成员表版本化快照落 data/,即本清单"赛道覆盖体检"的正式工具);④ 桥 akg_gate / + akg_score 两个新因子(三档+两段式,进 build all 日更;C 闸默认关 + =ENABLE_TRACK_GATE,yml 转正后开);⑤ requirements.txt 加 PyYAML(**桥机需 + rebuild 镜像**;不 rebuild 只影响 tracks 命令,其余照跑)。 + **工作哲学(07-30 用户令)**:数据渐进式暴露是常态——不等全量抽完、不等干净 + 截面、不等覆盖完备;缺锚降档不弃用,滞后交采信规则消化。 \ No newline at end of file diff --git a/factors.py b/factors.py index 49c4711..851e432 100644 --- a/factors.py +++ b/factors.py @@ -20,6 +20,9 @@ FACTORS = { "akg_heat": "t_factor_akg_heat", "akg_event": "t_factor_akg_event", "akg_transmission": "t_factor_akg_transmission", + # 合成层(07-30 拍板:两锚三档 + 两段式),build all 一并日更 + "akg_gate": "t_factor_akg_gate", + "akg_score": "t_factor_akg_score", } # ---- 事件极性草案(设计 §5.3,待 §9-5 确认)---- @@ -304,7 +307,145 @@ def _warn_upstream_truncation(tr: pd.DataFrame) -> None: f"若构建窗口含 07-24 及更早的旧制度日,此告警对旧日属预期)") +# ---------------------------------------------------------------- 合成层 +# 07-30 拍板:门槛=两锚三档(akg_gate),主榜组内排序=两段式(akg_score)。 +# 依据(07-29 截面权重体检):传导权重 0.5 时 top50 有 48 只传导票——现行加权 +# 事实上就是传导优先;两种排法 top20 仅重合 14/20,组内排序即榜单本身。 + +def _robust_z(v: pd.Series) -> pd.Series: + med = v.median() + mad = (v - med).abs().median() + if pd.isna(mad) or mad == 0: + sd = v.std() + return (v - v.mean()) / sd if sd and sd > 0 else v * 0.0 + return (v - med) / (1.4826 * mad) + + +def _heat_day_for(ds: str): + """热度 T+1 到达 → 取 <= ds 的最新热度日(与 probe / 设计 §3.4 末口径一致)。""" + df = db.read_mysql("heat", "SELECT MAX(trade_date) d FROM stock_fund_heat_scores " + "WHERE trade_date <= %s", (ds,)) + 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 _col(df: pd.DataFrame, name: str) -> pd.Series: + """子因子 builder 输出 → 前缀码索引的一列(同股同日取最大)。""" + if df is None or df.empty: + return pd.Series(dtype=float, name=name) + x = df.copy() + x["k"] = x["stock_code"].map(common.to_prefix) + return x.groupby("k")["factor_value"].max().astype(float).rename(name) + + +def _panel(ds: str) -> pd.DataFrame: + """单日截面:upside / heat / transmission 三列,索引=前缀码。""" + up = _col(build_upside(ds, ds), "upside") + hd = _heat_day_for(ds) + ht = _col(build_heat(hd, hd), "heat") if hd else pd.Series(dtype=float, name="heat") + tr = _col(build_transmission(ds, ds), "transmission") + return pd.concat([up, ht, tr], axis=1) + + +def _track_set(): + """已转正赛道的成员集合(前缀码)。C 闸未启用时返回 None。""" + if not config.ENABLE_TRACK_GATE: + return None + try: + import tracks + df, _missing = tracks.resolve_members(only_confirmed=True) + return {common.to_prefix(x) for x in df["ts_code"]} + except Exception as e: # noqa: BLE001 —— 解析失败按闸未启用降级,只告警不中止 + print(f" ⚠️ 赛道成员表解析失败,C 闸按未启用处理: {e!r}") + return None + + +def _gate_of(p: pd.DataFrame, track_set) -> pd.Series: + """三档(07-30 拍板): + 2 = 主榜:有业绩锚(券商覆盖且 upside>=0);C 闸开启后再交赛道成员; + 1 = 观察档:无业绩锚但有产业链锚(当前=在传导链上;C 转正后并入赛道成员); + 0 = 不采纳:两锚皆无,或有覆盖但 upside<0(贵了不买是绝对下限)。""" + covered = p["upside"].notna() + chain = p["transmission"].fillna(0.0) > 0 + if track_set: + in_track = pd.Series(p.index.isin(track_set), index=p.index) + chain = chain | in_track + g = pd.Series(0.0, index=p.index) + main = covered & (p["upside"] >= 0.0) + if track_set: + main = main & in_track + g[main] = 2.0 + g[(~covered) & chain] = 1.0 + return g + + +def _days_of(start, end): + days = [pd.Timestamp(d).date().isoformat() + for d in common.trading_days(start, end)] + if not days and start == end: + days = [start] # 日历查不到也照算单日:以调用方给的日期为准 + return days + + +def build_gate(start, end): + """akg_gate ∈ {0,1,2}:全池出行——让平台看得见门槛本身,不只是池内排序。 + history 模式逐日重算三路子因子,区间大会慢;建议 daily / 短区间。""" + ts = _track_set() + rows = [] + for ds in _days_of(start, end): + p = _panel(ds) + if p.empty: + continue + g = _gate_of(p, ts) + rows.append(pd.DataFrame({"trade_date": ds, "stock_code": g.index, + "factor_value": g.values})) + return pd.concat(rows, ignore_index=True) if rows else _EMPTY + + +def build_score(start, end): + """akg_score:先档后分,仅 gate>0 出行,数值直接可排序。 + + 主榜 = 200 + 传导档位×10 + 组内分:传导票按 log 传导分的组内中位数分成 + 强(2)/弱(1)两档,无传导 0 档;组内分 = 0.6·z(−热度) + 0.4·z(upside), + 即"还没热、还便宜"。观察档 = 100 + 0.6·z(传导) + 0.4·z(−热度)——无估值锚, + 低置信。组内分夹在 ±9.9 保证档位永不重叠;z 都在各自档内算。 + 缺失处置同设计 §3.4:传导缺=0,热度缺=档内中位数。""" + ts = _track_set() + out = [] + for ds in _days_of(start, end): + p = _panel(ds) + if p.empty: + continue + g = _gate_of(p, ts) + + P = p[g == 2.0].copy() # ---- 主榜 + if not P.empty: + P["transmission"] = P["transmission"].fillna(0.0) + P["heat"] = P["heat"].fillna(P["heat"].median()) + lt = np.log1p(P["transmission"]) + tier = pd.Series(0.0, index=P.index) + hit = P["transmission"] > 0 + if hit.any(): + med = lt[hit].median() + tier[hit] = np.where(lt[hit] >= med, 2.0, 1.0) + inner = 0.6 * (-_robust_z(P["heat"])) + 0.4 * _robust_z(P["upside"]) + score = 200.0 + tier * 10.0 + inner.clip(-9.9, 9.9) + out.append(pd.DataFrame({"trade_date": ds, "stock_code": P.index, + "factor_value": score.values})) + + O = p[g == 1.0].copy() # ---- 观察档 + if not O.empty: + O["heat"] = O["heat"].fillna(O["heat"].median()) + zt = _robust_z(np.log1p(O["transmission"].fillna(0.0))) + inner = 0.6 * zt + 0.4 * (-_robust_z(O["heat"])) + score = 100.0 + inner.clip(-9.9, 9.9) + out.append(pd.DataFrame({"trade_date": ds, "stock_code": O.index, + "factor_value": score.values})) + return pd.concat(out, ignore_index=True) if out else _EMPTY + + BUILDERS = { "akg_upside": build_upside, "akg_heat": build_heat, "akg_event": build_event, "akg_transmission": build_transmission, + "akg_gate": build_gate, "akg_score": build_score, } diff --git a/requirements.txt b/requirements.txt index f583d31..8319cc0 100644 --- a/requirements.txt +++ b/requirements.txt @@ -3,3 +3,4 @@ numpy>=1.24 psycopg[binary]>=3.1 PyMySQL>=1.1 python-dotenv>=1.0 +PyYAML>=6.0 diff --git a/run.py b/run.py index f529f10..b522eb1 100644 --- a/run.py +++ b/run.py @@ -4,7 +4,8 @@ python run.py apply-views [--dry-run] # 把插槽视图 DDL 应用到基座 PG python run.py probe # G1 体检(只读,详见 probe.py) python run.py freeze [--date D] # 输入冻结(G0.5,详见 freeze.py) - python run.py register # 注册四子因子到 factor_metadata + python run.py tracks # 赛道覆盖体检 + 成员表快照(G2) + python run.py register # 注册全部因子到 factor_metadata python run.py build all --mode history --start 2024-01-01 --end 2025-12-31 python run.py build akg_heat --mode daily --date 2026-07-24 python run.py build akg_event --mode history --start 2025-01-01 --end 2026-07-24 @@ -88,11 +89,19 @@ _META = { "akg_heat": ("astock-kg 热度", "生态日频资金热度分(0~1)"), "akg_event": ("astock-kg 事件", "利好利空事件时间衰减加权分(仅公告来源,单文档封顶)"), "akg_transmission": ("astock-kg 传导", "板块传导未动成员传导强度(distinct源数×(1-已动比例))"), + "akg_gate": ("astock-kg 门槛档位", + "三档置信门槛(07-30拍板): 2=主榜(券商覆盖且upside>=0; 赛道C闸未启用, " + "转正后再交赛道成员), 1=观察档(无券商覆盖但在图谱传导链上, 无估值锚, " + "低置信), 0=不采纳(两锚皆无, 或upside<0)。全池出行, 让平台看得见门槛"), + "akg_score": ("astock-kg 景气度漏斗", + "先档后分: 主榜=200+传导档位x10+组内分(0.6z(-热度)+0.4z(upside), " + "两段式07-30拍板); 观察档=100+0.6z(传导)+0.4z(-热度)。" + "仅gate>0出行, 数值直接可排序"), } def cmd_register(): - print("注册四子因子:") + print(f"注册因子(共 {len(_META)} 个):") for code, (name, desc) in _META.items(): common.register(code, name, factors.FACTORS[code], ["astock-kg", code.split("_", 1)[1]], desc) @@ -130,6 +139,7 @@ def main(): sub = ap.add_subparsers(dest="cmd", required=True) sub.add_parser("views") sub.add_parser("register") + sub.add_parser("tracks") # 赛道覆盖体检 + confirmed 成员表快照(G2) p = sub.add_parser("probe") p.add_argument("--section", choices=["all", "pools", "price", "upside", "corr"], default="all", @@ -161,6 +171,14 @@ def main(): freeze.snapshot(a.date) elif a.cmd == "register": cmd_register() + elif a.cmd == "tracks": + import tracks + tracks.coverage_report() + out, df, missing = tracks.snapshot(only_confirmed=True) + print(f"\nconfirmed 成员表快照: {out}" + f"({df['ts_code'].nunique()} 只,{len(df)} 行 股票×赛道)") + if missing: + print(f"⚠️ {len(missing)} 个主题在 industry_pools 里没找到(见体检表)") elif a.cmd == "build": cmd_build(a.factor, a.mode, a.start, a.end, a.date, do_freeze=not a.no_freeze) diff --git a/tracks.py b/tracks.py new file mode 100644 index 0000000..3ec07cd --- /dev/null +++ b/tracks.py @@ -0,0 +1,106 @@ +"""赛道映射(硬门槛 C 的实现载体,设计 §3.2)。 + +config/frontier_tracks.yml 是唯一事实源:每个赛道列出 kg_themes +(industry_pools 主题名)。本模块把它解析成 ts_code 级成员表,并落 +data/track_members_<日期>.csv 版本化快照——可 git diff、可审计: +每只股票能追溯到因哪个赛道、哪个主题入选。 + +当前唯一的映射路径是主题名(kg_segments / kg_concepts 要等基座的环节 +投影表建成,即三批-4),source_rule 统一记 'pool_theme'。 +""" +import datetime as dt +import json +import os + +import pandas as pd + +import config +import db + +try: + import yaml +except ImportError: # 镜像未装 PyYAML 时给出可执行的修复指令,而不是裸崩 + yaml = None + + +def _need_yaml(): + if yaml is None: + raise SystemExit( + "缺 PyYAML:requirements.txt 已加,请在桥机重建镜像——\n" + " docker compose build akg-factor-bridge && " + "docker compose up -d akg-factor-bridge") + + +def load_yml(path: str | None = None) -> dict: + _need_yaml() + with open(path or config.TRACKS_YML, "r", encoding="utf-8") as f: + return yaml.safe_load(f) + + +def resolve_members(only_confirmed: bool = True, path: str | None = None): + """yml + industry_pools 当前态 → (成员表, 未命中主题清单)。 + + 成员表列:ts_code, name, track, status, theme。同股同赛道多主题只留一行, + 同股跨赛道保留多行。未命中 = yml 里写了、industry_pools 里查无此主题 + (通常是主题改名或池尚未涌现,体检时重点看)。 + """ + d = load_yml(path) + pools = db.read_pg("SELECT theme, members FROM industry_pools") + by_theme = {} + for _, r in pools.iterrows(): + ms = r["members"] + if isinstance(ms, str): + ms = json.loads(ms) + by_theme[str(r["theme"]).strip()] = ms or [] + excl = {str(x).strip() for x in (d.get("exclude_themes") or [])} + rows, missing = [], [] + for tr in d.get("tracks") or []: + if only_confirmed and tr.get("status") != "confirmed": + continue + for theme in tr.get("kg_themes") or []: + t = str(theme).strip() + if t in excl: + continue + if t not in by_theme: + missing.append((tr["name"], t)) + continue + for m in by_theme[t]: + ts = (m or {}).get("ts_code") + if ts: + rows.append((ts, m.get("name"), tr["name"], + tr.get("status"), t)) + df = (pd.DataFrame(rows, columns=["ts_code", "name", "track", + "status", "theme"]) + .drop_duplicates(["ts_code", "track"])) + return df, missing + + +def snapshot(only_confirmed: bool = True, path: str | None = None): + """成员表落 data/ 版本化快照(含 source_rule / updated_at,审计列)。""" + df, missing = resolve_members(only_confirmed, path) + os.makedirs("data", exist_ok=True) + out = f"data/track_members_{dt.date.today().isoformat()}.csv" + (df.assign(layer="", source_rule="pool_theme", + updated_at=dt.datetime.now().isoformat(timespec="seconds")) + .to_csv(out, index=False)) + return out, df, missing + + +def coverage_report(path: str | None = None) -> pd.DataFrame: + """赛道覆盖体检(G1 清单的首件事):逐赛道成员数与主题命中率,含 candidate。""" + df, missing = resolve_members(only_confirmed=False, path=path) + d = load_yml(path) + print("赛道覆盖体检(industry_pools 当前态;成员数已去重):") + for tr in d.get("tracks") or []: + sub = df[df["track"] == tr["name"]] + n_theme = len(tr.get("kg_themes") or []) + miss = [t for name, t in missing if name == tr["name"]] + tag = "" if tr.get("status") == "confirmed" else "(candidate)" + line = (f" {tr['name']}{tag}: 成员 {sub['ts_code'].nunique()} 只" + f" | 主题命中 {sub['theme'].nunique()}/{n_theme}") + if miss: + line += f" | 未命中: {'、'.join(miss)}" + print(line) + conf = df[df["status"] == "confirmed"] + print(f" —— confirmed 合计(去重): {conf['ts_code'].nunique()} 只") + return df