#!/usr/bin/env python3 # -*- coding: utf-8 -*- r"""事实契约的**重算台账生成端**(用户令 2026-09-19:「要让观澜从重算台账生成 claim」)。 ## 背景(为什么换掉人裁底稿) 门户的「本场结论·契约生成」段、`/detail` 的问答/报告面(`/api/facts`)都读事实契约 (`outputs/<场>/guanlan/facts_contract_v0.json` + `derived/{detail_cards,portal_claims,qa_refs}.json`)。 契约的输入此前是**人裁底稿**:`sop/findings.json`(119 条结论)+ `paradigm_r1` 的 E3/E5/E8 三份底稿 —— 清过产物、只从 `data/raw` 重算的机器上它们必然不在位 ⇒ 契约面永远 503、门户结论段停在旧数。 用户令选定:**由观澜从重算台账生成 claim**,这三处随重算刷新。 ## 本器做什么 按 `scripts/guanlan_facts_contract.py` 读的四份输入的**字段结构**产出同结构件(契约构建器一行不改): outputs/<场>/sop/findings.json 结论本体(claim 列表) outputs/<场>/paradigm_r1/experiments/E5_coverage_table/底稿.json finding_rows: 模块映射 outputs/<场>/paradigm_r1/experiments/E3_candidate_closure/底稿.json rows: 归宿提议 outputs/<场>/paradigm_r1/experiments/E8_rudong_rebuild/底稿_s0.json findings_rows: 系统/台号 claim 全部由**重算产物**算出(每条都记 `source_refs` 的件与 sha16): · 报警集中台 windscada/alarms.parquet(逐台计数 + 同族码集中度) · 重复检修件 windscada/workorders.parquet(同台同部件重复更换) · 停机损失集中 windscada/loss_monthly.parquet(停机态损失占比) · 温度离群台 windscada/temp_bins.parquet(同工况同通道 dev 离群) · 功率曲线偏离台 windscada/powercurve_dev.parquet(同风速残差) · 振动面报警台 m5_cms_tcm/fusion_38.csv(融合级 定论/准定论·预警) ## 判据纪律(不造数) · 每条 claim 的 `verdict` 只取契约六枚举,且**只能用台账能支撑的那一档**:全部相对判据 + 无实物锚 ⇒ 最高 「候选」(振动面来自观澜自算融合, 同样封顶候选);口径性/样本性问题用「参考」或「INSUFFICIENT」。 · 逐条写 `mandatory_limitation`(缺什么证据)与 `falsifiability`(怎么证伪),而不是只给结论。 · `coverage` 写**实算**的窗天数与覆盖台数; 点名台号写进 `sample_level.named`。 · 与振动线/厂家正本无关: 本件是**观澜自算结论**, 不冒充人工裁决(generate 侧写在文件头与 meta)。 用法: python scripts/sop_findings_from_ledger.py # 生成四份输入件并自登记 python scripts/sop_findings_from_ledger.py --dry-run # 只报会生成几条 claim python scripts/sop_findings_from_ledger.py --status """ from __future__ import annotations import argparse import hashlib import json import pathlib import sys import time import pandas as pd ROOT = pathlib.Path(__file__).resolve().parents[1] sys.path.insert(0, str(ROOT)) sys.path.insert(0, str(ROOT / 'scripts')) from src import paths as P # noqa: E402 SIX = ('定论', '准定论·预警', '候选', '参考', 'INSUFFICIENT', '撤回') OUT_FILES = ('sop/findings.json', 'paradigm_r1/experiments/E5_coverage_table/底稿.json', 'paradigm_r1/experiments/E3_candidate_closure/底稿.json', 'paradigm_r1/experiments/E8_rudong_rebuild/底稿_s0.json') def _sha16(p: pathlib.Path) -> str: return hashlib.sha256(p.read_bytes()).hexdigest()[:16] def _win(t, days: int) -> tuple: """(起, 止) 与窗天数 —— 从数据实算, 不写死。""" end = pd.to_datetime(t.max()) start = end - pd.Timedelta(days=days) return start, end, int((end - start).days) def claims(store: pathlib.Path, m5: pathlib.Path) -> list[dict]: out: list[dict] = [] def add(cid, title, verdict, refs, *, named=(), systems=(), module='', window_days=None, n_machines=None, limitation='', falsify='', nature='', claim_class='相对判据', coverage_note=''): out.append(dict(id=cid, title=title, verdict=verdict, source_refs=refs, named=list(named), systems=list(systems), module=module, window_days=window_days, n_machines=n_machines, mandatory_limitation=limitation, falsifiability=falsify, problem_nature=nature, claim_class=claim_class, coverage_note=coverage_note)) # ── ① 报警集中台 ──────────────────────────────────────────────────────────── ap = store / 'alarms.parquet' if ap.is_file(): al = pd.read_parquet(ap) al['t_on'] = pd.to_datetime(al['t_on'], errors='coerce') s, e, days = _win(al['t_on'], 365) w = al[(al.t_on >= s) & (al.t_on <= e)] cnt = w.groupby('turbine').size().sort_values(ascending=False) if len(cnt) > 2: med = float(cnt.median()) top = cnt.head(3) named = [f'{t}({int(n)}条)' for t, n in top.items()] x = float(top.iloc[0]) / med if med else float('nan') refs = [dict(file=P.rel(ap), sha16=_sha16(ap))] code_top = (w[w.turbine == top.index[0]].groupby('code').size().sort_values(ascending=False) .head(3).to_dict()) if 'code' in w else {} add('RD-LEDGER-ALM-01', f'报警集中: {top.index[0]} {int(top.iloc[0])} 条(机群中位 {med:.0f} 条的 {x:.1f} 倍)', '候选', refs, named=named, systems=['报警面'], module='台账·报警', window_days=days, n_machines=int(w.turbine.nunique()), limitation='报警条数不等于缺陷严重度: 单码抖动(同一码反复报)会抬高条数; 未做同码基线归一。', falsify=f'若该台的高发码是全场共因(同一码在多台同量出现)或与限电/调度同向, 则该"集中"是工况不是缺陷 —— 判据: 逐码机群基率>{1 / max(len(cnt), 1):.0%}。', nature='口径性(计数面)', coverage_note=f'窗 {s:%Y-%m-%d}~{e:%Y-%m-%d}, 全员 {int(w.turbine.nunique())} 台', claim_class=f'计数面相对判据; 头部码 {code_top}') # ── ② 重复检修件 ──────────────────────────────────────────────────────────── wo = store / 'workorders.parquet' if wo.is_file(): d = pd.read_parquet(wo) comp_col = next((c for c in ('故障位置二级', '故障位置一级', '故障名称') if c in d.columns), None) tm = next((c for c in ('故障报出时间', '复位运行时间') if c in d.columns), None) unit = next((c for c in ('机组编号',) if c in d.columns), None) if comp_col and unit: d['_c'] = d[comp_col].astype(str).str.strip() g = d[d['_c'] != ''].groupby([unit, '_c']).size().sort_values(ascending=False) if len(g): (t0, c0), n0 = g.index[0], int(g.iloc[0]) tm_ = pd.to_datetime(d[tm], errors='coerce') if tm else None s, e, days = _win(tm_, 365) if tm_ is not None else ('', '', None) _norm = lambda x: (str(x) if str(x).upper().startswith('WTG') else f'WTG{int(str(x)):02d}' if str(x).isdigit() else str(x)) named = [f'{_norm(t)}·{c}({int(n)}次)' for (t, c), n in g.head(3).items()] add('RD-LEDGER-WO-01', f'重复检修: {_norm(t0)} 的「{c0}」累计 {n0} 次(窗内)', '候选', [dict(file=P.rel(wo), sha16=_sha16(wo))], named=named, systems=['传动链'], module='台账·工单', window_days=days, n_machines=int(d[unit].nunique()), limitation='检修台账记录的是"做过什么", 不等于"坏过几次": 同一部件的例行检查也会计次; ' '未做计划性/故障性停机区分(缺 work_type 字段)。', falsify='若这些次记录多为计划性检查/预防性更换, 则该"重复"无非计划性含义 —— 判据: 记录里标注为故障停机的比例。', nature='趋势性(重复发生)', coverage_note=f'窗 {s}~{e}' if s else '窗未标(无时间列)', claim_class='计数面相对判据; 部件名取台账原文') # ── ③ 停机损失集中 ───────────────────────────────────────────────────────── lm = store / 'loss_monthly.parquet' if lm.is_file(): d = pd.read_parquet(lm) st_col = 'state' if 'state' in d.columns else None if st_col and 'loss' in d.columns and 'turbine' in d.columns: stop = d[d[st_col].astype(str).str.contains('停机', na=False)] # ★单位: loss 列是 **kWh**(实测全场合计 4,428 万、停机态 2,273 万)—— 直接当 MWh 报会 # 把量级说大 1000 倍。这里统一换算成 MWh 再报, 分子分母同口径。 tot_stop = float(stop['loss'].sum()) / 1000.0 per = (stop.groupby('turbine')['loss'].sum() / 1000.0).sort_values(ascending=False) if len(per) and tot_stop > 0: sh = float(per.iloc[0]) / tot_stop named = [f'{t}({v:,.0f}MWh)' for t, v in per.head(3).items()] add('RD-LEDGER-LOSS-01', f'停机损失集中: {per.index[0]} 占全场停机损失 {sh:.0%}({per.iloc[0]:,.0f} MWh, 全场停机 {tot_stop:,.0f} MWh)', '参考', [dict(file=P.rel(lm), sha16=_sha16(lm))], named=named, systems=['发电量'], module='台账·损失', window_days=int(d['month'].nunique()) * 30 if 'month' in d else None, n_machines=int(d['turbine'].nunique()), limitation='损失是"损失了多少电", 不是"为什么损失"; 含调度令停机(非缺陷)与限电, 需与故障类分开看。', falsify='若该台损失主要来自调度令/限电(状态列已区分)而非故障停机, 则与设备无关。', nature='口径性(电量面)', coverage_note='按月汇总, 月份数取自产物') # ── ④ 温度离群台 ──────────────────────────────────────────────────────────── tb = store / 'temp_bins.parquet' if tb.is_file() and (m5 / 'fusion_38.csv').is_file(): d = pd.read_parquet(tb) if {'turbine', 'channel', 'pbin', 'med', 'n'}.issubset(d.columns): d = d[d['n'] >= 20] # 样本量门 (n 太少的中位不稳) if len(d): piv = d.pivot_table(index=['channel', 'pbin', 'win'], columns='turbine', values='med') dev = piv.sub(piv.median(axis=1), axis=0).abs() worst = dev.max(axis=1).sort_values(ascending=False) if len(worst): (ch, pb, win), val = worst.index[0], float(worst.iloc[0]) t0 = str(dev.loc[worst.index[0]].idxmax()) add('RD-LEDGER-TEMP-01', f'温度相对离群: {t0} 的 {ch} 在 {pb} 功率档偏离机群中位 {val:.1f}K(工况窗 {win})', '候选', [dict(file=P.rel(tb), sha16=_sha16(tb))], named=[f'{t0}·{ch} +{val:.1f}K'], systems=['温度面'], module='台账·温度 NBM', window_days=None, n_machines=int(d['turbine'].nunique()), limitation='同工况中位比对的是"机群", 若全机群共因升温和(环境/冷却水温)则判据失效; ' '未接绝对限值(无 ISO/厂商温度锚)。', falsify='若该通道全场中位同期同向上升(共模), 则该台"离群"不成立 —— 判据: 逐月机群中位趋势。', nature='趋势性(个体偏离)', coverage_note=f'{win} 工况段, 样本门 n≥20', claim_class='同工况机群残差; 功率档分箱') # ── ⑤ 功率曲线偏离台 ─────────────────────────────────────────────────────── pc = store / 'powercurve_dev.parquet' if pc.is_file(): d = pd.read_parquet(pc) if 'dev_w' in d.columns: d2 = d.assign(_a=pd.to_numeric(d['dev_w'], errors='coerce').abs()).sort_values('_a', ascending=False) r0 = d2.iloc[0] # ★单位: dev_w 是 **MW**(实测 -0.147 → -147 kW)—— 直接标 kW 会把量级说小 1000 倍。 _dv = float(pd.to_numeric(pd.Series([r0.get('dev_w')]), errors='coerce').iloc[0]) * 1000.0 _cnt = (d['判别'].astype(str).value_counts().to_dict() if '判别' in d.columns else {}) _bad = {k: v for k, v in _cnt.items() if k not in ('—', 'nan', '')} add('RD-LEDGER-PC-01', f'功率曲线偏离: {r0.get("turbine")} 同风速下功率偏差 {_dv:+,.0f} kW(判别: {r0.get("判别", "—")})' + (f'; 全场判别分布 {_bad}' if _bad else ''), '参考', [dict(file=P.rel(pc), sha16=_sha16(pc))], named=[f'{r0.get("turbine")} {float(r0.get("dev_w")):,.0f}kW'], systems=['发电量', '功率曲线'], module='台账·功率曲线', window_days=None, n_machines=int(len(d)), limitation='同风速残差受机位风资源(尾流/扇区)影响; 未做机位修正, 故只作"待核"不作缺陷结论。', falsify='若该台所处扇区长期受尾流(邻机偏置)影响, 残差可由机位解释 ⇒ 非机组问题。', nature='口径性(机位/风资源)', coverage_note='逐台一行(残差+判别列)') # ── ⑥ 振动面过闸线(观澜自算六层链) ─────────────────────────────────────── f38 = m5 / 'fusion_38.csv' l6p = m5 / 'model_run_l6.parquet' cand = pd.DataFrame() if l6p.is_file(): l6 = pd.read_parquet(l6p) if '定级' in l6.columns and len(l6): cand = l6[l6['定级'].astype(str).str.contains('候选|准定论|定论', na=False)] if len(cand): named = [f"{r['台']}·{str(r['线'])[:22]} {r['定级']}" for _, r in cand.head(5).iterrows()] add('RD-LEDGER-VIB-01', f"振动面过闸线: {len(cand)} 条特征线达「候选」及以上({', '.join(n.split(' ')[0] for n in named[:3])})", '候选', [dict(file=P.rel(l6p), sha16=_sha16(l6p))], named=named, systems=['振动面'], module='振动·六层链 L6 过闸线', window_days=None, n_machines=int(cand['台'].nunique()), limitation='机制未定 ⇒ **不命名部件**(滚道/圈侧须实物或换件闭环); 峰值拾取为本器口径' '(目标频率 ±2 bin 取最大), 无正样本锚 ⇒ 封顶「候选」, 不出「定论」。', falsify='若复测同一判据回落(或拆检未见对应损伤), 则该候选撤回 —— 判据: 下一窗同线定级 + 实物证据。', nature='趋势性(证据面)', coverage_note='L6 过闸线表; 闸分布见 m5_cms_tcm/model_run_summary.json', claim_class='观澜自算(非振动线正裁)') if not len(cand) and f38.is_file(): d = pd.read_csv(f38, dtype=str) if '融合' in d.columns: bad = d[d['融合'].astype(str).str.contains('候选|定论|准定论', na=False)] if len(bad): named = [f"{r['台']}·{r['融合']}" for _, r in bad.head(5).iterrows()] add('RD-LEDGER-VIB-01', f"振动面融合级: {len(bad)} 台达「候选」及以上({', '.join(named[:3])})", '候选', [dict(file=P.rel(f38), sha16=_sha16(f38))], named=named, systems=['振动面'], module='振动·六层链融合', window_days=None, n_machines=int(len(d)), limitation='机制未定 ⇒ 不命名部件; 无正样本锚 ⇒ 封顶「候选」。', falsify='复测同一判据回落或拆检未见损伤即撤回。', nature='趋势性(证据面)', coverage_note='融合表逐台一行', claim_class='观澜自算融合(非振动线正裁)') return out def build(farm: str | None = None, dry=False) -> int: root = P.out_root(farm) store, m5 = P.store(farm), P.m5(farm) cl = claims(store, m5) if not cl: print('[X] 一条 claim 也没算出来 —— 台账产物不在位? 先跑重算链(②三门台账 / ③SCADA 侧)') return 2 if dry: for c in cl: print(f' [{c["verdict"]:6}] {c["id"]:20} {c["title"][:70]}') print(f' (dry-run, 共 {len(cl)} 条;未写文件)') return 0 t0 = time.strftime('%Y-%m-%d %H:%M:%S') findings = dict( schema='guanlan-sop-findings/v1', generated_by='scripts/sop_findings_from_ledger.py (用户令 2026-09-19: 从重算台账生成 claim)', built=t0, n=len(cl), note='本件由**观澜自算**(重算台账 → claim), 不含人工裁决; 逐条 source_refs 记件与 sha16, ' 'mandatory_limitation/falsifiability 如实写明缺证据与证伪判据。', findings=[dict( id=c['id'], title=c['title'], verdict=c['verdict'], coverage=dict(window_days=c['window_days'], window_days_note=c['coverage_note'], per_day_evidence=None, n_machines=c['n_machines']), claim_class=c['claim_class'], module=c['module'], mandatory_limitation=c['mandatory_limitation'], falsifiability=c['falsifiability'], problem_nature=c['problem_nature'], source_chapter='重算台账', evidence=[c['title']], named=c['named'], systems=c['systems'], source_refs=c['source_refs']) for c in cl]) p1 = root / 'sop' / 'findings.json' p1.parent.mkdir(parents=True, exist_ok=True) p1.write_text(json.dumps(findings, ensure_ascii=False, indent=1), encoding='utf-8') base = root / 'paradigm_r1' / 'experiments' # E5 底稿: 模块映射(契约的 scope.module / module_map_how) (base / 'E5_coverage_table').mkdir(parents=True, exist_ok=True) (base / 'E5_coverage_table' / '底稿.json').write_text(json.dumps(dict( schema='guanlan-paradigm-e5/v1', generated_by='scripts/sop_findings_from_ledger.py (自算)', note='模块映射: 每条 claim 归到哪个分析模块(观澜自算口径)', finding_rows=[dict(i=i, module=c['module'], map_how='ledger:' + c['module']) for i, c in enumerate(cl)]), ensure_ascii=False, indent=1), encoding='utf-8') # E3 底稿: 归宿提议(契约的 closure.machine_proposed / human=空) (base / 'E3_candidate_closure').mkdir(parents=True, exist_ok=True) (base / 'E3_candidate_closure' / '底稿.json').write_text(json.dumps(dict( schema='guanlan-paradigm-e3/v1', generated_by='scripts/sop_findings_from_ledger.py (自算)', note='归宿提议: 机器只提"下一步该做什么", 人裁字段留空(不在本器职责内)', rows=[dict(i=i, human=None, proposed=('复测同一判据并核实物证据' if c['verdict'] == '候选' else '例行跟踪, 不单独排程')) for i, c in enumerate(cl)]), ensure_ascii=False, indent=1), encoding='utf-8') # E8 底稿_s0: 系统/台号(契约的 scope.systems/units 与 sample_level.named) (base / 'E8_rudong_rebuild').mkdir(parents=True, exist_ok=True) (base / 'E8_rudong_rebuild' / '底稿_s0.json').write_text(json.dumps(dict( schema='guanlan-paradigm-e8-s0/v1', generated_by='scripts/sop_findings_from_ledger.py (自算)', note='系统/台号: 由 claim 的点名台与系统列直出', findings_rows=[dict(i=i, n_tid=len(c['named']), tids=[n.split('(')[0].split('·')[0] for n in c['named']], systems=c['systems'], units=c['named']) for i, c in enumerate(cl)]), ensure_ascii=False, indent=1), encoding='utf-8') print(f'已写 {P.rel(p1)}({len(cl)} 条 claim)+ 3 份底稿(E5/E3/E8_s0)') for c in cl: print(f' [{c["verdict"]:6}] {c["id"]:20} {c["title"][:64]}') try: from src import derived_manifest as DM rels = {f: 'scripts/sop_findings_from_ledger.py (重算台账 → claim; 观澜自算)' for f in OUT_FILES} DM.record(root, rels, by='sop_findings_from_ledger') print(' 已自登记 → _derived_manifest.json') except Exception as e: print(f' [i] 自登记跳过: {type(e).__name__}: {e}') return 0 def status(farm: str | None = None) -> int: root = P.out_root(farm) p = root / 'sop' / 'findings.json' if not p.is_file(): print(json.dumps(dict(path=P.rel(p), exists=False, note='不在位 —— 门户结论段/问答/报告面会 503 或停在旧数; 跑本器可由重算台账生成'), ensure_ascii=False)) return 0 try: d = json.loads(p.read_text(encoding='utf-8')) except Exception as e: print(json.dumps(dict(path=P.rel(p), exists=True, err=str(e)[:120]))) return 1 print(json.dumps(dict(path=P.rel(p), exists=True, n=d.get('n'), built=d.get('built'), generated_by=d.get('generated_by'), kind=('观澜自算' if 'sop_findings_from_ledger' in str(d.get('generated_by')) else '人裁底稿')), ensure_ascii=False)) return 0 def main() -> int: ap = argparse.ArgumentParser(description='事实契约的 claim 生成端(重算台账 → findings/底稿)') ap.add_argument('--farm', default=None) ap.add_argument('--dry-run', action='store_true') ap.add_argument('--status', action='store_true') a = ap.parse_args() if a.status: return status(a.farm) return build(a.farm, a.dry_run) if __name__ == '__main__': for _s in (sys.stdout, sys.stderr): try: _s.reconfigure(errors='replace') except Exception: pass sys.exit(main())