"""四路子因子构造。输入日期区间 [start, end](YYYY-MM-DD),输出 DataFrame[trade_date, stock_code, factor_value]。stock_code 输出形态不限, common.write_factor 统一转前缀式。 覆盖范围由 config.SUBFACTOR_UNIVERSE 控制(pool=池内 / market=全市场,见 §4 评审); 传导不受该开关影响——它天生就是池内语义。 建模参数(EVENT_POLARITY / *_HALF_LIFE / *_WINDOW)是**因子决策**,见设计文档 §5.3 与 §9-5,此处取草案默认,待用户确认后调。 """ import numpy as np import pandas as pd import config import common import db FACTORS = { "akg_upside": "t_factor_akg_upside", "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 确认)---- EVENT_POLARITY = { "股份回购": 1.0, "重大合同中标": 1.0, "股权激励授予": 0.5, "诉讼仲裁": -1.0, "行政处罚": -1.0, "股权质押": -0.5, "发行上市": 0.0, "并购交割": 0.0, "other": 0.0, # 增减持 / 业绩预告:符号取决于 direction(下 EVENT_DIR) "增减持": 0.0, "业绩预告": 0.0, } EVENT_DIR = {"预增": 1.0, "预减": -1.0, "增持": 1.0, "减持": -1.0} EVENT_HALF_LIFE = 10 # 交易日 EVENT_WINDOW = 60 # 交易日(超窗不计) _EMPTY = pd.DataFrame(columns=["trade_date", "stock_code", "factor_value"]) # ---------------------------------------------------------------- 热度 def build_heat(start, end): """热度 = stock_fund_heat_scores 最新批次 score(0~1)。stock_code 已前缀式。""" df = db.read_mysql("heat", """SELECT s.trade_date, s.stock_code, s.score AS factor_value FROM stock_fund_heat_scores s JOIN (SELECT trade_date, MAX(batch_no) bn FROM stock_fund_heat_scores WHERE trade_date BETWEEN %s AND %s GROUP BY trade_date) m ON m.trade_date = s.trade_date AND m.bn = s.batch_no""", (start, end)) if df.empty: return _EMPTY df["stock_code"] = df["stock_code"].astype(str).str.strip() df["factor_value"] = pd.to_numeric(df["factor_value"], errors="coerce") df = common.universe_filter(df, col="stock_code", as_prefix=True) return df[["trade_date", "stock_code", "factor_value"]] def _price_code_col() -> str: """gp_day_data 的代码列名(平台实测 = symbol,非 ts_code)。 用 LIMIT 1 探列名,而不是拿全区间查询去试错——原来那种写法在 history 模式下 第一次试探就会拉一次全区间数据。""" cands = [config.PRICE_CODE_COL] + [c for c in ("symbol", "ts_code") if c != config.PRICE_CODE_COL] last = None for c in cands: try: db.read_mysql("price", f"SELECT `{c}` FROM gp_day_data LIMIT 1") print(f" (现价用 gp_day_data.{c})") return c except Exception as e: # noqa: BLE001 —— 列名不对就换下一个候选 last = e raise RuntimeError(f"gp_day_data 代码列都不行(试了 {cands}): {last!r}") def _read_gp_price(start, end): """gp_day_data 现价,**按月分块**读取。 原来一次拉全区间:`--mode history --start 2006-01-01` 会把千万级行拉进 pandas。 不在 SQL 里按代码过滤——代码形态(600000.SH / SH600000 / 600000)两边不一致, SQL 侧过滤容易全空且难排查,统一折前缀式后在 pandas 侧过滤。 """ col = _price_code_col() parts, cur = [], pd.Timestamp(start) endts = pd.Timestamp(end) step = max(1, config.PRICE_CHUNK_DAYS) while cur <= endts: hi = min(cur + pd.Timedelta(days=step - 1), endts) parts.append(db.read_mysql( "price", f"SELECT `timestamp` AS trade_date, `{col}` AS ts_code, close " f"FROM gp_day_data WHERE `timestamp` BETWEEN %s AND %s", (cur.date().isoformat(), hi.date().isoformat()))) cur = hi + pd.Timedelta(days=1) if not parts: return pd.DataFrame(columns=["trade_date", "ts_code", "close"]) return pd.concat(parts, ignore_index=True) # ---------------------------------------------------------------- 预期空间 def build_upside(start, end): """upside = 一致预期目标价中枢 / 当日现价 − 1(as-of:现价日取 asof<=当日最新一致预期)。 ⚠️ 两个已知口径特征(评审 §6.3,不是 bug,但读数时要知道): · consensus_daily 的 target_mid_avg 是**90 天内全部研报行的简单平均** (不按机构去重、不按时间加权),且实测 max_price 非空仅 0.4%、min_price 31% —— 目标价中枢主要由单值目标价构成; · gp_day_data 是**前复权**、锚在最新日,每次除权历史 close 会被整体重写, 所以 upside 的历史值不可复现 —— 用 freeze.py 冻结当时用到的 close。 """ cons = db.read_pg( "SELECT ts_code, asof_date, target_mid_avg FROM v_factor_consensus " "WHERE asof_date <= %s", (end,)) if cons.empty: return _EMPTY cons["k"] = cons["ts_code"].map(common.to_prefix) # 注:consensus_daily 本就只对池成员聚合(基座 _pool_ts_codes), # 故 market 模式下 upside 覆盖不会真的变宽——这里过滤只为口径一致。 cons = common.universe_filter(cons, col="k", as_prefix=True) if cons.empty: return _EMPTY price = _read_gp_price(start, end) if price.empty: return _EMPTY price["close"] = pd.to_numeric(price["close"], errors="coerce") price["k"] = price["ts_code"].map(common.to_prefix) price = price[(price["close"] > 0)].dropna(subset=["close"]) cons = cons[cons["k"].isin(set(price["k"]))] if cons.empty: return _EMPTY # 强制两侧键同分辨率 datetime64[ns](pandas 2.x 不同来源可能 us/ns 混, # merge_asof 会报 incompatible merge keys) cons["asof_date"] = pd.to_datetime(cons["asof_date"]).astype("datetime64[ns]") price["trade_date"] = pd.to_datetime(price["trade_date"]).astype("datetime64[ns]") left = price[["trade_date", "k", "close"]].sort_values("trade_date") right = cons[["asof_date", "k", "target_mid_avg"]].sort_values("asof_date") m = pd.merge_asof(left, right, left_on="trade_date", right_on="asof_date", by="k", direction="backward") # 每股取 asof<=当日最新目标价 m = m.dropna(subset=["target_mid_avg", "close"]) if m.empty: return _EMPTY m["factor_value"] = m["target_mid_avg"].astype(float) / m["close"] - 1.0 m = m.rename(columns={"k": "stock_code"}) return m[["trade_date", "stock_code", "factor_value"]] # ---------------------------------------------------------------- 事件 def _polarity(event_type, direction): d = (direction or "").strip() if d in EVENT_DIR: return EVENT_DIR[d] return EVENT_POLARITY.get(event_type, 0.0) def build_event(start, end): """事件分 = Σ 近窗口内事件 极性 × 时间衰减(exp(-交易日龄·ln2/半衰期))。 ts_code 取文档锚(v_factor_events 已解析)。无事件的股当天不出行(= 缺 → 合成侧填 0)。 """ look = (pd.Timestamp(start) - pd.Timedelta(days=EVENT_WINDOW * 2)).date() try: ev = db.read_pg( "SELECT ts_code, disclosure_date, event_type, direction, " " confidence, doc_id, source_type " "FROM v_factor_events WHERE disclosure_date BETWEEN %s AND %s", (look, end)) except Exception as e: # noqa: BLE001 —— 视图还是 v1(无 doc_id/source_type)时退回 print(f" (v_factor_events 无 doc_id/source_type,退回旧列——建议先更新视图: {e!r})") ev = db.read_pg( "SELECT ts_code, disclosure_date, event_type, direction " "FROM v_factor_events WHERE disclosure_date BETWEEN %s AND %s", (look, end)) if ev.empty: return _EMPTY ev = common.universe_filter(ev, col="ts_code", as_prefix=False).copy() if ev.empty: return _EMPTY # ---- 年报污染防护(评审 §6.5)------------------------------------------ # v_factor_events 的锚 documents.meta->>'company_ts_code' 年报同样有,而一份年报 # 能抽十几条 EVENT,且含**历史**诉讼/处罚 —— 会在年报披露日形成巨大负值尖峰。 # 基座 hotspot._pick_event_anomalies 为此专门做了防刷屏(other 不进 / 同主体同 # 类型只取最新 / 单主体≤2),桥侧原来零保护。三道,从强到弱: if "source_type" in ev.columns: keep = ev["source_type"].astype(str).isin(config.EVENT_SOURCE_TYPES) if (~keep).any(): drop_by = ev.loc[~keep, "source_type"].value_counts().to_dict() print(f" (事件:按 source_type 剔除 {int((~keep).sum())} 条 {drop_by})") ev = ev[keep] if ev.empty: print(" ⚠️ 按 source_type 白名单过滤后为空——核对 EVENT_SOURCE_TYPES") return _EMPTY if "confidence" in ev.columns: # 同键取置信度最高的一条 ev = ev.sort_values("confidence", ascending=False, na_position="last") ev = ev.drop_duplicates(["ts_code", "event_type", "direction", "disclosure_date"]) if "doc_id" in ev.columns: # 单文档封顶 before = len(ev) ev = ev.groupby("doc_id", group_keys=False).head(config.EVENT_MAX_PER_DOC) if len(ev) < before: print(f" (事件:单文档封顶 {config.EVENT_MAX_PER_DOC} 条,剔除 {before - len(ev)} 条)") ev["event_type"] = ev["event_type"].fillna("") ev["direction"] = ev["direction"].fillna("") # 多数事件无 direction(NULL→NaN) ev["pol"] = [_polarity(t, d) for t, d in zip(ev["event_type"], ev["direction"])] ev = ev[ev["pol"] != 0.0] if ev.empty: return _EMPTY cal = common.trading_days(start, end) if not cal: print(f" ⚠️ {start}~{end} 无交易日 → 事件因子空转" f"(这是「日历为空」,不是「没有事件」)") return _EMPTY cal = pd.DatetimeIndex(cal) cal_i = cal.values.astype("datetime64[ns]").astype("int64") # (D,) decay = np.log(2) / EVENT_HALF_LIFE rows = [] for ts, g in ev.groupby("ts_code"): disc = pd.to_datetime(g["disclosure_date"]).values \ .astype("datetime64[ns]").astype("int64") # (E,) pol = g["pol"].to_numpy(dtype=float) # 自然日龄 → 交易日龄近似 ×(5/7)(v1 近似,见 README 待优化项) age = (cal_i[:, None] - disc[None, :]) / 86_400e9 * (5.0 / 7.0) # (D,E) m = (age >= 0) & (age <= EVENT_WINDOW) if not m.any(): continue age_safe = np.where(m, age, 0.0) # 先夹再 exp,防 exp(超大正数) 溢出 vals = np.where(m, pol[None, :] * np.exp(-age_safe * decay), 0.0).sum(axis=1) for i in np.nonzero(vals)[0]: rows.append((cal[i].date(), ts, float(vals[i]))) return (pd.DataFrame(rows, columns=["trade_date", "stock_code", "factor_value"]) if rows else _EMPTY) # ---------------------------------------------------------------- 传导 def build_transmission(start, end): """传导分 = 指向该股所在环节的 **distinct 源数** ×(1 − 已动比例); 同股同日多候选取最大。 口径修正(评审 硬伤2):原用 n_paths = jsonb_array_length(paths),但 · graph_store.cascade() 的变长边 `*1..N` **每种长度各返回一条路径** (A→B 与 A→X→B 同时出现、末节点都是 B),跨路径不去重; · graph_store.transmission_targets() 的 updown/supply/drives 三桶之间也不去重。 于是同一 source 对同一 target 重复计入,而 n_paths 是 factor_value 的主量级。 改用 distinct source 数,也更贴合传导模块自述的语义「多源汇聚 = 传导逻辑更硬」。 """ cols_new = ("scan_date, target, ts_code, n_sources, n_paths, moved_ratio, " "members_total, moved, n_quiet_stored, mkt_trade_date") try: tr = db.read_pg( f"SELECT {cols_new} FROM v_factor_transmission " f"WHERE scan_date BETWEEN %s AND %s", (start, end)) col_val = "n_sources" except Exception as e: # noqa: BLE001 —— 视图还是 v1 时退回,但明确告警 print(f" ⚠️ v_factor_transmission 还是旧版(无 n_sources)——退回 n_paths," f"该口径把重复路径计入了强度,请尽快更新视图: {e!r}") tr = db.read_pg( "SELECT scan_date, ts_code, n_paths, moved_ratio " "FROM v_factor_transmission WHERE scan_date BETWEEN %s AND %s", (start, end)) col_val = "n_paths" if tr.empty: return _EMPTY uni = common.load_universe() # 传导恒按池过滤(天生池内语义) tr = tr[tr["ts_code"].isin(uni)].copy() if tr.empty: return _EMPTY _warn_upstream_truncation(tr) tr["factor_value"] = (pd.to_numeric(tr[col_val], errors="coerce").fillna(0.0) * (1.0 - pd.to_numeric(tr["moved_ratio"], errors="coerce").fillna(0.0))) g = (tr.groupby(["scan_date", "ts_code"])["factor_value"].max().reset_index() .rename(columns={"scan_date": "trade_date", "ts_code": "stock_code"})) return g[["trade_date", "stock_code", "factor_value"]] def _warn_upstream_truncation(tr: pd.DataFrame) -> None: """把上游「静默失真」变成显式告警——07-28 按基座二批后的新制度重写判据。 基座二批后:topic_context cap=1000、scan 大主题闸=600、quiet 全量落库、 mkt_trade_date 落库。旧判据(members_total>=30 撞 cap30 / n_quiet_stored>=12 撞截断)在新制度下全是假阳性(07-27 实测 6/7 条误报),废弃。 采信规则(stale 剔除/打折/标记)待三批-5 拍板,当前全量采信、只告警。 """ if "mkt_trade_date" in tr.columns and tr["mkt_trade_date"].notna().any(): bad = tr[tr["mkt_trade_date"].astype(str) != tr["scan_date"].astype(str)] if not bad.empty: print(f" ⚠️ {bad['scan_date'].nunique()} 个 scan_date 的 movers 快照日与 " f"scan_date 不符(17:30 sync_market 晚点)——这些日的传导项不可信") if "members_total" in tr.columns: hit = tr[pd.to_numeric(tr["members_total"], errors="coerce") > config.UPSTREAM_MEMBER_CAP] if not hit.empty: n = hit["target"].nunique() if "target" in hit else len(hit) print(f" ⚠️ {n} 个环节 members_total > {config.UPSTREAM_MEMBER_CAP}" f"(基座大主题闸值)——基座制度回退或闸失效,查 transmission.scan") if {"n_quiet_stored", "members_total", "moved"} <= set(tr.columns): mt = pd.to_numeric(tr["members_total"], errors="coerce") mv = pd.to_numeric(tr["moved"], errors="coerce") nq = pd.to_numeric(tr["n_quiet_stored"], errors="coerce") bad = tr[(mt - mv) != nq] if not bad.empty: n = bad["target"].nunique() if "target" in bad else len(bad) print(f" ⚠️ {n} 个候选 n_quiet_stored ≠ members_total − moved ——" f"quiet 截断复活或台账/视图不一致(07-27 新制度起不该发生;" 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) _RISK_PREFIXES = ("*ST", "ST", "SST", "S*ST", "退市") def _risk_set() -> set: """重大风险名单(07-31 拍板:默认杜绝):按证券简称识别 ST/*ST/退市族。 名称来自 industry_pools 成员表;开关 config.EXCLUDE_RISK_NAMES(默认开)。 加载失败只告警不中止——宁可当天不过滤,也不让计划断产。""" if not config.EXCLUDE_RISK_NAMES: return set() import json try: pools = db.read_pg("SELECT members FROM industry_pools") except Exception as e: # noqa: BLE001 print(f" ⚠️ 风险名单加载失败(本轮不过滤): {e!r}") return set() out = set() for _, r in pools.iterrows(): ms = r["members"] if isinstance(ms, str): ms = json.loads(ms) for m in ms or []: ts = (m or {}).get("ts_code") name = str((m or {}).get("name") or "").replace(" ", "") if ts and (name.startswith(_RISK_PREFIXES) or name.endswith("退")): out.add(common.to_prefix(ts)) if out: print(f" (重大风险闸: {len(out)} 只 ST/退市族挡在档位外," f"EXCLUDE_RISK_NAMES=0 可关)") return out 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, risk_set=None) -> pd.Series: """三档(07-30 拍板): 2 = 主榜:有业绩锚(券商覆盖且 upside>=0);C 闸开启后再交赛道成员; 1 = 观察档:无业绩锚但有产业链锚(当前=在传导链上;C 转正后并入赛道成员); 0 = 不采纳:两锚皆无,有覆盖但 upside<0(贵了不买是绝对下限), 或命中重大风险闸(ST/退市族,07-31 拍板默认杜绝)。""" 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 if risk_set: g[p.index.isin(risk_set)] = 0.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() rs = _risk_set() rows = [] for ds in _days_of(start, end): p = _panel(ds) if p.empty: continue g = _gate_of(p, ts, rs) 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 + 传导档位×20 + 组内分:传导票按 log 传导分的组内中位数分成 强(2)/弱(1)两档,无传导 0 档;组内分 = 0.6·z(−热度) + 0.4·z(upside), 即"还没热、还便宜"。观察档 = 100 + 0.6·z(传导) + 0.4·z(−热度)——无估值锚, 低置信。**档位步长 20、组内分夹 ±9.9**:组内分全幅 19.8 小于步长, 先档后分的词典序才真正成立(步长 10 会档间重叠,07-30 交付当天修正)。 z 都在各自档内算;缺失处置同设计 §3.4:传导缺=0,热度缺=档内中位数。""" ts = _track_set() rs = _risk_set() out = [] for ds in _days_of(start, end): p = _panel(ds) if p.empty: continue g = _gate_of(p, ts, rs) 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 * 20.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, }