akg-factor-bridge/chain_diag.py

495 lines
25 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""产业链细化 · 诊断(交付一,纯只读)——六个读数一次跑出。
方案出处astock-kg/docs/产业链细化方案_定稿_2026-08-05.md 第三节。
口径基准 = 桥的消费面:环节投影视图 v_factor_segment_members / v_factor_segment_edges
(基座每晨 06:10 删后插刷新)+ frontier_tracks.yml 三路映射 + 平台档位表。
诊断读的就是选股链路实际吃到的数据,不直接打图——图与投影的差异是另一类问题。
产出目录 data/chain_diag_<日期>/(同日重跑整目录覆盖,幂等):
01_链序未明环节榜.csv —— 补边的靶子,按上市成员数×热度排
02_未归属环节榜.csv —— 三路映射都够不着的环节
03_锚覆盖矩阵.csv —— 十赛道 × 三路现状
03b_未映射链名榜.csv —— kg_chains 扩充拍板的弹药(含赛道提示列)
03c_深海检索.csv —— chain/segment_name 的深海关键词检索
04_错杀对照组_主榜侧_<档位日>.csv / 04b_错杀对照组_观察侧_<档位日>.csv —— 冻结基线
05_碎片化读数.txt —— 成员分布 / 链名覆盖 / 度分布 / 连通分量
06_研报盘点.csv —— 在库研报逐篇产出读数(单篇判收的基线)
06b_研报补给清单.csv —— 十赛道逐条缺口 + 建议检索词(拿去找报告)
汇总.md —— 头条数字 + 拍板建议(同时打印到终端)
跑法【桥机 factorevaluation-UTC · ~/akg-factor-bridge】
docker compose exec -T akg-factor-bridge python chain_diag.py
纯只读PG / MySQL 全部 SELECT绝不写库产物只落容器内 data/ 目录
(代码卷挂载,宿主 ~/akg-factor-bridge/data/ 直接可见)。
个别数据源不可用(如档位表尚无当日行)对应读数跳过并在汇总标 ⚠️,其余照出。
"""
from __future__ import annotations
import datetime as dt
import os
import re
import pandas as pd
import common
import db
import tracks
_LINES: list[str] = [] # 汇总累积器(写 汇总.md + 打印)
def say(msg: str = "") -> None:
_LINES.append(msg)
print(msg)
def _median(vals: list[float]) -> float | None:
xs = sorted(v for v in vals if v is not None)
return xs[len(xs) // 2] if xs else None
# ---------------------------------------------------------------- 数据装载
def _load_projection():
mem = db.read_pg(
"SELECT segment_name, ts_code, member_name, chain "
"FROM v_factor_segment_members")
mem["segment_name"] = mem["segment_name"].astype(str).str.strip()
mem["chain"] = mem["chain"].fillna("").astype(str).str.strip()
edg = db.read_pg(
"SELECT src_segment, dst_segment, rel_type, chain "
"FROM v_factor_segment_edges")
edg["chain"] = edg["chain"].fillna("").astype(str).str.strip()
return mem, edg
def _load_heat() -> dict[str, float]:
"""最新日最新批次热度分,键=前缀码(表内本就是 SZ002625 形态)。拿不到给空。"""
df = db.read_mysql("heat", """
SELECT s.stock_code, s.score FROM stock_fund_heat_scores s
JOIN (SELECT trade_date, MAX(batch_no) bn FROM stock_fund_heat_scores
WHERE trade_date = (SELECT MAX(trade_date) FROM stock_fund_heat_scores)
GROUP BY trade_date) m
ON m.trade_date = s.trade_date AND m.bn = s.batch_no""")
return {str(r.stock_code).strip(): float(r.score)
for r in df.itertuples() if pd.notna(r.score)}
def _anchor_maps(yml: dict):
"""yml → (环节名→赛道, 链名→赛道, 赛道→主题键集, 赛道名→别名等关键词)。"""
seg2track: dict[str, str] = {}
chain2track: dict[str, str] = {}
kw: dict[str, list[str]] = {}
for tr in yml.get("tracks") or []:
name = tr["name"]
for s in tr.get("kg_segments") or []:
seg2track.setdefault(str(s).strip(), name)
for c in tr.get("kg_chains") or []:
chain2track.setdefault(str(c).strip(), name)
kw[name] = ([name] + [str(a) for a in tr.get("aliases") or []]
+ [str(t) for t in tr.get("kg_themes") or []]
+ [str(c) for c in tr.get("kg_chains") or []])
return seg2track, chain2track, kw
# ---------------------------------------------------------------- 读数一/二共用的环节聚合
def _segment_table(mem: pd.DataFrame, heat: dict[str, float]):
"""环节聚合:上市成员集合、链名集合、热度中位、代表成员。"""
rows = {}
for r in mem.itertuples():
seg = r.segment_name
d = rows.setdefault(seg, {"listed": set(), "chains": set(), "names": []})
if pd.notna(r.ts_code) and str(r.ts_code).strip():
ts = str(r.ts_code).strip()
if ts not in d["listed"]:
d["listed"].add(ts)
d["names"].append((str(r.member_name), ts))
if r.chain:
d["chains"].add(r.chain)
out = {}
for seg, d in rows.items():
hs = [heat.get(common.to_prefix(ts)) for ts in d["listed"]]
named = sorted(d["names"],
key=lambda x: -(heat.get(common.to_prefix(x[1])) or 0.0))
out[seg] = {
"listed": len(d["listed"]), "listed_set": d["listed"],
"chains": sorted(d["chains"]), "heat_med": _median(hs),
"top_members": "".join(n for n, _ in named[:3]),
}
return out
# ---------------------------------------------------------------- 读数一
def diag_unordered(segs: dict, edg: pd.DataFrame, seg2track, chain2track,
outdir: str) -> None:
touched = set(edg["src_segment"]) | set(edg["dst_segment"])
rows = []
for seg, d in segs.items():
if seg in touched:
continue
anchored = seg2track.get(seg) or next(
(chain2track[c] for c in d["chains"] if c in chain2track), "")
rows.append((seg, d["listed"], "|".join(d["chains"]),
d["heat_med"], d["top_members"], anchored))
df = (pd.DataFrame(rows, columns=["环节", "上市成员数", "链名", "热度中位",
"代表成员", "已锚赛道"])
.sort_values(["上市成员数", "热度中位"], ascending=[False, False],
na_position="last"))
df.to_csv(os.path.join(outdir, "01_链序未明环节榜.csv"),
index=False, encoding="utf-8-sig")
n_all, n_zero = len(segs), len(df)
w_all = sum(d["listed"] for d in segs.values())
w_zero = sum(r[1] for r in rows)
say(f"一、链序未明环节榜:{n_zero}/{n_all} 个环节零链序边"
f"{n_zero / max(1, n_all):.0%};按上市成员关系加权 "
f"{w_zero / max(1, w_all):.0%}。头部成员≥5{sum(1 for r in rows if r[1] >= 5)}"
f"——这就是补边靶子的规模。")
# ---------------------------------------------------------------- 读数二
def diag_unmapped(segs: dict, seg2track, chain2track, track_all: set,
outdir: str) -> None:
rows = []
for seg, d in segs.items():
anchored = seg in seg2track or any(c in chain2track for c in d["chains"])
if anchored:
continue
inter = len(d["listed_set"] & track_all)
ratio = inter / d["listed"] if d["listed"] else 0.0
status = "无任何归属" if inter == 0 else f"仅成员经主题({ratio:.0%})"
rows.append((seg, d["listed"], "|".join(d["chains"]), status,
d["heat_med"], d["top_members"]))
df = pd.DataFrame(rows, columns=["环节", "上市成员数", "链名", "归属状态",
"热度中位", "代表成员"])
df["_o"] = (df["归属状态"] != "无任何归属").astype(int)
df = (df.sort_values(["_o", "上市成员数"], ascending=[True, False])
.drop(columns="_o"))
df.to_csv(os.path.join(outdir, "02_未归属环节榜.csv"),
index=False, encoding="utf-8-sig")
n_none = int((df["归属状态"] == "无任何归属").sum())
say(f"二、未归属环节榜:结构锚够不着 {len(df)} 个环节,其中 {n_none} 个连成员"
f"都不经任何映射主题两头落空。头部成员≥5"
f"{int((df[df['归属状态'] == '无任何归属']['上市成员数'] >= 5).sum())} 个。")
# ---------------------------------------------------------------- 读数三
def diag_matrix(yml: dict, mem: pd.DataFrame, chain2track, kw, outdir: str) -> None:
themes_in_pool = {str(t).strip()
for t in db.read_pg("SELECT theme FROM industry_pools")["theme"]}
chains_have = set(mem.loc[mem["chain"] != "", "chain"])
segs_have = set(mem["segment_name"])
df_all, _missing = tracks.resolve_members(only_confirmed=False, dedup=False)
rows = []
for tr in yml.get("tracks") or []:
name = tr["name"]
sub = df_all[df_all["track"] == name]
n_seg_keys = len(tr.get("kg_segments") or [])
n_chain_keys = len(tr.get("kg_chains") or [])
n_theme_keys = len(tr.get("kg_themes") or [])
rows.append((
name,
n_theme_keys,
sum(1 for t in tr.get("kg_themes") or [] if str(t).strip() in themes_in_pool),
n_chain_keys,
sum(1 for c in tr.get("kg_chains") or [] if str(c).strip() in chains_have),
n_seg_keys,
sum(1 for s in tr.get("kg_segments") or [] if str(s).strip() in segs_have),
sub["ts_code"].nunique(),
sub[sub["source_rule"] != "pool_theme"]["ts_code"].nunique(),
))
df = pd.DataFrame(rows, columns=["赛道", "主题键", "主题命中", "链锚键", "链锚命中",
"环节锚键", "环节锚命中", "成员数", "图谱锚成员数"])
df.to_csv(os.path.join(outdir, "03_锚覆盖矩阵.csv"),
index=False, encoding="utf-8-sig")
dz = df[(df["链锚键"] == 0) & (df["环节锚键"] == 0)]["赛道"].tolist()
say(f"三、锚覆盖矩阵:链锚合计 {int(df['链锚键'].sum())}、环节锚合计 "
f"{int(df['环节锚键'].sum())};两类结构锚双零的赛道:{''.join(dz) or ''}")
# 03b 未映射链名榜kg_chains 扩充的弹药
ch = (mem[mem["chain"] != ""]
.groupby("chain")
.agg(上市成员数=("ts_code", lambda s: s.dropna().nunique()),
环节数=("segment_name", "nunique"))
.reset_index().rename(columns={"chain": "链名"}))
ch["已映射赛道"] = ch["链名"].map(lambda c: chain2track.get(c, ""))
ch["赛道提示"] = ch["链名"].map(
lambda c: "".join(sorted({t for t, ks in kw.items()
if any(k and (k in c or c in k) for k in ks)})))
ch = ch.sort_values(["已映射赛道", "上市成员数"], ascending=[True, False])
ch.to_csv(os.path.join(outdir, "03b_未映射链名榜.csv"),
index=False, encoding="utf-8-sig")
n_un = int((ch["已映射赛道"] == "").sum())
n_hint = int(((ch["已映射赛道"] == "") & (ch["赛道提示"] != "")).sum())
say(f" 未映射链名 {n_un} 个(其中 {n_hint} 个带赛道提示,是 kg_chains 扩充"
f"首批拍板对象);空串链名成员边占比见读数五。")
# 03c 深海检索yml 挂账待办)
pats = ["深海", "海洋", "水下", "海底"]
got = []
for p in pats:
hit_c = mem[mem["chain"].str.contains(p, na=False)]
for c, g in hit_c.groupby("chain"):
got.append(("链名", c, p, g["ts_code"].dropna().nunique()))
hit_s = mem[mem["segment_name"].str.contains(p, na=False)]
for s, g in hit_s.groupby("segment_name"):
got.append(("环节名", s, p, g["ts_code"].dropna().nunique()))
dfc = (pd.DataFrame(sorted(set(got)),
columns=["位置", "名称", "命中词", "上市成员数"])
.sort_values(["位置", "上市成员数"], ascending=[True, False]))
dfc.to_csv(os.path.join(outdir, "03c_深海检索.csv"),
index=False, encoding="utf-8-sig")
say(f" 深海检索:{''.join(pats)} 共命中 {len(dfc)}"
f"(深海科技锚定拍板一并处理)。")
# ---------------------------------------------------------------- 读数四
def diag_misskill(mem: pd.DataFrame, outdir: str) -> None:
d = db.read_mysql("factor", "SELECT MAX(trade_date) d FROM t_factor_akg_gate")
v = None if d.empty else d.iloc[0, 0]
if v is None or pd.isna(v):
say("四、错杀对照组:档位表尚无数据,跳过(部署后重跑本诊断补冻结)。")
return
gd = pd.Timestamp(v).date().isoformat()
g = db.read_mysql("factor", "SELECT stock_code, factor_value "
"FROM t_factor_akg_gate WHERE trade_date=%s", (gd,))
up = db.read_mysql("factor", "SELECT stock_code, factor_value "
"FROM t_factor_akg_upside WHERE trade_date=%s", (gd,))
gate = {str(r.stock_code).strip(): float(r.factor_value) for r in g.itertuples()}
upside = {str(r.stock_code).strip(): float(r.factor_value) for r in up.itertuples()}
# 图谱证据按"该股自己的成员边"逐行归集——链名取本股边上的 chain 修饰。
# 首版从环节聚合继承整个环节的链名集合,串味成"3D打印、6G"满屏08-05 实测),勿回退。
stock_ev: dict[str, dict] = {} # 前缀码 → 图谱证据
for r in mem.itertuples():
if pd.isna(r.ts_code) or not str(r.ts_code).strip():
continue
k = common.to_prefix(str(r.ts_code).strip())
e = stock_ev.setdefault(k, {"name": "", "segs": set(), "chains": set()})
e["segs"].add(r.segment_name)
if r.chain:
e["chains"].add(r.chain)
names = {}
try:
pools = db.read_pg("SELECT members FROM industry_pools")
import json as _json
for _, r in pools.iterrows():
ms = r["members"]
if isinstance(ms, str):
ms = _json.loads(ms)
for m in ms or []:
if (m or {}).get("ts_code"):
names[common.to_prefix(m["ts_code"])] = m.get("name") or ""
except Exception as e: # noqa: BLE001 —— 简称拿不到不影响榜单
say(f" ⚠️ 成员简称加载失败(榜单缺简称列): {e!r}")
def _row(k):
e = stock_ev[k]
nm = names.get(k, "")
risk = "风险股" if re.match(r"^(\*?S?ST|退市)", nm.replace(" ", "")) else ""
return (k, nm, risk, len(e["segs"]),
"".join(sorted(e["segs"])[:3]),
"".join(sorted(e["chains"])[:3]))
a_rows = [(_row(k) + (round(upside[k], 4),))
for k, gv in gate.items()
if gv == 0.0 and k in upside and upside[k] >= 0 and k in stock_ev]
dfa = (pd.DataFrame(a_rows, columns=["代码", "简称", "风险", "环节数",
"环节", "链名", "upside"])
.sort_values(["环节数", "upside"], ascending=[False, False]))
dfa.to_csv(os.path.join(outdir, f"04_错杀对照组_主榜侧_{gd}.csv"),
index=False, encoding="utf-8-sig")
b_rows = [_row(k) for k, gv in gate.items()
if gv == 0.0 and k not in upside and k in stock_ev]
dfb = (pd.DataFrame(b_rows, columns=["代码", "简称", "风险", "环节数",
"环节", "链名"])
.sort_values("环节数", ascending=False))
dfb.to_csv(os.path.join(outdir, f"04b_错杀对照组_观察侧_{gd}.csv"),
index=False, encoding="utf-8-sig")
say(f"四、错杀对照组(档位日 {gd},已冻结):主榜侧 {len(dfa)}"
f"有覆盖、upside≥0、有图谱环节证据、却 gate=0——赛道映射够不着它们"
f"观察侧 {len(dfb)} 只(无覆盖、有环节证据、没进观察档)。"
f"其中风险股 {int((dfa['风险'] != '').sum()) + int((dfb['风险'] != '').sum())} 只属正当拦截,读榜时剔除。")
# ---------------------------------------------------------------- 读数五
def diag_fragmentation(segs: dict, mem: pd.DataFrame, edg: pd.DataFrame,
outdir: str) -> None:
out = []
buckets = [(0, 0), (1, 1), (2, 2), (3, 5), (6, 10), (11, 30), (31, 10 ** 9)]
cnt = {b: 0 for b in buckets}
for d in segs.values():
for lo, hi in buckets:
if lo <= d["listed"] <= hi:
cnt[(lo, hi)] += 1
break
out.append("上市成员数分布(环节个数):")
for (lo, hi), n in cnt.items():
label = f"{lo}" if lo == hi else (f"{lo}-{hi}" if hi < 10 ** 9 else f"{lo}")
out.append(f" 成员 {label:>5} {n}")
n_rows = len(mem)
n_empty_chain = int((mem["chain"] == "").sum())
ups = edg[edg["rel_type"] == "SEGMENT_UPSTREAM_OF"]
drv = edg[edg["rel_type"] == "DRIVES"]
out.append(f"\n链名修饰:成员边 {n_rows} 行,空串链名 {n_empty_chain}"
f"{n_empty_chain / max(1, n_rows):.0%}"
f"非空链名 {mem.loc[mem['chain'] != '', 'chain'].nunique()} 个。")
ups_empty = int((ups["chain"] == "").sum())
out.append(f"链序边SEGMENT_UPSTREAM_OF {len(ups)} 条(其中无链名 {ups_empty}"
f"{ups_empty / max(1, len(ups)):.0%}——环节遍升级后新边应带链名)"
f";环节级 DRIVES {len(drv)} 条。")
deg: dict[str, int] = {}
for r in ups.itertuples():
deg[r.src_segment] = deg.get(r.src_segment, 0) + 1
deg[r.dst_segment] = deg.get(r.dst_segment, 0) + 1
hubs = sorted(deg.items(), key=lambda kv: -kv[1])[:10]
out.append("\n链序度最高的环节(枢纽):" +
"".join(f"{s}({n})" for s, n in hubs))
parent: dict[str, str] = {}
def find(x: str) -> str:
while parent.get(x, x) != x:
parent[x] = parent.get(parent[x], parent[x])
x = parent[x]
return x
for r in ups.itertuples():
a, b = find(r.src_segment), find(r.dst_segment)
parent.setdefault(a, a)
parent.setdefault(b, b)
if a != b:
parent[a] = b
comp: dict[str, list[str]] = {}
nodes = set(ups["src_segment"]) | set(ups["dst_segment"])
for n in nodes:
comp.setdefault(find(n), []).append(n)
sizes = sorted(comp.values(), key=len, reverse=True)
out.append(f"\n连通分量(仅上下游边参与,无向):{len(sizes)} 个块,"
f"最大 {len(sizes[0]) if sizes else 0} 个环节。")
for i, c in enumerate(sizes[:10], 1):
out.append(f"{i}{len(c)}{''.join(sorted(c)[:6])}"
+ ("" if len(c) > 6 else ""))
txt = "\n".join(out)
with open(os.path.join(outdir, "05_碎片化读数.txt"), "w", encoding="utf-8") as f:
f.write(txt + "\n")
say(f"五、碎片化读数:链序骨架 {len(sizes)} 个连通块(最大 "
f"{len(sizes[0]) if sizes else 0} 环节);成员边空串链名占比 "
f"{n_empty_chain / max(1, n_rows):.0%}。明细见 05_碎片化读数.txt。")
# ---------------------------------------------------------------- 读数六
def diag_research(yml: dict, mem: pd.DataFrame, edg: pd.DataFrame, kw,
outdir: str) -> None:
try:
docs = db.read_pg("""
SELECT d.doc_id::text AS doc_id, d.title, d.disclosure_date,
count(*) FILTER (WHERE c.predicate = 'IN_SEGMENT') AS n_in_segment,
count(*) FILTER (WHERE c.predicate = 'SEGMENT_UPSTREAM_OF') AS n_upstream,
count(*) FILTER (WHERE c.predicate = 'DRIVES') AS n_drives
FROM documents d
LEFT JOIN claims c ON c.doc_id = d.doc_id
WHERE d.source_type = 'research_report'
GROUP BY 1, 2, 3 ORDER BY d.disclosure_date""")
except Exception as e: # noqa: BLE001 —— documents/claims 读不到就退化为提示
say(f"六、研报盘点PG documents/claims 读取失败({e!r})——"
f"改在基座机跑同名 SQL见 汇总.md 附注)。")
return
docs.to_csv(os.path.join(outdir, "06_研报盘点.csv"),
index=False, encoding="utf-8-sig")
chains_series = mem.loc[mem["chain"] != "", ["chain", "ts_code"]]
rows = []
for tr in yml.get("tracks") or []:
name = tr["name"]
# 标题命中关键词短英文缩写AI/6G在标题里过度匹配只留 ≥3 字符
# 或含中文的键;这是提示列口径,不是归属判定。
ks = [k for k in kw[name]
if k and (len(k) >= 3 or any(ord(ch) > 127 for ch in k))]
n_docs = int(docs["title"].fillna("").map(
lambda t, ks=ks: any(k in t for k in ks)).sum()) if len(docs) else 0
tr_chains = {str(c).strip() for c in tr.get("kg_chains") or []}
n_edges = int(edg["chain"].isin(tr_chains).sum()) if tr_chains else 0
n_members = (chains_series[chains_series["chain"].isin(tr_chains)]
["ts_code"].dropna().nunique()) if tr_chains else 0
gap = ("缺深度研报" if n_docs == 0
else ("链序未立" if n_edges < 5 else "初步成形"))
hint = "".join(dict.fromkeys(
[name] + [str(a) for a in tr.get("aliases") or []]))
rows.append((name, n_docs, n_edges, n_members, gap,
f"{hint}」+「产业链」组合:深度/全景/梳理/图谱/框架"))
df = pd.DataFrame(rows, columns=["赛道", "标题命中研报数", "链内链序边",
"链锚上市成员", "缺口判定", "建议检索词"])
df.to_csv(os.path.join(outdir, "06b_研报补给清单.csv"),
index=False, encoding="utf-8-sig")
n_lack = int((df["缺口判定"] == "缺深度研报").sum())
say(f"六、研报补给清单:在库研报 {len(docs)}"
f"IN_SEGMENT 合计 {int(docs['n_in_segment'].sum()) if len(docs) else 0}"
f"上下游合计 {int(docs['n_upstream'].sum()) if len(docs) else 0}——单篇产出基线);"
f"十赛道中 {n_lack} 条标题层面零研报。逐条缺口与检索词见 06b。")
say(" 投递约定:找到的报告放 MinIO inbox/research/,文件名带披露日"
f"YYYY-MM-DD 或紧凑八位);部署交付二后走 targeted 队列即到即抽。")
# ---------------------------------------------------------------- 主流程
def main() -> int:
today = dt.date.today().isoformat()
outdir = os.path.join("data", f"chain_diag_{today}")
os.makedirs(outdir, exist_ok=True)
say(f"产业链细化诊断 @ {today}(只读;口径=环节投影视图+三路映射+档位表)")
say("")
yml = tracks.load_yml()
mem, edg = _load_projection()
if mem.empty:
say("❌ 环节投影为空——基座 06:10 投影任务没跑或视图未建,先修再诊。")
return 2
seg2track, chain2track, kw = _anchor_maps(yml)
try:
heat = _load_heat()
except Exception as e: # noqa: BLE001
say(f"⚠️ 热度加载失败(榜单热度列为空,排序退化为纯成员数): {e!r}")
heat = {}
segs = _segment_table(mem, heat)
try:
track_df, _missing = tracks.resolve_members(only_confirmed=True)
track_all = set(track_df["ts_code"].astype(str))
except Exception as e: # noqa: BLE001
say(f"⚠️ 赛道成员表解析失败(读数二按纯结构锚判定): {e!r}")
track_all = set()
for name, fn in [
("读数一", lambda: diag_unordered(segs, edg, seg2track, chain2track, outdir)),
("读数二", lambda: diag_unmapped(segs, seg2track, chain2track, track_all, outdir)),
("读数三", lambda: diag_matrix(yml, mem, chain2track, kw, outdir)),
("读数四", lambda: diag_misskill(mem, outdir)),
("读数五", lambda: diag_fragmentation(segs, mem, edg, outdir)),
("读数六", lambda: diag_research(yml, mem, edg, kw, outdir)),
]:
try:
fn()
except Exception as e: # noqa: BLE001 —— 单读数失败不拖死整诊
say(f"⚠️ {name} 失败(其余照出): {e!r}")
say("")
say("下一步(方案第三节判收):把本目录整包发回拍板——首批建议看 "
"03b 未映射链名榜头部、02 未归属榜头部与 06b 补给清单;"
"样板研报入库后重跑本诊断,对比 01/05/06 的前后读数。")
say("\n附注:读数六若在桥机因权限失败,到基座机跑等价 SQL")
say(" 【基座机 tlai4090 · ~/project/astock-kg】")
say(" docker compose exec -T postgres psql -U akg -d akg -c \""
"SELECT d.title, count(*) FILTER (WHERE c.predicate='SEGMENT_UPSTREAM_OF') n_up "
"FROM documents d LEFT JOIN claims c ON c.doc_id=d.doc_id "
"WHERE d.source_type='research_report' GROUP BY 1 ORDER BY n_up DESC;\"")
with open(os.path.join(outdir, "汇总.md"), "w", encoding="utf-8") as f:
f.write(f"# 产业链细化诊断汇总({today}\n\n"
+ "\n".join(_LINES) + "\n")
print(f"\n产物目录: {outdir}/(汇总.md + 各榜单 CSV")
return 0
if __name__ == "__main__":
raise SystemExit(main())