495 lines
25 KiB
Python
495 lines
25 KiB
Python
"""产业链细化 · 诊断(交付一,纯只读)——六个读数一次跑出。
|
||
|
||
方案出处: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())
|