farm_pipeline.py 38 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582
  1. # -*- coding: utf-8 -*-
  2. """farm_pipeline.py — ON-2 编排器: contract_gate → derive_10min → discriminators → findings.json.
  3. 把 Batch2 三块积木 (CL-1 清洗门 / §3.4 10min 派生 / §4 稳定判别器) 串成 farm-agnostic 管线,
  4. 新场 onboard 后一条命令跑通"清洗→派生→判别→findings", 替代 per-farm 重写分析脚本 (DP-10 类).
  5. 设计:
  6. - 纯编排, 不读盘 (raw_df 由调用方按场加载; 库不 hardcode F: 路径/编码/列名 — 那是 per-farm 噪声)。
  7. - 判别器→finding 的映射经"适配器" (_finding_*): 套 SOP §0.2 降级表 + §4 判据出 verdict/claim_class。
  8. - 出 findings.json 前逐条过 schemas.validate_finding (六枚举/coverage/falsifiability), issue 不静默 — 落 _validation。
  9. 接口: run_pipeline(raw_df, config) -> dict{clean_report, df10, discriminators, findings, validation}。
  10. config 见 docstring of run_pipeline; dogfood 实例见 scripts/dogfood_farm_pipeline_jtxq.py。
  11. """
  12. from __future__ import annotations
  13. try:
  14. from app_common.app_common_guanlan.api import install_root as _install_root
  15. except ImportError: # 理论不可达;包结构异常时回退到按位置上跳
  16. from pathlib import Path as _P
  17. def _install_root(_f): return _P(_f).resolve().parents[4]
  18. from src import paths as P
  19. import json
  20. from datetime import datetime
  21. from pathlib import Path
  22. import pandas as pd
  23. from .contract_gate import apply_contract_gate
  24. from .derive10min import derive_10min
  25. from . import discriminators as D
  26. from . import schemas
  27. ROOT = _install_root(__file__)
  28. # ============================================================ 判别器 → finding 适配器
  29. # 每个适配器: 入 (判别器结果 dict, coverage dict) → 出 findings_item (过 validate_finding)。
  30. # verdict 取舍套 §0.2 降级表; 不在此造数字, 只把判别器结构化结论翻译成 finding 语言。
  31. def _finding_curtail(res, cov, *, title='限电 vs 故障 同步性判别'):
  32. """§4.3 限电同步性 → finding。三证(同步台数/频率偏移/季节集中)齐 = 定论电网限电。"""
  33. v = res.get('verdict')
  34. if v == '无限电候选':
  35. verdict = '参考'
  36. fals = '若后续窗出现 ≥半数台同时 curtail 态 则存在限电候选, 需重判'
  37. elif v and str(v).startswith('疑误派生'):
  38. verdict = 'INSUFFICIENT'
  39. fals = ('curtail 信号常年on(>50%全时)→疑派生阈值高于铭牌功率设定(正常次额定误标限电); '
  40. '须核派生口径(corr(P,ws)是否随风/真设定值)后重判, 当前不可判限电')
  41. elif v and str(v).startswith('限电信号存疑'):
  42. verdict = 'INSUFFICIENT'
  43. fals = ('curtail 标志在匹配中风速段未压低功率(setpoint 标签 artifact, pset 多个正常值都满发); '
  44. '须核派生口径或换真限电信号(真停机限电应在同风段把功率压数百kW)后重判, 当前不可判限电')
  45. elif v == '电网同步限电':
  46. # 三证: sync_fraction(主证) + freq_offset(脱网) + season_peak(集中). 主证强 + 任一旁证 → 定论
  47. strong = (res.get('sync_fraction') or 0) >= 0.5
  48. corrob = (res.get('freq_offset_hz') not in (None, 0)) or ((res.get('season_peak_ratio') or 0) > 2)
  49. verdict = '定论' if (strong and corrob) else '准定论·预警'
  50. fals = '若 curtail 窗实为 per 机散布(中位同时台数≈1)∨电网频率无偏移 则非电网限电(是 per 机故障)'
  51. elif v == '电网限电(部分调度)':
  52. # 2026-07-26 新增档: sync 不过半 (电网只调度部分机组) 但旁证 >=2 票。
  53. # 不给"定论" —— 主证(同步过半)未成立, 靠旁证推断, 上限=准定论·预警。
  54. verdict = '准定论·预警'
  55. fals = ('若中位同时台数降到 ≈2(真散布故障量级) ∨ 季节集中消失(峰/均<2) ∨ 电网频率无偏移 '
  56. '则改判 per 机故障; 若后续出现 ≥半数台同步则升为定论级电网限电')
  57. elif v and str(v).startswith('INSUFFICIENT'):
  58. # 2026-07-26 新增档: 落在"部分调度 vs 散布故障"模糊带 (median_simul≈3 且旁证不足)。
  59. # ★旧版会把这类强行判成 per机故障(散布) -> 去追不存在的故障。宁可显式说不知道。
  60. verdict = 'INSUFFICIENT'
  61. fals = ('当前证据不足以分"部分电网调度"与"per机散布故障": 中位同时台数落在模糊带(≈3), '
  62. '旁证(季节集中/频率偏移)不足。须补**独立限电台账**或**调度指令记录**才可判; '
  63. '不得据此报机组缺陷, 亦不得据此报电网限电')
  64. else: # per机故障(散布)
  65. verdict = '候选'
  66. fals = '若 curtail 候选窗全场同步(≥半数台同时)∨频率偏移显著 则改判电网限电'
  67. return {
  68. 'title': title, 'verdict': verdict, 'claim_class': '限电',
  69. 'coverage': {**cov, 'per_day_evidence': res.get('evidence', '')},
  70. 'basis': {'mast_usability': 'NONE', 'window_months': cov.get('window_months'),
  71. 'status_source': '状态码 curtail 候选态占比(per 10min 窗) + 电网频率 + 月度集中'},
  72. 'falsifiability': fals,
  73. 'detail': f"判别器 curtail_vs_fault_sync: verdict={v}; "
  74. f"同步窗 {res.get('n_curtail_windows')} / 总机 {res.get('n_total')}",
  75. 'affected_machines': None, 'problem_nature': '外部(电网)', # ① 源头落盘 (限电=场级外部, 无问题机; render 渲"全场(电网级)")
  76. 'discriminator': res,
  77. }
  78. def _finding_yaw(res, cov, *, title='静态偏航对风(对风角度列)'):
  79. """§4.4① 对风角度列可用性 → finding。不可用(任意零点 artifact)→ INSUFFICIENT。"""
  80. usable = res.get('usable')
  81. affected, reco = None, None # ① 源头落盘 (偏置台/校正建议)
  82. if usable is None: # 有效台<2: 数字字段为 None, evidence 只用 reason 不内插 None
  83. verdict, fals = 'INSUFFICIENT', f"有效台<2 不足以判 inter-turbine 散布: {res.get('reason')}"
  84. ev = res.get('reason', '有效台<2 (per台样本不足)')
  85. elif usable is False:
  86. verdict = 'INSUFFICIENT'
  87. fals = f"列不可用(任意零点 artifact): {res.get('reason')}; 真对风需 native 对风角列或现场独立量"
  88. ev = res.get('reason', '列不可用 (任意零点 artifact)')
  89. else:
  90. # 列可用. ⚠ "列可用"≠"全场对中健康": 单台真失配(|中位|>5°)不可被"定论健康"掩盖(Critic逮昌平坳#36+8.8°)
  91. ptm = res.get('per_turbine_median') or {}
  92. offset = sorted(((k, v) for k, v in ptm.items() if abs(v) > 5.0), key=lambda x: -abs(x[1]))
  93. if offset: # 有真失配台 → 候选偏置台 (可回收电量), 不笼统判健康
  94. verdict = '候选'
  95. affected = [k for k, _ in offset] # ① 偏置台(>5°)直接落盘
  96. reco = '现场独立量(LiDAR/偏航阶跃)核实并下发零点校正'
  97. top = ', '.join(f'#{k} {v:+.1f}°' for k, v in offset[:3])
  98. ev = (f"对风角列可用(极差 {res.get('inter_turbine_spread')}°), 多数台对中聚集; "
  99. f"但 {len(offset)} 台中位偏置>5°: {top} = 候选偏航偏置台(单台静态偏置~cos²–³ 可回收电量, 须现场独立量校正)")
  100. fals = ('若偏置台经现场独立量(LiDAR/偏航阶跃)证实<3° 则为列标定噪声非真偏置; '
  101. '若 inter-turbine 极差>30°∨|中位|>45° 则列不可用(任意零点)')
  102. else: # 全台 ±5° 内 → 相对对中健康 (定论级 relative; 绝对零点仍需 LiDAR)
  103. verdict = '定论'
  104. fals = '若某台对风角 |中位|>45° ∨ inter-turbine 极差>30° (任意零点阈, SOP §4.4①) 则偏置台/列不可用, 需现场独立量核'
  105. ev = (f"对风角列可用; 全台 per机中位聚集 ±5° 内 (极差 {res.get('inter_turbine_spread')}° / "
  106. f"最大|中位| {res.get('max_abs_median')}°, n={res.get('n_turbines')}台); {res.get('reason')}")
  107. return {
  108. 'title': title, 'verdict': verdict, 'claim_class': '相对排名',
  109. 'coverage': {**cov, 'per_day_evidence': ev},
  110. # ⚠ 2026-07-16 去硬编码假声明 (独立critic af45d34f: 原两字段是模板常量非per场派生, jishan_p1自相矛盾坐实):
  111. # mast_usability 判别器根本不查mast(_run_yaw不接mast参数)→ 勿据此判无塔(5/11场实有测风塔);
  112. # status_source 未验证native vs wind_dir−nacelle重构 → 真列源须读per场清洗层映射, 不在此断言。
  113. 'basis': {'mast_usability': 'NOT_CHECKED (判别器不接mast参数; 勿判无塔 — 部分场实有测风塔)',
  114. 'window_months': cov.get('window_months'),
  115. 'status_source': '对风角度列 per机中位 (列源=清洗层映射, 未在此验证native/重构 — 见per场清洗config)',
  116. 'curtail_strip': {'method': '对风角与限功率正交(限电降功率不改对中, §4.3)', 'n_removed': 'N/A (正交, 无需剔)'}},
  117. 'falsifiability': fals,
  118. 'detail': '相对对中健康=per机中位聚集; ⚠绝对偏航零点(可下发校正)仍须 LiDAR/测风塔/偏航阶跃(SOP §4.4 铁律)',
  119. 'affected_machines': affected, 'problem_nature': '性能', 'recommendation': reco, # ① 源头落盘
  120. 'discriminator': res,
  121. }
  122. def _finding_nbm(res, cov, *, title='温度 NBM 残差第二道'):
  123. """§4.6 温度 NBM 残差 → finding。flagged 非空=候选(需绝对温度+持续性二判); 空=参考(筛过非保证健康)。"""
  124. flagged = res.get('flagged') or []
  125. affected, reco = None, None # ① 源头落盘 (候选偏热台)
  126. if not res.get('per_turbine'):
  127. verdict, ev = 'INSUFFICIENT', res.get('note', '样本不足, NBM 无法拟合')
  128. fals = '若补足样本(per台≥阈)后拟合稳定则可判'
  129. elif flagged:
  130. verdict = '候选'
  131. affected = list(flagged) # ① 候选偏热台直接落盘
  132. reco = '现场核30秒实际温度(非仅告警计数) + 盯持续 NBM z>2 偏热台'
  133. top = res['per_turbine'][0]
  134. ev = (f"NBM 残差 fleet MAD-z 升候选 {len(flagged)} 台: {flagged}; "
  135. f"头名 {top['tid']} z={top['z']} (残差均 {top['resid_mean']}°, EWMA超 {top.get('ewma_exc')})")
  136. fals = '若候选台绝对温度未超同型分位 ∧ EWMA 持续超占比低 则为基线噪声非真发热(防 C 系低基线污染)'
  137. else:
  138. verdict = '参考'
  139. ev = f"NBM 残差第二道无 z≥{2.0} 候选台 (n={res.get('n_total')}台); 残差离群筛过"
  140. fals = '若个别台绝对温度持续超同型高分位则仍需单独绝对判(残差筛过≠绝对健康)'
  141. # §4.6c 参考通道守卫: 死 amb 台已剔出/降级 INSUFFICIENT (防"参考冻结"伪装"部件热", F0810001022 戒)
  142. if res.get('ref_dead'):
  143. ev += (' | ⚠参考(舱外温度)冻结剔出 ' + str(len(res['ref_dead'])) + ' 台(INSUFFICIENT, 非候选): '
  144. + '; '.join(f"{x['tid']}({x['reason']})" for x in res['ref_dead']))
  145. if res.get('ref_fallback'):
  146. ev += ' | 参考死→备选顶替: ' + '; '.join(f"{x['tid']}→{x['col']}" for x in res['ref_fallback'])
  147. return {
  148. 'title': title, 'verdict': verdict, 'claim_class': '相对排名',
  149. 'coverage': {**cov, 'per_day_evidence': ev},
  150. 'basis': {'mast_usability': 'NONE', 'window_months': cov.get('window_months'),
  151. 'status_source': 'pooled OLS 残差(comp~P+rotor+Tmed) fleet MAD-z + EWMA 持续性',
  152. 'curtail_strip': {'method': 'NBM 回归 P 项吸收负荷(含限电降功率)共模, 限电效应回归控制非剔除(§4.3/§4.6)', 'n_removed': 'N/A (回归控制)'}},
  153. 'falsifiability': fals,
  154. 'detail': '残差第二道(吸收功率/转速/季节共模)+绝对温度双判; 单冷 outlier 不驱动 corr(SOP §4.6)',
  155. 'affected_machines': affected, 'problem_nature': '可靠性/部件', 'recommendation': reco, # ① 源头落盘
  156. 'discriminator': res,
  157. }
  158. def _finding_ntf(res, cov, *, title='机舱风 NTF / 自由来流达成率'):
  159. """§4.5 NTF + 空间代表性 → finding。单 mast 非空间代表 → 绝对达成只报区间(INSUFFICIENT 绝对)。"""
  160. rep = res.get('spatial_representative')
  161. mast = str(cov.get('mast_usability', 'NONE')).upper()
  162. # 绝对达成(机舱风 self_ref)须双门: mast PASS ∧ 空间代表. 任一不齐 → INSUFFICIENT 绝对
  163. # (空间代表只证机舱/自由来流比值内部一致, 与测风塔本身物理是否 PASS 无关 — 审查 HIGH-1)
  164. if rep is True and mast == 'PASS':
  165. verdict = '准定论·预警'
  166. fals = '若 NTF 随扇区极差>0.1 则为方向加权混合非机型 NTF, 绝对达成降区间'
  167. else:
  168. verdict = 'INSUFFICIENT'
  169. fals = (f'绝对达成须 mast PASS ∧ 空间代表双门 (当前 mast={mast}/空间代表={rep}) → 只报区间; '
  170. '若两门齐(各机比∈[0.9,1.1]∧扇区极差≤0.1∧mast PASS)则可定量')
  171. ev = (f"NTF(机舱/自由来流)中位 {res.get('ntf_median')}; reads_low={res.get('reads_low')}; "
  172. f"空间代表={rep}; {res.get('caveat')}")
  173. return {
  174. 'title': title, 'verdict': verdict, 'claim_class': '绝对性能',
  175. 'coverage': {**cov, 'per_day_evidence': ev},
  176. 'basis': {'mast_usability': cov.get('mast_usability', 'NONE'),
  177. 'window_months': cov.get('window_months'),
  178. 'status_source': '测风塔自由来流 vs 机舱风速计 NTF',
  179. 'curtail_strip': {'method': '功率曲线达成率须剔限电(§4.3 限电=性能头号混杂); 限电态由 curtail 判别器标记剔', 'n_removed': '限电样本(601态等, 见限电finding)'}},
  180. 'falsifiability': fals,
  181. 'detail': '绝对达成率天然 self_ref(机舱风), 必自由来流修正; 单 mast 空间代表性 caveat(SOP §4.5)',
  182. 'affected_machines': None, 'problem_nature': '性能', # ① 源头落盘 (场级数据局限, 无问题机)
  183. 'recommendation': '补测风塔/独立风速量 → 绝对达成率升定论',
  184. 'discriminator': res,
  185. }
  186. # ──────────── ④部件健康 适配器 (2026-06-20 loop训练抽库 + 真数据RV-2 + 独立Critic HOLD)。
  187. # 全部 **affected_machines 从 flagged 取, 禁 verdict-first** (Critic 逮 cooling verdict 会被整体 load_driven 淹没真衰减台)。
  188. def _finding_matched_load(res, cov, *, title='温度 matched-load 双判 (NBM盲区兜底)'):
  189. """§4.6/3.1 matched-load 双判 → finding。flagged(同功率bin ΔT-z ∨ 多部件簇状)=候选; 空=参考(筛过非健康)。"""
  190. flagged = res.get('flagged') or []
  191. affected, reco = None, None
  192. if not res.get('n_components'):
  193. verdict, ev = 'INSUFFICIENT', res.get('note', '无有效部件温列/功率带样本不足')
  194. fals = '若补足部件温列 + 功率带重叠样本则可判'
  195. elif flagged:
  196. verdict, affected = '候选', list(flagged)
  197. reco = '现场IR/CMS核 + 与 nbm_residual(极值)交叉对照 (matched-load 兜底中度/多部件/簇状升温)'
  198. top = res['per_turbine'].get(flagged[0], {})
  199. ev = f"matched-load 双判升候选 {len(flagged)} 台: {flagged}; 头名 {flagged[0]} {top.get('reason')}"
  200. fals = '若候选台同功率bin ΔT 经现场核为传感器标定差/无绝对失控 则非真发热; ⚠非重叠功率带台 ΔT路径致盲(须fleet功率带重叠)'
  201. else:
  202. verdict = '参考'
  203. ev = f"matched-load 双判无升候选 (n_components={res.get('n_components')}); NBM盲区(中度/多部件/簇状)兜底筛过"
  204. fals = '若个别台同功率bin多部件持续偏热则仍需单独 matched-load 核 (筛过≠绝对健康)'
  205. if res.get('dead_channel'): # §4.6c 冻结部件温列已逐台剔出排名+簇状 (F0810001022 戒)
  206. ev += ' | ⚠§4.6c 冻结部件温列剔出: ' + '; '.join(f"{x['comp']}:{x['tids']}" for x in res['dead_channel'])
  207. return {
  208. 'title': title, 'verdict': verdict, 'claim_class': '相对排名',
  209. 'coverage': {**cov, 'per_day_evidence': ev},
  210. 'basis': {'mast_usability': 'NONE', 'window_months': cov.get('window_months'),
  211. 'status_source': '同功率bin 部件温−fleet中位 ΔT-z + 多部件簇状同热占比',
  212. 'curtail_strip': {'method': 'matched-load 同功率bin比已控负荷(含限电降功率), bin内消(§4.6)', 'n_removed': 'N/A (同bin控制)'}},
  213. 'falsifiability': fals,
  214. 'detail': 'NBM极值检测器盲区兜底(逮中度+4-6°C/多部件/簇状, 凉水泉2208F戒); ΔT相对fleet非OEM绝对阈, 上限候选; 真数据RV-2复现凉水泉',
  215. 'affected_machines': affected, 'problem_nature': '可靠性/部件', 'recommendation': reco,
  216. 'discriminator': res,
  217. }
  218. def _finding_cooling(res, cov, *, title='冷却系统效能'):
  219. """§4.1/3.5 冷却效能 → finding。⚠ affected 必从 flagged(C层真衰减台)取, 禁 verdict-first(整体可能 load_driven 淹没)。"""
  220. flagged = res.get('flagged') or []
  221. affected, reco = None, None
  222. if res.get('verdict') == 'INSUFFICIENT' or not res.get('per_turbine'):
  223. verdict, ev = 'INSUFFICIENT', res.get('evidence', '样本不足')
  224. fals = '若补足同功率bin样本则可判'
  225. elif flagged:
  226. verdict, affected = '候选', list(flagged)
  227. reco = '散热器/水泵/风扇检修 + 30秒实温 vs F级155°C 核(非仅告警计数)'
  228. ev = res.get('evidence', '')
  229. fals = '若候选台同功率温偏经现场核为传感器/工况差 则非冷却衰减; 真过热须30秒实温超部件告警阈坐实'
  230. else:
  231. verdict = '参考' # ⚠Critic: flagged=[] = 同功率温检测下限内无可检出 ≠ 冷却健康定论(小样本可漏中度偏热台, A11型)
  232. ev = (res.get('evidence', '') + f" | causal={res.get('causal')}/cooling_type={res.get('cooling_type')}; "
  233. "⚠flagged=[]=检测下限内无可检出, 非'冷却健康'定论(小样本可漏中度偏热台)")
  234. fals = '若同功率bin某台温水平系统偏高(z≥2∧ΔT≥2°C)则改候选; "高温降容"若同风速高温组功率反更高则负荷驱动非衰减'
  235. if res.get('dead_channel'): # §4.6c 冻结冷却水温台已剔出排名 (F0810001022 戒)
  236. ev += ' | ⚠§4.6c 冻结冷却水温台剔出排名(INSUFFICIENT): ' + str(res['dead_channel'])
  237. return {
  238. 'title': title, 'verdict': verdict, 'claim_class': '相对排名',
  239. 'coverage': {**cov, 'per_day_evidence': ev},
  240. 'basis': {'mast_usability': 'NONE', 'window_months': cov.get('window_months'),
  241. 'status_source': '同功率bin冷却温−fleet中位 MAD-z(C层) + 因果方向守卫(高温组功率) + 释放曲线分型',
  242. 'curtail_strip': {'method': '同功率bin比已控负荷; 因果守卫剔负荷驱动温升(非冷却衰减)', 'n_removed': 'N/A (同bin控制+因果守卫)'}},
  243. 'falsifiability': fals,
  244. 'detail': '因果反置守卫(功率↑→温升=负荷驱动非降容, 郭家店1338→197MWh戒) + 低温冷启vs高温降容释放曲线; affected从C层flagged取(禁verdict-first, Critic); 真数据RV-2复现郭家店',
  245. 'affected_machines': affected, 'problem_nature': '可靠性/部件', 'recommendation': reco,
  246. 'discriminator': res,
  247. }
  248. def _finding_vibration(res, cov, *, title='振动 level 趋势 (粗筛)'):
  249. """§4.6/3.2 振动 level 趋势 → finding。flagged=[(测点,台)] → 候选(affected=去重台号); 退化列/无效→INSUFFICIENT。"""
  250. flagged = res.get('flagged') or []
  251. affected, reco = None, None
  252. if str(res.get('verdict', '')).startswith('INSUFFICIENT'):
  253. verdict, ev = 'INSUFFICIENT', res.get('caveat', '退化列/样本不足/<2台')
  254. fals = '若补足有效振动列(非退化)+≥2台则可判'
  255. elif flagged:
  256. verdict = '候选'
  257. affected = sorted({t for _c, t in flagged})
  258. reco = 'CMS高频频谱诊断定位(轴承BPFO/BPFI/齿轮啮合/叶片1P); 无CMS→建议加装振动测点'
  259. ev = '振动 level 离群/爬升候选 (分测点 fleet MAD-z / 月度斜率持续): ' + ', '.join(f'{c}:{t}' for c, t in flagged[:4])
  260. fals = '若候选台同测点同工况健康台同样高则共模(机型/工况)非本台; 仅10min level无频谱→出不了部件根因(只候选)'
  261. else:
  262. verdict, ev = '参考', '分测点振动 level 各测点平稳 (无 z≥2/爬升); 10min level 粗筛过'
  263. fals = '若某测点 level z≥2∧月度爬升持续则改候选; level粗筛过≠CMS频谱健康'
  264. return {
  265. 'title': title, 'verdict': verdict, 'claim_class': '相对排名',
  266. 'coverage': {**cov, 'per_day_evidence': ev},
  267. 'basis': {'mast_usability': 'NONE', 'window_months': cov.get('window_months'),
  268. 'status_source': '分测点 fleet null MAD-z(禁混池) + 月度Theil-Sen斜率+持续性',
  269. 'curtail_strip': {'method': '振动level与限功率正交(限电不直接改振动level)', 'n_removed': 'N/A (正交)'}},
  270. 'falsifiability': fals,
  271. 'detail': '10min level 粗筛(verdict上限预警); 部件级须CMS频谱; 分测点禁混池+退化列守卫+MAD地板(SOP §4.6 R5)',
  272. 'affected_machines': affected, 'problem_nature': '可靠性/部件', 'recommendation': reco,
  273. 'discriminator': res,
  274. }
  275. def _finding_blade(res, cov, *, title='叶片不平衡 1P (CMS/秒级层)'):
  276. """§4.1/3.2b 叶片1P 重力感知门 → finding。raw 通道→INSUFFICIENT(重力/未标定); 补偿通道 flagged→候选。"""
  277. flagged = res.get('flagged') or []
  278. affected, reco = None, None
  279. if 'INSUFFICIENT' in str(res.get('verdict', '')):
  280. verdict, ev = 'INSUFFICIENT', res.get('caveat', 'raw通道1P不可信(重力/未标定)')
  281. fals = '若用切向/已补偿通道(gravity_compensated=True)或矢量扣平衡参考窗重力相量则可跑fleet-z; raw挥舞/摆振向1P=重力/标定不一不可作不平衡(路由modal-z)'
  282. elif flagged:
  283. verdict, affected = '候选', list(flagged)
  284. reco = '现场叶片检查(质量配平/气动外形) + CMS 持续监视'
  285. ev = f"已补偿通道 1P fleet MAD-z 升候选 {len(flagged)} 台: {flagged}"
  286. fals = '若候选台 1P 经阶次跟踪/现场核落回fleet则原是泄漏/标定artifact; 极低1P台疑死/低增益通道(非健康)'
  287. else:
  288. verdict, ev = '参考', '已补偿通道 1P fleet 中位无离群'
  289. fals = '若某台 1P z≥2 则改候选'
  290. return {
  291. 'title': title, 'verdict': verdict, 'claim_class': '相对排名',
  292. 'coverage': {**cov, 'per_day_evidence': ev},
  293. 'basis': {'mast_usability': 'NONE', 'window_months': cov.get('window_months'),
  294. 'status_source': '振动加速度阶次跟踪 1P 幅值 fleet MAD-z(切向/补偿通道); raw通道重力门拒',
  295. 'curtail_strip': {'method': 'CMS振动1P与限功率正交', 'n_removed': 'N/A (正交)'}},
  296. 'falsifiability': fals,
  297. 'detail': '重力感知门: raw挥舞/摆振向1P=重力主导(deg2去不掉)/标定不一(跨台bimodal)→INSUFFICIENT路由modal-z; 仅补偿通道跑fleet-z(SOP §3.2b, REJECT改写+真数据RV-2)',
  298. 'affected_machines': affected, 'problem_nature': '可靠性/部件', 'recommendation': reco,
  299. 'discriminator': res,
  300. }
  301. def _finding_pitch(res, cov, *, title='变桨卡桨 (秒级层)'):
  302. """§4.1/3.6 变桨卡桨 → finding。集距死→INSUFFICIENT(不可分); flagged(运动段三叶同步+lag一致)→候选。"""
  303. flagged = res.get('flagged') or []
  304. affected, reco = None, None
  305. if res.get('collective_pitch'):
  306. verdict, ev = 'INSUFFICIENT', res.get('caveat', '集距死: 三桨同值, per叶不可分')
  307. fals = '若秒级独立三叶桨距通道可分(非集距)则可跑卡桨筛查; 集距场叶片问题转振动1P(blade_imbalance)'
  308. elif flagged:
  309. verdict, affected = '候选', list(flagged)
  310. reco = '变桨系统检查 + 取脂送检(磨损金属Fe/Cu/水分)'
  311. top = next((r for r in res.get('per_turbine', []) if r.get('tid') == flagged[0]), {})
  312. ev = f"运动段三叶同步极差离群+lag一致性确认 卡滞候选 {len(flagged)} 台: {flagged}; {top.get('note', '')}"
  313. fals = '若候选台经现场核三叶同步/速率正常则为噪声; SCADA只筛卡死(stick), 渐进磨损(wear)盲须油脂化验(F29/F35双null)'
  314. else:
  315. verdict, ev = '参考', f"运动段三叶同步无卡滞候选; {res.get('wear_blind_caveat', '')}"
  316. fals = '若某叶运动段持续滞后(sync离群+lag一致)则改候选; 渐进磨损 SCADA 盲(须化验)'
  317. return {
  318. 'title': title, 'verdict': verdict, 'claim_class': '相对排名',
  319. 'coverage': {**cov, 'per_day_evidence': ev},
  320. 'basis': {'mast_usability': 'NONE', 'window_months': cov.get('window_months'),
  321. 'status_source': '运动段三叶桨距同步极差p95 + lag一致性 fleet MAD-z(秒级独立三叶)',
  322. 'curtail_strip': {'method': '卡桨与限功率正交(运动段同步性)', 'n_removed': 'N/A (正交)'}},
  323. 'falsifiability': fals,
  324. 'detail': 'SCADA只筛卡死(stick)非渐进磨损(wear≠stick, F29/F35双null须化验); 集距死场不可分INSUFFICIENT(SOP §3.6, 真数据RV-2复现zyx集距)',
  325. 'affected_machines': affected, 'problem_nature': '可靠性/部件', 'recommendation': reco,
  326. 'discriminator': res,
  327. }
  328. # kind → (判别器函数 wrapper, 适配器). 判别器 wrapper 接 (df10, args) 出 res dict。
  329. def _run_curtail(df10, a, n_total):
  330. return D.curtail_vs_fault_sync(
  331. df10, time_col='bin', turbine_col=a['turbine_col'],
  332. curtail_col=a['curtail_col'], freq_col=a.get('freq_col'), gen_col=a.get('gen_col'),
  333. n_total=n_total, sync_thresh=a.get('sync_thresh', 0.5))
  334. def _run_yaw(df10, a, n_total):
  335. return D.yaw_column_usability(
  336. df10, turbine_col=a['turbine_col'], yaw_err_col=a['yaw_err_col'],
  337. spread_thresh=a.get('spread_thresh', 30.0), max_abs_thresh=a.get('max_abs_thresh', 45.0),
  338. min_per_turbine=a.get('min_per_turbine', 200))
  339. def _run_nbm(df10, a, n_total):
  340. gen_mask = None
  341. gfc = a.get('gen_frac_col')
  342. if gfc and gfc in df10:
  343. gen_mask = df10[gfc] > a.get('gen_frac_thresh', 0.8)
  344. return D.nbm_residual(
  345. df10, comp_col=a['comp_col'], p_col=a['p_col'], rotor_col=a['rotor_col'],
  346. turbine_col=a['turbine_col'], time_col='bin',
  347. amb_col=a.get('amb_col'), gen_mask=gen_mask, z_thresh=a.get('z_thresh', 2.0),
  348. min_per_turbine=a.get('min_per_turbine', 1000),
  349. amb_fallback_cols=a.get('amb_fallback_cols'), # §4.6c 死参考备选: 舱内/塔底环境温度
  350. amb_valid_range=tuple(a.get('amb_valid_range', (-40.0, 55.0))),
  351. # 冷浸闭环: developing 命中自动跑冷浸审(真发热vs传感器漂移物理判别, 平陆20#教训)。
  352. # 配置缺省=None → 与旧行为等价(developing 标 not_checked 须人工); 配了 ref/rotor 则自动审。
  353. cold_soak_ref_col=a.get('cold_soak_ref_col'),
  354. cold_soak_rotor_col=a.get('cold_soak_rotor_col'))
  355. def _run_matched_load(df10, a, n_total):
  356. gen_mask = None
  357. gfc = a.get('gen_frac_col')
  358. if gfc and gfc in df10:
  359. gen_mask = df10[gfc] > a.get('gen_frac_thresh', 0.8)
  360. comps = [c for c in a['comp_cols'] if c in df10.columns]
  361. return D.matched_load_overheat(
  362. df10, comp_cols=comps, p_col=a['p_col'], turbine_col=a['turbine_col'], time_col='bin',
  363. rated=a['rated'], gen_mask=gen_mask, dt_z_thresh=a.get('dt_z_thresh', 3.0),
  364. p_lo_frac=a.get('p_lo_frac', 0.30), p_hi_frac=a.get('p_hi_frac', 0.95),
  365. bin_kw=a.get('bin_kw', 100.0), min_cell=a.get('min_cell', 30), min_bins=a.get('min_bins', 3))
  366. def _run_cooling(df10, a, n_total):
  367. if a['temp_col'] not in df10.columns:
  368. return {'causal': None, 'cooling_type': None, 'per_turbine': [], 'flagged': [],
  369. 'n_total': 0, 'verdict': 'INSUFFICIENT', 'evidence': f"温度列 {a['temp_col']} 不存在"}
  370. return D.cooling_efficiency(
  371. df10, temp_col=a['temp_col'], p_col=a['p_col'], turbine_col=a['turbine_col'], time_col='bin',
  372. ws_col=a.get('ws_col'), rated=a.get('rated'), z_thresh=a.get('z_thresh', 2.0),
  373. min_per_turbine=a.get('min_per_turbine', 1000), min_cell=a.get('min_cell', 200))
  374. def _run_vibration(df10, a, n_total):
  375. vibs = [c for c in a['vib_cols'] if c in df10.columns]
  376. if not vibs:
  377. return {'per_channel': {}, 'flagged': [], 'verdict': 'INSUFFICIENT(无有效振动列)', 'ceiling': '预警', 'caveat': '无振动列'}
  378. return D.vibration_trend(
  379. df10, vib_cols=vibs, turbine_col=a['turbine_col'], time_col='bin', z_thresh=a.get('z_thresh', 2.0))
  380. _DISPATCH = {
  381. 'curtail_sync': (_run_curtail, _finding_curtail),
  382. 'yaw_usability': (_run_yaw, _finding_yaw),
  383. 'nbm': (_run_nbm, _finding_nbm),
  384. 'matched_load': (_run_matched_load, _finding_matched_load), # 3.1 NBM盲区兜底 (10min 温度)
  385. 'cooling': (_run_cooling, _finding_cooling), # 3.5 冷却效能 (10min 温度/功率/风速)
  386. 'vibration': (_run_vibration, _finding_vibration), # 3.2 振动level趋势 (10min 振动RMS列, 有则)
  387. # 'ntf' 走单独路径 (需 nacelle/freestream 对齐 array, 非 df10 列) — 见 run_pipeline ntf 块
  388. }
  389. # 秒级/CMS 层适配器: blade_imbalance_1p(CMS波形) / pitch_stick(秒级独立三叶) 须 10min 之外的数据层,
  390. # 不进 _DISPATCH(run_pipeline 只携 df10); 由专用秒级 runner 调判别器后用这些适配器出标准 finding(flagged-based)。
  391. _SEC_ADAPTERS = {'blade': _finding_blade, 'pitch': _finding_pitch}
  392. def _window_days(df10, time_col='bin'):
  393. if time_col not in df10:
  394. return None
  395. t = pd.to_datetime(df10[time_col], errors='coerce')
  396. if t.notna().sum() == 0:
  397. return None
  398. # floor 到 ~1 个 10min bin (0.007 天): 单/少 bin 窗有数据必 truthy, 不被 validate_finding
  399. # 误判"缺/空"(审查 HIGH-2); 0.007 诚实=至少一个 10min 窗, 多日时取真实跨度
  400. span = (t.max() - t.min()).total_seconds() / 86400.0
  401. return max(round(span, 3), 0.007)
  402. def run_pipeline(raw_df, config, *, write=True, out_root=None):
  403. """ON-2 编排: raw long df → clean(CL-1) → 10min(§3.4) → discriminators(§4) → findings.json。
  404. raw_df: 长表 (所有机台堆叠), 含 config['time_col'] / config['machine_col'] / 原始信号列。
  405. 按场加载 (编码/路径/列名是 per-farm 噪声, 不进库)。
  406. config 必填:
  407. farm outputs/<farm>/sop/ 落盘
  408. contract 契约 yaml 路径 (或 dict); section 契约 section 名
  409. time_col 原始时间列 (raw, 派生后变 'bin')
  410. machine_col 机号列
  411. derive: {signal_cols, power_col?, vib_cols?, status_col?, code_sets?, min_valid?}
  412. discriminators: [{kind, args, title?}, ...] kind ∈ {curtail_sync, yaw_usability, nbm}
  413. config 可选:
  414. ntf: {nacelle_ws, freestream_ws, turbine?, sector?} — array 对齐好的机舱/自由来流(见 dogfood 脚本)
  415. window_months coverage 用 (单季→资源类降级判据)
  416. 返回 dict: {clean_report, df10(DataFrame), discriminators(list), findings(list), validation(dict)}。
  417. write=True 落 outputs/<farm>/sop/findings.json (+ clean_gate/cleaned 由 contract_gate 落)。
  418. """
  419. # ON-0 焊点1 (§0.8; 2026-07-29): 出数前必过前置锁定门。
  420. # **默认 warn 不硬拦** —— 实测 4 个锁 3 个不可用 (1 YAML 坏 / 2 空壳), 现在硬接会卡死所有场;
  421. # 锁债清完后置 ON0_STRICT=1 翻硬 (判据见 SOP §0.8 焊点1 条)。新写的 per 场分析脚本
  422. # 请**直接调 assert_frozen(farm)** (默认 strict=True) —— 那才是焊点1 的主用法。
  423. import os as _os
  424. from src.sop.analysis_lock import assert_frozen as _af
  425. _af(config.get("farm", "?"), strict=_os.environ.get("ON0_STRICT") == "1")
  426. farm = config['farm']
  427. section = config.get('section', 'turbine')
  428. time_col = config['time_col']
  429. machine_col = config['machine_col']
  430. target = Path(out_root) if out_root else P.sop(farm)
  431. # ---- stage 1: CL-1 清洗门 (sentinel/range/死列/ffill + 清洗快照指纹)
  432. clean_df, clean_report = apply_contract_gate(
  433. raw_df, config['contract'], section, farm=farm if not out_root else None,
  434. out_dir=str(target) if out_root else None,
  435. machine_col=machine_col, time_col=time_col, persist=config.get('persist', True))
  436. # ---- stage 2: §3.4 10min 真时均派生
  437. dv = config['derive']
  438. df10 = derive_10min(
  439. clean_df, time_col=time_col, signal_cols=dv['signal_cols'], machine_col=machine_col,
  440. power_col=dv.get('power_col'), vib_cols=dv.get('vib_cols', ()),
  441. status_col=dv.get('status_col'), code_sets=dv.get('code_sets'),
  442. min_valid=dv.get('min_valid', 30),
  443. freq=dv.get('freq', '10min')) # ON-2 默认10min; 15min采样场(如肥城上汽900S)传 freq='15min'
  444. n_total = int(df10[machine_col].nunique())
  445. wdays = _window_days(df10)
  446. cov_base = {'n_machines': n_total, 'window_days': wdays,
  447. 'window_months': config.get('window_months')}
  448. # ---- stage 3: §4 判别器 → finding
  449. findings, discr_dump = [], []
  450. for spec in config.get('discriminators', []):
  451. kind = spec['kind']
  452. run_fn, adapt_fn = _DISPATCH[kind]
  453. res = run_fn(df10, spec.get('args', {}), n_total)
  454. discr_dump.append({'kind': kind, 'result': res})
  455. kw = {'title': spec['title']} if spec.get('title') else {}
  456. findings.append(adapt_fn(res, cov_base, **kw))
  457. # ntf 单独 (array 入, 非 df10 列)
  458. if config.get('ntf'):
  459. nt = config['ntf']
  460. res = D.ntf_correction(nt['nacelle_ws'], nt['freestream_ws'],
  461. turbine=nt.get('turbine'), sector=nt.get('sector'))
  462. discr_dump.append({'kind': 'ntf', 'result': res})
  463. cov_ntf = {**cov_base, 'mast_usability': nt.get('mast_usability', 'NONE')}
  464. findings.append(_finding_ntf(res, cov_ntf,
  465. **({'title': nt['title']} if nt.get('title') else {})))
  466. # ---- stage 4: 校验 + 落盘 (issue 不静默)
  467. val_issues = {}
  468. for i, it in enumerate(findings):
  469. iss = schemas.validate_finding(it, tag=f"finding[{i}]({it.get('title','')[:16]})")
  470. if iss:
  471. val_issues[i] = iss
  472. cons = schemas.validate_conservation(findings, n_total, window_days=wdays)
  473. perm = schemas.validate_per_machine(findings) # ① 逐台字段契约核 (affected_machines/problem_nature/recommendation)
  474. validation = {'per_finding': val_issues, 'conservation': cons, 'per_machine': perm,
  475. 'all_pass': not val_issues and not cons and not perm}
  476. doc = {
  477. '_pipeline': 'src/sop/farm_pipeline.run_pipeline (ON-2)',
  478. '_generated': datetime.now().isoformat(timespec='seconds'),
  479. '_farm': farm, '_n_machines': n_total, '_window_days': wdays,
  480. '_validation': validation,
  481. 'findings': findings,
  482. }
  483. # ON-0 焊点2 (§0.8): 产物落盘**前**再核锁哈希 —— 锁若中途被改, 这批数已不对应声明配置, 不该落盘。
  484. # 放 write 之前是刻意的: 放之后再发现, 错的产物已经在盘上了。
  485. from src.sop.analysis_lock import assert_unchanged as _au
  486. _au(config.get("farm", "?"), stage="findings 落盘前")
  487. if write:
  488. target.mkdir(parents=True, exist_ok=True)
  489. (target / 'findings.json').write_text(
  490. json.dumps(doc, ensure_ascii=False, indent=2, default=str), encoding='utf-8')
  491. # ---- stage 4: 柱1 生产器自动跑 (findings → cases 跨批次持续档, 2026-07-11 接线) ----
  492. # 每场 onboard 一条命令到 cases; 之前每场手动 `python -m src.sop.cases <farm>` = 会忘。
  493. # config['produce_cases']=False 可关 (如只想重算 findings 不动案卷)。
  494. cases_out = None
  495. if write and config.get('produce_cases', True):
  496. cases_out = produce_cases_stage(doc, target)
  497. return {'clean_report': clean_report, 'df10': df10,
  498. 'discriminators': discr_dump, 'findings': findings, 'validation': validation,
  499. 'cases': cases_out}
  500. def produce_cases_stage(findings_doc, out_dir):
  501. """柱1 stage 4 (可独测): findings doc → produce_cases(跨批次持续档) → validate → cases.json。
  502. 校验不过 → **响亮 raise, cases.json 不更新** (守护失败必响亮; findings 已落盘, 案卷保持上批)。
  503. 返回 {n_cases, n_prev, states}。
  504. """
  505. from collections import Counter
  506. from src.sop.cases import produce_cases
  507. from src.sop.schemas import validate_cases
  508. out_dir = Path(out_dir)
  509. cpath = out_dir / 'cases.json'
  510. prev = json.loads(cpath.read_text(encoding='utf-8')) if cpath.exists() else []
  511. cases = produce_cases(findings_doc, prev, findings_doc['_generated'],
  512. str(findings_doc['_generated'])[:10])
  513. issues = validate_cases(cases)
  514. if issues:
  515. raise RuntimeError(f'produce_cases validate FAIL ({len(issues)} issue), cases.json 未更新: '
  516. f'{issues[:5]}')
  517. cpath.write_text(json.dumps(cases, ensure_ascii=False, indent=2), encoding='utf-8')
  518. return {'n_cases': len(cases), 'n_prev': len(prev),
  519. 'states': dict(Counter(c['state'] for c in cases))}