| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354 |
- #!/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())
|