692 lines
34 KiB
Markdown
692 lines
34 KiB
Markdown
# 评审问题的修复方案(2026-07-26)
|
||
|
||
> 按依赖顺序分三批。**第一批**(桥侧 + 视图)不依赖任何决策、不动基座,今天就能跑;
|
||
> **第二批**(基座侧)是确定性缺陷修复,需要动 astock-kg 代码;
|
||
> **第三批**(赛道门槛 C 与权重结构)需要你先拍板,我给出两种拍法各自的代码形态。
|
||
>
|
||
> 附件:`astock_kg_slot_views_v2.sql`(完整替换)、`freeze.py`(新文件,直接放桥根目录)。
|
||
> 所有改动都保持"桥只连三库、基座不算因子"的边界。
|
||
|
||
---
|
||
|
||
## 0. 数据安全性核对:已抽取的成果一份都不会白费
|
||
|
||
已逐表核实(`db/postgres/init.sql` + 全仓 grep),结论先行:
|
||
**本方案的全部改动都不触碰 `claims` / `documents`,不需要 replay,不需要重抽。**
|
||
|
||
### 0.1 每一处改动写了什么
|
||
|
||
| 改动 | 写入对象 | 对已有数据的影响 |
|
||
|---|---|---|
|
||
| 1.1 视图替换 | 视图(不含数据) | 无 |
|
||
| 1.1 `ALTER ... ADD COLUMN IF NOT EXISTS mkt_trade_date` | `transmission_candidates` 加可空列 | 无(已存行取 NULL) |
|
||
| 1.1 `CREATE INDEX IF NOT EXISTS` | 索引 | 无 |
|
||
| 1.2–1.8 桥侧补丁 | 平台 MySQL `t_factor_*` | 只动平台侧派生表,可随时重算 |
|
||
| 1.9 `freeze.py` | 桥工程 `data/frozen/` | 纯新增文件 |
|
||
| 2.1 `topic_context` 加 `ORDER BY` | **只读 Cypher** | 无 |
|
||
| 2.2 `transmission.scan` | `transmission_candidates`(按 `scan_date` 删后插的日台账) | 只影响新扫描日 |
|
||
| 2.3 `hotspot._latest_mkt` 加参数 | `hotspot_candidates`(同为日台账) | 只影响新扫描日 |
|
||
| 2.4 `sync_consensus(asof)` | `consensus_daily`(主键 `ts_code, asof_date`) | 回填是**新增** asof 行;不覆盖历史 |
|
||
| 2.5 `industry_pools_history` | 新表 | 纯新增 |
|
||
|
||
**没有一处 DELETE / UPDATE / TRUNCATE 落在抽取产物上。** 全仓 grep 也确认:
|
||
`claims` / `documents` / `mkt_daily` / `company_master` 全库没有任何删除或清空语句
|
||
(`claims` 的 DDL 注释就写着 `---- 断言表:只追加 ----`,并有
|
||
`CONSTRAINT uq_claim_dedup UNIQUE (dedup_key)` 保证幂等)。
|
||
|
||
### 0.2 为什么架构上就不会白费
|
||
|
||
抽取的产物是 `claims` + `documents`,它们是**事实源**;Neo4j 里的 belief 图是**投影**。
|
||
`replay_all()` 的 docstring 说得很明白:"从 claim 层全量重建 belief 图:
|
||
只跑融合规则,**不重新抽取、不调用任何 LLM**。因此确定、便宜、稳定。"
|
||
|
||
所以就算将来真要动本体/闸门/融合规则,代价也只是一次 replay(12297 条约 80 秒),
|
||
claims 无损。真正昂贵的 re-extraction 是另一件事,`claims.extractor_model` 这一列
|
||
就是专门为"按模型版本圈定重抽范围"留的——架构本来就是为这种演进设计的。
|
||
|
||
**而本方案连 replay 都不需要**:2.1 是只读查询,2.2–2.5 全部写 PG 应用层表,
|
||
没有一处改动落在融合规则或 Neo4j 投影路径上。
|
||
|
||
### 0.3 需要你知道的三点(不是"白费",但要说清)
|
||
|
||
1. **已存的 3 天传导台账修不"好",但也不用修。**
|
||
`transmission_candidates` 现有的 07-11~07-23 三天,`quiet` 是被 `[:12]` 截断的、
|
||
`moved_ratio` 建立在任意 ≤30 抽样上。修法上线后新扫描才正确。
|
||
**好消息**:`mkt_daily` 是持久快照表(`trade_date` 在主键里、全库无删除),
|
||
所以历史 movers 其实**是可重建的** —— 2.2/2.3 改完之后,
|
||
`transmission.scan(scan_date=D, allow_stale=True)` 能对任意有 mkt 快照的 D 重扫。
|
||
(这也顺带修正了 v1.1 文档里"历史 movers 难重建"的判断——难的不是 movers,
|
||
是当时的**图谱状态**。)
|
||
**但我建议不要重扫这三天**:会用今天的图谱去套那三天,而 44079 份公告正在
|
||
以 200 份/30 分钟入库,图谱这两周变化不小 —— 那就是标准的成员性前视。
|
||
这三天目前也没被任何东西消费,直接从修好之后重新起算最干净。
|
||
|
||
2. **新鲜度断言会让传导"少产出",这是有意的。**
|
||
2.2 加的断言在 `mkt_daily` 最新 stock 快照 ≠ `scan_date` 时**中止扫描**。
|
||
在公告批量入库期间,如果 beat 被饿到(就是 `announcement-corpus-bulk-load`
|
||
记的那个单队列风险),17:30 快照晚点 → 当天传导直接没有产出。
|
||
这比"用 T−1 的 movers 生成带今日戳的候选"好,但你要预期到它会偶发触发,
|
||
也意味着 live 传导史积累得慢一点。队列分离那次修复的实机验证,
|
||
现在多了一个理由要尽快做完。
|
||
|
||
3. **唯一可能产生 LLM 花费的是"提环节覆盖",但它也不浪费任何东西。**
|
||
如果 G1 体检发现赛道覆盖太薄,对策是 `segment_backfill`(环节专项遍存量重抽)。
|
||
它零新语料、`dedup_key` 幂等、只追加不删旧 claims —— 花的是 brain 调用,
|
||
不是把已有成果作废。
|
||
|
||
### 0.4 落地前后的自查(跑一遍留个数,最稳)
|
||
|
||
```bash
|
||
# 改动前后各跑一次,两次输出应完全一致(除 transmission_candidates 的列数 +1)
|
||
docker exec -i akg-postgres psql -U akg -d akg <<'SQL'
|
||
SELECT 'claims' t, count(*) n, max(ingestion_date)::date latest FROM claims
|
||
UNION ALL SELECT 'documents', count(*), max(ingested_at)::date FROM documents
|
||
UNION ALL SELECT 'entity_links', count(*), max(created_at)::date FROM entity_links
|
||
UNION ALL SELECT 'company_master', count(*), max(synced_at)::date FROM company_master
|
||
UNION ALL SELECT 'consensus_daily', count(*), max(asof_date) FROM consensus_daily
|
||
UNION ALL SELECT 'mkt_daily', count(*), max(trade_date) FROM mkt_daily
|
||
UNION ALL SELECT 'industry_pools', count(*), max(refreshed_at)::date FROM industry_pools;
|
||
|
||
-- 顺带回答一个对回填很关键的问题:mkt_daily 的 stock 切片到底存了多少天?
|
||
-- (它决定了传导"理论上"能重扫到哪一天,也决定了 §7 回填的真实边界)
|
||
SELECT kind, min(trade_date), max(trade_date), count(DISTINCT trade_date) days
|
||
FROM mkt_daily GROUP BY kind ORDER BY kind;
|
||
SQL
|
||
```
|
||
|
||
最后那条查询的结果请一并发我 —— 如果 `kind='stock'` 的天数明显多于 3 天,
|
||
设计文档 §7 里"传导只 live 累积"这一条的**边界可以往前推**(图谱漂移仍在,
|
||
但至少多了一个"用当时 movers + 今天图谱"的、带明确标注的近似区间可选)。
|
||
|
||
---
|
||
|
||
## 第一批 · 桥侧 + 视图(零决策,今天可落)
|
||
|
||
### 1.1 `sql/astock_kg_slot_views.sql` → 用附件 `astock_kg_slot_views_v2.sql` 整体替换
|
||
|
||
三处变更:`n_paths` → `n_sources`(修重复计数);暴露 `target/members_total/moved/
|
||
n_quiet_stored/mkt_trade_date`(让桥能自检上游截断与快照新鲜度);`v_factor_events`
|
||
补 `doc_id/source_type/tier`(供年报封顶)。文件里带一段幂等 `ALTER TABLE`,一起跑即可。
|
||
|
||
```bash
|
||
docker exec -i akg-postgres psql -U akg -d akg < sql/astock_kg_slot_views.sql
|
||
```
|
||
|
||
**数据安全性**:只建视图 + 加一个可空列 + 加索引,**不写不改不删任何一行数据**。
|
||
`claims` / `documents` 完全不动,无需 replay、无需重抽。
|
||
|
||
✅ 列名已核实(不用再探):`documents.source_type`,取值
|
||
`annual_report / research_report / announcement / news / interactive_qa`
|
||
(`db/postgres/init.sql:12`)。`claims` 有 `tier / confidence / dedup_key / doc_id /
|
||
subject_norm`(`init.sql:22-58, 263`),`entity_links(alias_norm, canonical_id, entity_type)`、
|
||
`company_master(ts_code, short_name)` 均存在(`init.sql:230-251`)——
|
||
第三批 3.1 的救急 SQL 依赖的列全部对得上。
|
||
|
||
### 1.2 `common.py` — 交易日历改用行情表(**修硬伤 3**)
|
||
|
||
```python
|
||
def trading_days(start: str, end: str) -> list:
|
||
"""目标区间交易日历。
|
||
|
||
⚠️ 曾用热度表(stock_fund_heat_scores) —— 它最早只有 2026-03-26 前后,
|
||
导致 `build akg_event --mode history --start 2024-01-01` 静默返回空
|
||
(cal 为空 → _EMPTY → 只打印"无数据(跳过)",不报错)。
|
||
改用行情表 gp_day_data(5584 天,覆盖全历史)。"""
|
||
df = db.read_mysql(
|
||
"price",
|
||
"SELECT DISTINCT `timestamp` AS trade_date FROM gp_day_data "
|
||
"WHERE `timestamp` BETWEEN %s AND %s ORDER BY `timestamp`",
|
||
(start, end))
|
||
if df.empty:
|
||
print(f" ⚠️ 交易日历为空(gp_day_data 在 {start}~{end} 无数据)"
|
||
f"——上层会跳过该区间,请核对区间与行情表覆盖")
|
||
return list(pd.to_datetime(df["trade_date"]))
|
||
```
|
||
|
||
### 1.3 `common.py` — 分块提交(**修放量会炸**)
|
||
|
||
```python
|
||
def write_factor(table: str, df: pd.DataFrame, mode: str = "daily") -> None:
|
||
"""df[trade_date, stock_code, factor_value] → 幂等写因子表。
|
||
|
||
分块提交(评审):热度全史约 322×1674≈54 万行,原来单事务 executemany
|
||
走 ShardingSphere 代理有风险。改成 DELETE 一个事务 + INSERT 按块提交。
|
||
代价是中途失败会留下部分区间——但整个写入按区间幂等,重跑即修复。"""
|
||
if df is None or df.empty:
|
||
print(f" {table}: 无数据(跳过)")
|
||
return
|
||
df = df.dropna(subset=["trade_date", "stock_code", "factor_value"]).copy()
|
||
df["stock_code"] = df["stock_code"].map(to_prefix)
|
||
df["trade_date"] = pd.to_datetime(df["trade_date"]).dt.date
|
||
df = df.drop_duplicates(["trade_date", "stock_code"], keep="last")
|
||
if df.empty:
|
||
print(f" {table}: 清洗后无数据")
|
||
return
|
||
dmin, dmax = df["trade_date"].min(), df["trade_date"].max()
|
||
ensure_table(table)
|
||
rows = list(df[["trade_date", "stock_code", "factor_value"]]
|
||
.itertuples(index=False, name=None))
|
||
chunk = config.WRITE_CHUNK_ROWS
|
||
with db.factor_conn() as conn:
|
||
with conn.cursor() as cur:
|
||
cur.execute(f"DELETE FROM {table} WHERE trade_date BETWEEN %s AND %s",
|
||
(dmin, dmax))
|
||
conn.commit()
|
||
for i in range(0, len(rows), chunk):
|
||
with conn.cursor() as cur:
|
||
cur.executemany(
|
||
f"INSERT INTO {table} (trade_date, stock_code, factor_value) "
|
||
f"VALUES (%s,%s,%s)", rows[i:i + chunk])
|
||
conn.commit()
|
||
if len(rows) > chunk:
|
||
print(f" ...{min(i + chunk, len(rows))}/{len(rows)}")
|
||
print(f" {table}: 写入 {len(rows)} 行, 日期 {dmin}~{dmax}")
|
||
```
|
||
|
||
顶部加 `import config`。
|
||
|
||
### 1.4 `common.py` — universe 过滤改成可切换(**评审 §4**)
|
||
|
||
```python
|
||
def universe_filter(df: pd.DataFrame, col: str = "stock_code",
|
||
as_prefix: bool = True) -> pd.DataFrame:
|
||
"""按 config.SUBFACTOR_UNIVERSE 决定子因子是否受覆盖池限制。
|
||
|
||
默认 'pool'(保持现状)。**建议改 'market'**:industry_pools 只存最新态、
|
||
每周一 refresh_pools 自动长大,用它过滤历史子因子会引入成员性前视
|
||
(§2.2 批传导用的同一条论证,上升一层),且 z 统计量每周一结构性跳变。
|
||
子因子表的定位是"可独立观察的仪表",过滤该发生在消费端(akg_score)而不是这里。"""
|
||
if config.SUBFACTOR_UNIVERSE != "pool":
|
||
return df
|
||
uni = load_universe()
|
||
if as_prefix:
|
||
uni = {to_prefix(x) for x in uni}
|
||
return df[df[col].astype(str).str.strip().isin(uni)]
|
||
```
|
||
|
||
`build_heat` / `build_upside` / `build_event` 里的 `uni = ...` + `isin(uni)` 都换成调它。
|
||
`build_transmission` **不换** —— 传导天生就是池内语义。
|
||
|
||
### 1.5 `config.py` — 追加四个旋钮
|
||
|
||
```python
|
||
# 子因子是否受覆盖池限制:pool(现状)| market(评审建议,避免成员性前视)
|
||
SUBFACTOR_UNIVERSE = os.environ.get("SUBFACTOR_UNIVERSE", "pool").lower()
|
||
# 行情读取按月分块(防 `--start 2006` 把千万行拉进 pandas)
|
||
PRICE_CHUNK_DAYS = int(os.environ.get("PRICE_CHUNK_DAYS", "31"))
|
||
# 因子表写入分块行数
|
||
WRITE_CHUNK_ROWS = int(os.environ.get("WRITE_CHUNK_ROWS", "50000"))
|
||
# 单文档最多贡献几条事件(年报能抽十几条,见评审 §6.5)
|
||
EVENT_MAX_PER_DOC = int(os.environ.get("EVENT_MAX_PER_DOC", "3"))
|
||
# 事件只认哪些来源(documents.source_type)。默认只认公告——年报里的"历史诉讼"
|
||
# 会被记成披露日的当日负面事件,是 S2 否决闸最危险的假信号来源。
|
||
# 想放宽就填 "announcement,annual_report"。
|
||
EVENT_SOURCE_TYPES = {
|
||
s.strip() for s in
|
||
os.environ.get("EVENT_SOURCE_TYPES", "announcement").split(",") if s.strip()}
|
||
```
|
||
|
||
`.env.example` 同步加这四行(都给默认值,可不填)。
|
||
|
||
### 1.6 `factors.py` — 行情按月分块(**修放量会炸**)
|
||
|
||
```python
|
||
def _read_gp_price(start, end):
|
||
"""gp_day_data 现价。两处修正(评审):
|
||
① 按月分块 —— 原来一次拉全区间,`--mode history --start 2006-01-01`
|
||
会把千万级行拉进 pandas;
|
||
② 不在 SQL 里按代码过滤 —— 代码形态(600000.SH / SH600000 / 600000)
|
||
两边不一致,SQL 过滤容易全空且难排查,改在 pandas 侧折前缀后过滤。"""
|
||
cands = [config.PRICE_CODE_COL] + [c for c in ("symbol", "ts_code")
|
||
if c != config.PRICE_CODE_COL]
|
||
col, last = None, None
|
||
for c in cands: # 先用 LIMIT 1 探列名,别用全区间去试错
|
||
try:
|
||
db.read_mysql("price", f"SELECT `{c}` FROM gp_day_data LIMIT 1")
|
||
col = c
|
||
break
|
||
except Exception as e: # noqa: BLE001
|
||
last = e
|
||
if col is None:
|
||
raise RuntimeError(f"gp_day_data 代码列都不行(试了 {cands}): {last!r}")
|
||
print(f" (upside 现价用 gp_day_data.{col})")
|
||
|
||
parts, cur = [], pd.Timestamp(start)
|
||
endts = pd.Timestamp(end)
|
||
while cur <= endts:
|
||
hi = min(cur + pd.Timedelta(days=config.PRICE_CHUNK_DAYS - 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)
|
||
return (pd.concat(parts, ignore_index=True) if parts
|
||
else pd.DataFrame(columns=["trade_date", "ts_code", "close"]))
|
||
```
|
||
|
||
### 1.7 `factors.py` — 传导改用 `n_sources` + 上游截断自检(**修硬伤 2**)
|
||
|
||
```python
|
||
def build_transmission(start, end):
|
||
"""传导分 = 指向该股所在环节的 **distinct 源数** ×(1 − 已动比例);
|
||
同股同日多候选取最大。
|
||
|
||
口径修正(评审 硬伤2):原用 n_paths = jsonb_array_length(paths),
|
||
但 cascade() 的变长边 *1..N 会把 A→B 与 A→X→B 各返回一条,
|
||
transmission_targets() 的 updown/supply/drives 三桶之间也不去重,
|
||
于是同一 source 对同一 target 重复计入。改用 distinct source 数,
|
||
也更贴合传导模块自述的"多源汇聚 = 传导逻辑更硬"。"""
|
||
tr = db.read_pg(
|
||
"SELECT scan_date, target, ts_code, n_sources, n_paths_raw, moved_ratio, "
|
||
" members_total, n_quiet_stored, mkt_trade_date "
|
||
"FROM v_factor_transmission WHERE scan_date BETWEEN %s AND %s",
|
||
(start, end))
|
||
if tr.empty:
|
||
return _EMPTY
|
||
|
||
# ---- 上游失真自检:把三个"静默失真"变成显式告警(评审 硬伤1、§6.2)----
|
||
if 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 快照晚点)——这些日的传导项不可信")
|
||
if (tr["members_total"] >= 30).any():
|
||
n = tr.loc[tr["members_total"] >= 30, "target"].nunique()
|
||
print(f" ⚠️ {n} 个环节 members_total>=30,撞上游 topic_context cap"
|
||
f"——moved_ratio 建立在任意 ≤30 抽样上(先修基座再放量)")
|
||
if (tr["n_quiet_stored"] >= 12).any():
|
||
n = tr.loc[tr["n_quiet_stored"] >= 12, "target"].nunique()
|
||
print(f" ⚠️ {n} 个环节 quiet 存满 12 条,撞 quiet[:12] 截断"
|
||
f"——真实未动成员更多,因子覆盖被展示逻辑锁住")
|
||
|
||
tr["factor_value"] = (tr["n_sources"].astype(float)
|
||
* (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"]]
|
||
```
|
||
|
||
### 1.8 `factors.py` — 事件的年报污染防护(**S2 前必修,现在顺手做**)
|
||
|
||
`build_event` 里读完 `ev` 之后、算 `pol` 之前插入:
|
||
|
||
```python
|
||
# ---- 年报污染防护(评审 §6.5)----
|
||
# v_factor_events 的锚 documents.meta->>'company_ts_code' 年报也有,而一份年报
|
||
# 能抽十几条 EVENT,且含**历史**诉讼/处罚 —— 会在年报披露日形成巨大负值尖峰。
|
||
# hotspot._pick_event_anomalies 为此专门做了防刷屏(other 不进 / 同主体同类型
|
||
# 去重 / 单主体≤2),桥侧原来零保护。三道,从强到弱:
|
||
if "source_type" in ev.columns:
|
||
# documents.source_type ∈ annual_report / research_report / announcement
|
||
# / news / interactive_qa (init.sql:12 已核实)
|
||
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 "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"])
|
||
# 单文档封顶:一份文档最多贡献 EVENT_MAX_PER_DOC 条
|
||
if "doc_id" in ev.columns:
|
||
ev = ev.groupby("doc_id", group_keys=False).head(config.EVENT_MAX_PER_DOC)
|
||
if ev.empty:
|
||
return _EMPTY
|
||
```
|
||
|
||
`build_event` 的 SQL 同步改成 `SELECT ts_code, disclosure_date, event_type, direction,
|
||
confidence, doc_id, doc_type FROM v_factor_events WHERE ...`。
|
||
|
||
另外 `cal` 为空时别静默返回(1.2 已在 `trading_days` 里加了告警,这里再补一句):
|
||
|
||
```python
|
||
cal = common.trading_days(start, end)
|
||
if not cal:
|
||
print(f" ⚠️ {start}~{end} 无交易日 → 事件因子空转(不是"没有事件")")
|
||
return _EMPTY
|
||
```
|
||
|
||
### 1.9 `freeze.py` — 新文件,直接放桥根目录(**评审 §5,最高优先级**)
|
||
|
||
见附件。这一条是**唯一有时间不可逆性的**:传导、池、Segment 边、前复权 close
|
||
都是过期即不可复原,`§10` 里"GRU 复活需 ≥1 年 live 传导史"这个计时器只在
|
||
开始存快照那天启动。它同时一次性满足判收标准里的"重跑不漂移"与"赛道成员可审计"。
|
||
|
||
`run.py` 三处小改:
|
||
|
||
```python
|
||
# ① 顶部 subparser 加
|
||
sub.add_parser("freeze").add_argument("--date")
|
||
|
||
# ② cmd_build 收集产物并在末尾冻结
|
||
def cmd_build(which, mode, start, end, date, do_freeze=True):
|
||
...
|
||
frames = {}
|
||
for code in codes:
|
||
...
|
||
df = factors.BUILDERS[code](start, end)
|
||
frames[f"factor_{code}"] = df
|
||
common.write_factor(factors.FACTORS[code], df, mode)
|
||
if do_freeze and mode == "daily":
|
||
import freeze
|
||
freeze.snapshot(start, extra_frames=frames)
|
||
|
||
# ③ dispatch
|
||
elif a.cmd == "freeze":
|
||
import freeze
|
||
freeze.snapshot(a.date)
|
||
```
|
||
|
||
`build` 子命令加一个 `--no-freeze`(回填时不必每天冻结)。
|
||
`docker-compose.yml` 的 `volumes` 已经挂了 `.:/app`,`data/frozen/` 天然落在宿主仓库里;
|
||
记得 `.gitignore` 里决定要不要跟踪(我建议**跟踪 manifest.json、忽略数据文件**,
|
||
manifest 才是可 diff 的审计线)。
|
||
|
||
### 1.10 `probe.py` — 加两节(G1 需要)
|
||
|
||
```python
|
||
def probe_corr():
|
||
"""池内三项相关矩阵(评审 §2)——若 corr(z_H, z_V) > 0.5,
|
||
"三项加权"实际是两项,§9-8 的权重讨论要重开。"""
|
||
_sec("corr · 池内 (z_T, z_H, z_V) 相关矩阵")
|
||
# 取最新可算日:传导有值的最近 scan_date
|
||
d = db.read_pg("SELECT max(scan_date) FROM v_factor_transmission").iloc[0, 0]
|
||
... # 三路各取当日截面 → 外连接 → 标准化 → .corr(method="spearman")
|
||
# 同时打印:|P| 规模、|P ∩ 传导>0|、传导为 0 的占比
|
||
```
|
||
|
||
`SECTIONS` 里注册 `"corr": probe_corr`,`run.py` 的 `--section` choices 加 `corr`。
|
||
(`tracks` 节等第三批的 yml 有草案了再加。)
|
||
|
||
---
|
||
|
||
## 第二批 · 基座侧(确定性缺陷修复,不算"新增因子计算")
|
||
|
||
> 这四处都是"让上游可复现/不撒谎",不涉及打分,不破 §4 铁律。
|
||
> `graph_store` 那处我是通过子代理读到的 Cypher,落地前请你对一眼实际代码。
|
||
|
||
### 2.1 `graph_store.topic_context` — 加 `ORDER BY`、cap 参数化(**修硬伤 1 的根**)
|
||
|
||
```cypher
|
||
MATCH (c)-[m:IN_SEGMENT {status:'active'}]->(t:Segment {segment_name:$id})
|
||
RETURN coalesce(c.name, c.ts_code, c.entity_key) AS name, c.ts_code AS ts_code
|
||
ORDER BY coalesce(c.ts_code, c.entity_key) -- ★ 新增:无序 → 同日重跑结果会变
|
||
LIMIT $cap
|
||
```
|
||
|
||
`cap` 默认从 30 提到 200(或让调用方传)。**这一条不改,硬伤 1 修不掉** ——
|
||
`transmission` 的 `members_total / moved_ratio` 会一直建立在任意 ≤30 抽样上。
|
||
|
||
### 2.2 `transmission.scan` — 存全量 quiet + 新鲜度断言 + 记快照日
|
||
|
||
```python
|
||
def scan(scan_date=None, annotate=True, max_candidates=12, member_cap=200):
|
||
from app.analysis import hotspot
|
||
from app.store import claim_store, graph_store
|
||
|
||
sd = date.fromisoformat(scan_date) if scan_date else date.today()
|
||
|
||
# ---- 0. 新鲜度记账(评审 §6.2)----
|
||
# ★ 用户决定(2026-07-26):**照常落库,只打告警 + 记 mkt_trade_date**,不中止。
|
||
# 理由:保覆盖。传导每日只有几十行,公告批量入库期间 beat 一旦被饿就整天空白,
|
||
# 代价太大;改由桥侧按 mkt_trade_date 自行判断要不要采信(factors.py 已实现告警)。
|
||
# 背景:movers/quiet 来自 mkt_daily 的 **最新** 快照,不是 sd 的快照。若 17:30
|
||
# sync_market 晚点,会用 T−1 的 movers 写成 scan_date=T;反过来重跑历史日
|
||
# 会用今天的快照 = 直接前视。落这一列后,下游至少看得见。
|
||
with get_pool().connection() as conn:
|
||
row = conn.execute(
|
||
"SELECT max(trade_date) FROM mkt_daily WHERE kind='stock'").fetchone()
|
||
mkt_td = row[0] if row else None
|
||
if mkt_td != sd:
|
||
logger.warning(
|
||
"传导扫描:mkt_daily 最新 stock 快照 %s ≠ scan_date %s —— movers 非当日,"
|
||
"候选照常落库但已记 mkt_trade_date,下游(因子桥)会据此告警。"
|
||
"若非预期,补跑 sync_market 后重扫本日即可(同日幂等)。", mkt_td, sd)
|
||
...
|
||
for t in targets.values():
|
||
ctx = graph_store.topic_context(t["target_type"], t["target"], cap=member_cap)
|
||
members = [m for m in ctx["members"] if m.get("ts_code")]
|
||
if not members:
|
||
continue
|
||
moved = [m for m in members if m["ts_code"] in movers]
|
||
quiet = [m for m in members if m["ts_code"] not in movers]
|
||
# ★ 原来是 `for m in quiet[:12]` —— 那个 12 是给旁批/展示用的,
|
||
# 却把因子的结构上限锁在 12×12=144 行/日(实测 83 吻合)。
|
||
# 落库存全量,旁批仍只喂前 12(见下面 _annotate 的入参)。
|
||
quiet_rows, upsides, over = [], [], 0
|
||
for m in quiet: # ← 去掉 [:12]
|
||
...
|
||
```
|
||
|
||
`_annotate` 调用处把喂给 brain 的 quiet 截到 12:`_annotate(out, quiet_show=12)`,
|
||
prompt 里只列前 12(旁批的信息密度需求和因子的覆盖需求本来就是两件事)。
|
||
|
||
落库 INSERT 加 `mkt_trade_date`:
|
||
|
||
```python
|
||
conn.execute("ALTER TABLE transmission_candidates "
|
||
"ADD COLUMN IF NOT EXISTS mkt_trade_date DATE") # 幂等,跟 _DDL 一起
|
||
...
|
||
"""INSERT INTO transmission_candidates
|
||
(scan_date, rank, target, target_type, paths, members_total, moved,
|
||
moved_ratio, quiet, pricing, annotation, mkt_trade_date)
|
||
VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s)"""
|
||
```
|
||
|
||
### 2.3 `hotspot._latest_mkt / _movers_set` — 接受 `trade_date`
|
||
|
||
```python
|
||
def _latest_mkt(kind: str, trade_date: Optional[date] = None):
|
||
"""trade_date 给定时取该日快照;不给才退回 max(trade_date)。
|
||
(评审 §6.2:原来恒取最新,导致 scan_date 只是个标签。)"""
|
||
with get_pool().connection() as conn:
|
||
if trade_date is None:
|
||
row = conn.execute(
|
||
"SELECT max(trade_date) FROM mkt_daily WHERE kind=%s", (kind,)).fetchone()
|
||
td = row[0] if row else None
|
||
else:
|
||
td = trade_date
|
||
...
|
||
|
||
def _movers_set(trade_date: Optional[date] = None) -> set[str]:
|
||
_, stocks = _latest_mkt("stock", trade_date)
|
||
return {x["code"] for x in stocks if (x.get("pct_change") or 0) >= 3}
|
||
```
|
||
|
||
`hotspot.scan` / `transmission.scan` 里的 `_movers_set()` 都传 `sd`。
|
||
这样 2.2 的断言就成了双保险(断言拦"快照缺",传参拦"取错日")。
|
||
|
||
### 2.4 `sync_consensus` — 参数化 `asof`(**§9-6 甲案,约 5 行**)
|
||
|
||
```python
|
||
def sync_consensus(lookback_days: int = 90,
|
||
asof: Optional[str] = None) -> dict[str, Any]:
|
||
"""...(原 docstring 保留)
|
||
|
||
asof(评审 §6.4 / 设计 §9-6 甲案):给定时按该日重放同一套聚合逻辑,
|
||
live 与 history 共用一份代码、零漂移。不给 = today(原行为)。"""
|
||
ref = date.fromisoformat(asof) if asof else date.today()
|
||
pool_codes = _pool_ts_codes()
|
||
if not pool_codes:
|
||
return {"skipped": "无已注册股票池"}
|
||
since = (ref - timedelta(days=lookback_days)).isoformat()
|
||
...
|
||
cur.execute(
|
||
"""SELECT ts_code, report_date, org_name, quarter, eps, np,
|
||
max_price, min_price, rating
|
||
FROM gp_report_rc
|
||
WHERE report_date >= %s AND report_date <= %s -- ★ 上界必加
|
||
AND ts_code IN %s""",
|
||
(since, ref.isoformat(), tuple(pool_codes)))
|
||
...
|
||
# today = date.today() → 改成用 ref
|
||
```
|
||
|
||
**★ 上界 `report_date <= ref` 是这条改动里最关键的一行** —— 不加,历史回填会把
|
||
未来研报算进当时的一致预期,就是教科书式前视。live 模式下 today 本就是上界,
|
||
所以原代码没有它也对;参数化之后必须补。
|
||
|
||
两个已知残余(回填时在文档里写明,不影响本次改动):
|
||
|
||
- `_pool_ts_codes()` 用的是**今天的池** → 历史回填仍带成员性前视(同评审 §4);
|
||
- `target_mid_avg` 是 90 天内全部研报行的简单平均,不按机构去重、不按时间加权 ——
|
||
历史与 live 一致,所以不影响可比性,但要知道这个"中枢"的语义比字面弱。
|
||
|
||
回填跑法:
|
||
|
||
```bash
|
||
# 逐日重放(一天约 800 行,500 天约 40 万行,consensus_daily 主键天然幂等)
|
||
docker compose exec backend python -c "
|
||
from app.store import market_snapshot as ms
|
||
import pandas as pd
|
||
for d in pd.bdate_range('2025-01-01','2026-07-24'):
|
||
print(d.date(), ms.sync_consensus(asof=d.date().isoformat()))"
|
||
```
|
||
|
||
### 2.5(可选,但强烈建议)`industry_pools` 开始存历史
|
||
|
||
```sql
|
||
CREATE TABLE IF NOT EXISTS industry_pools_history (
|
||
snapshot_date DATE NOT NULL,
|
||
theme TEXT NOT NULL,
|
||
members JSONB NOT NULL,
|
||
stats JSONB,
|
||
PRIMARY KEY (snapshot_date, theme)
|
||
);
|
||
```
|
||
|
||
`claim_store.upsert_pool` 末尾顺手写一行(或 `refresh_pools` 跑完整表快照一次)。
|
||
成本几乎为零,解开的是评审 §4 那一整类问题(universe 成员性前视 + z 统计量周一跳变)。
|
||
桥侧 `freeze.py` 已经在冻结当日 universe,两边互为备份。
|
||
|
||
---
|
||
|
||
## 第三批 · 需要你拍板的两件(我给出两种拍法的代码形态)
|
||
|
||
### 3.1 赛道门槛 C 的数据源:三条路,我推荐 (c)
|
||
|
||
| | 做法 | 成本 | 代价 |
|
||
|---|---|---|---|
|
||
| (a) | 从 `claims` 重建(一条 SQL) | 最低,今天就能跑 | 融合前口径:无 status/supersede/tier 裁决,会捞到已被顶替的历史归属;object 侧无 norm,同义环节不合并 |
|
||
| (b) | 桥直连 Neo4j | 中 | 破坏"桥只连三库",裁决逻辑复制进桥 |
|
||
| (c) | **基座投影表 + 第五/六个只读视图** | 中,但一次性 | 无。是投影不是打分,不破 §4 |
|
||
|
||
**我的建议是 (a) 先用于 G1 体检、(c) 作为 S1 的正解。** 体检只需要回答"这个赛道在图谱里
|
||
有没有料",(a) 的超集口径完全够用而且更保守(宁可高估覆盖,也不要因为口径太严把
|
||
本来有料的赛道判成空)。而一旦进公式,成员性必须走裁决后的口径。
|
||
|
||
(a) 的 SQL(**列名待核实**,先跑第一段确认再跑第二段):
|
||
|
||
```sql
|
||
-- 先确认列名
|
||
SELECT column_name FROM information_schema.columns
|
||
WHERE table_name IN ('claims','entity_links','company_master')
|
||
ORDER BY table_name, ordinal_position;
|
||
|
||
-- G1 体检用:环节 → 上市成员(融合前口径,仅供体检,不进公式)
|
||
CREATE OR REPLACE VIEW v_probe_segment_members AS
|
||
SELECT c.object_id AS segment_name,
|
||
c.qualifiers->>'chain' AS chain,
|
||
COALESCE(el.canonical_id, cm.ts_code) AS ts_code,
|
||
c.tier,
|
||
count(*) AS n_claims
|
||
FROM claims c
|
||
LEFT JOIN entity_links el ON el.alias_norm = c.subject_norm
|
||
LEFT JOIN company_master cm ON cm.short_name = c.subject_id
|
||
WHERE c.predicate = 'IN_SEGMENT' AND c.object_type = 'Segment'
|
||
GROUP BY 1,2,3,4;
|
||
|
||
-- 环节上下游边(单表零 join)
|
||
CREATE OR REPLACE VIEW v_probe_segment_edges AS
|
||
SELECT DISTINCT c.subject_id AS up, c.object_id AS down, c.tier
|
||
FROM claims c
|
||
WHERE c.predicate = 'SEGMENT_UPSTREAM_OF' AND c.object_type = 'Segment';
|
||
```
|
||
|
||
(c) 的形状:基座在 Neo4j 投影 beat 之后,把
|
||
`MATCH (co:Company)-[e:IN_SEGMENT {status:'active'}]->(s:Segment)` 与
|
||
`MATCH (a:Segment)-[e:SEGMENT_UPSTREAM_OF {status:'active'}]->(b:Segment)`
|
||
两条 Cypher 的结果物化到 PG 两张表,视图名就叫
|
||
`v_factor_segment_members(segment_name, chain, ts_code, tier, updated_at)` /
|
||
`v_factor_segment_edges(up, down, tier, updated_at)` ——
|
||
`freeze.py` 已经预留了对这两个视图的冻结(视图不存在时静默跳过,不报错)。
|
||
|
||
**顺带三件与 C 有关的定调建议:**
|
||
|
||
1. **`layer` 写在 yml 里人工指定,加一列 `layer_source`。** 四层枚举(材料→设备→
|
||
制造→应用)系统里根本不存在,且 `enums._SEGMENT_BAD_SUBSTR` 把"上游/中游/下游/
|
||
环节"列为环节名**禁用词素**(注释:"链位置由图上边表达")。S1 公式又不用 layer ——
|
||
**降为可选元数据,别让它阻塞 G2**。
|
||
2. **C 做两级并强制标 `source_rule`**:`graph`(强,可审计到具体 claim)/
|
||
`fallback`(弱,概念标签或申万行业白名单兜底)。否则核聚变/低空经济/商业航天
|
||
一旦查不到,C 要么塌成空集、要么退化成"电子+计算机"。
|
||
3. **覆盖为空的赛道直接进采集清单** —— 这就是基座「需求闭环」的正用法,
|
||
比放弃赛道有价值(`claim_store.open_demand("research_report", "Concept", <赛道>)`)。
|
||
|
||
### 3.2 权重结构:两段式(推荐)还是真混合
|
||
|
||
传导池内 95% 为 0,z-score 后 `0.5·z_T` 的组间落差压过另两项全幅
|
||
(数值验证:200 只候选 / 10 只有传导 → top20 里传导票 9.8/10)。所以现在的
|
||
0.5/0.3/0.2 事实上是"传导票优先,组内再比冷和便宜"。二选一:
|
||
|
||
**A(推荐)· 承认两段式,代码更短、行为与文字一致**
|
||
|
||
```python
|
||
def akg_score(pool: pd.DataFrame) -> pd.Series:
|
||
"""景气度漏斗的池内排序(评审 §2 方案 A)。
|
||
|
||
结构 = 传导档位(主键) + 组内「冷 & 便宜」(次键)。
|
||
这不是"降低传导权重",恰恰是把"传导第一"写实:
|
||
0.5/0.3/0.2 在 95% 稀疏下已经等价于此,写明比藏在 z-score 里好 ——
|
||
可解释、可调试、一眼看出是哪一段在起作用。
|
||
"""
|
||
t = np.log1p(pool["transmission"].fillna(0.0))
|
||
# 档位:0 = 无传导;1 = 有传导且强度在有传导组的下半;2 = 上半
|
||
tier = pd.Series(0, index=pool.index, dtype=float)
|
||
hit = t > 0
|
||
if hit.any():
|
||
tier[hit] = np.where(t[hit] >= t[hit].median(), 2.0, 1.0)
|
||
|
||
zH = -_robust_z(pool["heat"].fillna(pool["heat"].median()))
|
||
zV = _robust_z(pool["upside"])
|
||
tiebreak = 0.6 * zH + 0.4 * zV # 次序仍是「还没热」>「便宜」
|
||
|
||
# 档间不可逆(这正是"传导第一"的含义),档内连续排序
|
||
span = 10.0 # > tiebreak 的理论幅度(约 ±6)
|
||
return tier * span + tiebreak
|
||
```
|
||
|
||
好处:档位与次键各自可独立观察;调试时 `groupby(tier)` 直接看出每档的表现;
|
||
把 §9-9("还没热"该不该做条件项)一并解决了 —— 它现在天然是组内量。
|
||
|
||
**B · 真要连续混合**:`w_T` 降到 **0.10~0.15**,或对 `z_T` 做有界变换:
|
||
|
||
```python
|
||
zT = np.clip((t - t.mean()) / (t.std() or 1.0), -3, 3) / 3 * 1.5 # 压到与另两项同量级
|
||
score = 0.5 * zT + 0.3 * zH + 0.2 * zV # 实测能把 top20 命中从 9.8 拉到 7.5
|
||
```
|
||
|
||
**无论选哪个,都建议同时注册 `akg_gate`**(0/1,全覆盖池出行)。否则平台的 IC/分层
|
||
只看得到"池内排序"的价值,**完全看不到两道门槛的价值**——而门槛才是这套逻辑的主体。
|
||
两个因子、两个仪表,各答一个问题,"哪一项在拖后腿"一眼可见。
|
||
|
||
---
|
||
|
||
## 落地顺序建议
|
||
|
||
```
|
||
今天 1.1 视图替换 → 1.2/1.3/1.5/1.6/1.7/1.8 桥侧补丁 → 1.9 freeze.py
|
||
跑一次:run.py views → run.py freeze → run.py build all --mode daily
|
||
(freeze 的 warnings 就是硬伤 1 的实测量级,正好拿来定 2.1/2.2 的 cap)
|
||
本周 2.1~2.3 基座三处(改完 replay 不需要,只影响新扫描)→ 2.5 池历史表
|
||
1.10 probe 加 corr 节 → 3.1(a) 的两个体检视图 → 跑 run.py probe
|
||
⇒ G1 体检报告(含三项相关矩阵 + 赛道命中量级)发我
|
||
再定 3.1 的 C 口径与赛道清单、3.2 的权重结构 → G2/G3
|
||
2.4 sync_consensus 参数化 + 逐日回填 → G4
|
||
```
|
||
|
||
**如果只能做一件:`freeze.py`。** 别的都能补,快照不存就永久没了。
|