diff --git a/chain_diag.py b/chain_diag.py new file mode 100644 index 0000000..b6fa508 --- /dev/null +++ b/chain_diag.py @@ -0,0 +1,490 @@ +"""产业链细化 · 诊断(交付一,纯只读)——六个读数一次跑出。 + +方案出处: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(segs: dict, 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()} + + stock_ev: dict[str, dict] = {} # 前缀码 → 图谱证据 + for seg, dd in segs.items(): + for ts in dd["listed_set"]: + k = common.to_prefix(ts) + e = stock_ev.setdefault(k, {"name": "", "segs": [], "chains": set()}) + e["segs"].append(seg) + e["chains"] |= set(dd["chains"]) + 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(set(e["segs"])), + "、".join(sorted(set(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(segs, 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())