34 KiB
评审问题的修复方案(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 需要你知道的三点(不是"白费",但要说清)
-
已存的 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 加的断言在
mkt_daily最新 stock 快照 ≠scan_date时中止扫描。 在公告批量入库期间,如果 beat 被饿到(就是announcement-corpus-bulk-load记的那个单队列风险),17:30 快照晚点 → 当天传导直接没有产出。 这比"用 T−1 的 movers 生成带今日戳的候选"好,但你要预期到它会偶发触发, 也意味着 live 传导史积累得慢一点。队列分离那次修复的实机验证, 现在多了一个理由要尽快做完。 -
唯一可能产生 LLM 花费的是"提环节覆盖",但它也不浪费任何东西。 如果 G1 体检发现赛道覆盖太薄,对策是
segment_backfill(环节专项遍存量重抽)。 它零新语料、dedup_key幂等、只追加不删旧 claims —— 花的是 brain 调用, 不是把已有成果作废。
0.4 落地前后的自查(跑一遍留个数,最稳)
# 改动前后各跑一次,两次输出应完全一致(除 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,一起跑即可。
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)
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 — 分块提交(修放量会炸)
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)
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 — 追加四个旋钮
# 子因子是否受覆盖池限制: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 — 行情按月分块(修放量会炸)
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)
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 之前插入:
# ---- 年报污染防护(评审 §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 里加了告警,这里再补一句):
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 三处小改:
# ① 顶部 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 需要)
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 的根)
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 + 新鲜度断言 + 记快照日
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:
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
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 行)
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 一致,所以不影响可比性,但要知道这个"中枢"的语义比字面弱。
回填跑法:
# 逐日重放(一天约 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 开始存历史
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(列名待核实,先跑第一段确认再跑第二段):
-- 先确认列名
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 有关的定调建议:
layer写在 yml 里人工指定,加一列layer_source。 四层枚举(材料→设备→ 制造→应用)系统里根本不存在,且enums._SEGMENT_BAD_SUBSTR把"上游/中游/下游/ 环节"列为环节名禁用词素(注释:"链位置由图上边表达")。S1 公式又不用 layer —— 降为可选元数据,别让它阻塞 G2。- C 做两级并强制标
source_rule:graph(强,可审计到具体 claim)/fallback(弱,概念标签或申万行业白名单兜底)。否则核聚变/低空经济/商业航天 一旦查不到,C 要么塌成空集、要么退化成"电子+计算机"。 - 覆盖为空的赛道直接进采集清单 —— 这就是基座「需求闭环」的正用法,
比放弃赛道有价值(
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(推荐)· 承认两段式,代码更短、行为与文字一致
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 做有界变换:
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。 别的都能补,快照不存就永久没了。