# -*- coding: utf-8 -*- """farm_pipeline.py — ON-2 编排器: contract_gate → derive_10min → discriminators → findings.json. 把 Batch2 三块积木 (CL-1 清洗门 / §3.4 10min 派生 / §4 稳定判别器) 串成 farm-agnostic 管线, 新场 onboard 后一条命令跑通"清洗→派生→判别→findings", 替代 per-farm 重写分析脚本 (DP-10 类). 设计: - 纯编排, 不读盘 (raw_df 由调用方按场加载; 库不 hardcode F: 路径/编码/列名 — 那是 per-farm 噪声)。 - 判别器→finding 的映射经"适配器" (_finding_*): 套 SOP §0.2 降级表 + §4 判据出 verdict/claim_class。 - 出 findings.json 前逐条过 schemas.validate_finding (六枚举/coverage/falsifiability), issue 不静默 — 落 _validation。 接口: run_pipeline(raw_df, config) -> dict{clean_report, df10, discriminators, findings, validation}。 config 见 docstring of run_pipeline; dogfood 实例见 scripts/dogfood_farm_pipeline_jtxq.py。 """ from __future__ import annotations try: from app_common.app_common_guanlan.api import install_root as _install_root except ImportError: # 理论不可达;包结构异常时回退到按位置上跳 from pathlib import Path as _P def _install_root(_f): return _P(_f).resolve().parents[4] from src import paths as P import json from datetime import datetime from pathlib import Path import pandas as pd from .contract_gate import apply_contract_gate from .derive10min import derive_10min from . import discriminators as D from . import schemas ROOT = _install_root(__file__) # ============================================================ 判别器 → finding 适配器 # 每个适配器: 入 (判别器结果 dict, coverage dict) → 出 findings_item (过 validate_finding)。 # verdict 取舍套 §0.2 降级表; 不在此造数字, 只把判别器结构化结论翻译成 finding 语言。 def _finding_curtail(res, cov, *, title='限电 vs 故障 同步性判别'): """§4.3 限电同步性 → finding。三证(同步台数/频率偏移/季节集中)齐 = 定论电网限电。""" v = res.get('verdict') if v == '无限电候选': verdict = '参考' fals = '若后续窗出现 ≥半数台同时 curtail 态 则存在限电候选, 需重判' elif v and str(v).startswith('疑误派生'): verdict = 'INSUFFICIENT' fals = ('curtail 信号常年on(>50%全时)→疑派生阈值高于铭牌功率设定(正常次额定误标限电); ' '须核派生口径(corr(P,ws)是否随风/真设定值)后重判, 当前不可判限电') elif v and str(v).startswith('限电信号存疑'): verdict = 'INSUFFICIENT' fals = ('curtail 标志在匹配中风速段未压低功率(setpoint 标签 artifact, pset 多个正常值都满发); ' '须核派生口径或换真限电信号(真停机限电应在同风段把功率压数百kW)后重判, 当前不可判限电') elif v == '电网同步限电': # 三证: sync_fraction(主证) + freq_offset(脱网) + season_peak(集中). 主证强 + 任一旁证 → 定论 strong = (res.get('sync_fraction') or 0) >= 0.5 corrob = (res.get('freq_offset_hz') not in (None, 0)) or ((res.get('season_peak_ratio') or 0) > 2) verdict = '定论' if (strong and corrob) else '准定论·预警' fals = '若 curtail 窗实为 per 机散布(中位同时台数≈1)∨电网频率无偏移 则非电网限电(是 per 机故障)' elif v == '电网限电(部分调度)': # 2026-07-26 新增档: sync 不过半 (电网只调度部分机组) 但旁证 >=2 票。 # 不给"定论" —— 主证(同步过半)未成立, 靠旁证推断, 上限=准定论·预警。 verdict = '准定论·预警' fals = ('若中位同时台数降到 ≈2(真散布故障量级) ∨ 季节集中消失(峰/均<2) ∨ 电网频率无偏移 ' '则改判 per 机故障; 若后续出现 ≥半数台同步则升为定论级电网限电') elif v and str(v).startswith('INSUFFICIENT'): # 2026-07-26 新增档: 落在"部分调度 vs 散布故障"模糊带 (median_simul≈3 且旁证不足)。 # ★旧版会把这类强行判成 per机故障(散布) -> 去追不存在的故障。宁可显式说不知道。 verdict = 'INSUFFICIENT' fals = ('当前证据不足以分"部分电网调度"与"per机散布故障": 中位同时台数落在模糊带(≈3), ' '旁证(季节集中/频率偏移)不足。须补**独立限电台账**或**调度指令记录**才可判; ' '不得据此报机组缺陷, 亦不得据此报电网限电') else: # per机故障(散布) verdict = '候选' fals = '若 curtail 候选窗全场同步(≥半数台同时)∨频率偏移显著 则改判电网限电' return { 'title': title, 'verdict': verdict, 'claim_class': '限电', 'coverage': {**cov, 'per_day_evidence': res.get('evidence', '')}, 'basis': {'mast_usability': 'NONE', 'window_months': cov.get('window_months'), 'status_source': '状态码 curtail 候选态占比(per 10min 窗) + 电网频率 + 月度集中'}, 'falsifiability': fals, 'detail': f"判别器 curtail_vs_fault_sync: verdict={v}; " f"同步窗 {res.get('n_curtail_windows')} / 总机 {res.get('n_total')}", 'affected_machines': None, 'problem_nature': '外部(电网)', # ① 源头落盘 (限电=场级外部, 无问题机; render 渲"全场(电网级)") 'discriminator': res, } def _finding_yaw(res, cov, *, title='静态偏航对风(对风角度列)'): """§4.4① 对风角度列可用性 → finding。不可用(任意零点 artifact)→ INSUFFICIENT。""" usable = res.get('usable') affected, reco = None, None # ① 源头落盘 (偏置台/校正建议) if usable is None: # 有效台<2: 数字字段为 None, evidence 只用 reason 不内插 None verdict, fals = 'INSUFFICIENT', f"有效台<2 不足以判 inter-turbine 散布: {res.get('reason')}" ev = res.get('reason', '有效台<2 (per台样本不足)') elif usable is False: verdict = 'INSUFFICIENT' fals = f"列不可用(任意零点 artifact): {res.get('reason')}; 真对风需 native 对风角列或现场独立量" ev = res.get('reason', '列不可用 (任意零点 artifact)') else: # 列可用. ⚠ "列可用"≠"全场对中健康": 单台真失配(|中位|>5°)不可被"定论健康"掩盖(Critic逮昌平坳#36+8.8°) ptm = res.get('per_turbine_median') or {} offset = sorted(((k, v) for k, v in ptm.items() if abs(v) > 5.0), key=lambda x: -abs(x[1])) if offset: # 有真失配台 → 候选偏置台 (可回收电量), 不笼统判健康 verdict = '候选' affected = [k for k, _ in offset] # ① 偏置台(>5°)直接落盘 reco = '现场独立量(LiDAR/偏航阶跃)核实并下发零点校正' top = ', '.join(f'#{k} {v:+.1f}°' for k, v in offset[:3]) ev = (f"对风角列可用(极差 {res.get('inter_turbine_spread')}°), 多数台对中聚集; " f"但 {len(offset)} 台中位偏置>5°: {top} = 候选偏航偏置台(单台静态偏置~cos²–³ 可回收电量, 须现场独立量校正)") fals = ('若偏置台经现场独立量(LiDAR/偏航阶跃)证实<3° 则为列标定噪声非真偏置; ' '若 inter-turbine 极差>30°∨|中位|>45° 则列不可用(任意零点)') else: # 全台 ±5° 内 → 相对对中健康 (定论级 relative; 绝对零点仍需 LiDAR) verdict = '定论' fals = '若某台对风角 |中位|>45° ∨ inter-turbine 极差>30° (任意零点阈, SOP §4.4①) 则偏置台/列不可用, 需现场独立量核' ev = (f"对风角列可用; 全台 per机中位聚集 ±5° 内 (极差 {res.get('inter_turbine_spread')}° / " f"最大|中位| {res.get('max_abs_median')}°, n={res.get('n_turbines')}台); {res.get('reason')}") return { 'title': title, 'verdict': verdict, 'claim_class': '相对排名', 'coverage': {**cov, 'per_day_evidence': ev}, # ⚠ 2026-07-16 去硬编码假声明 (独立critic af45d34f: 原两字段是模板常量非per场派生, jishan_p1自相矛盾坐实): # mast_usability 判别器根本不查mast(_run_yaw不接mast参数)→ 勿据此判无塔(5/11场实有测风塔); # status_source 未验证native vs wind_dir−nacelle重构 → 真列源须读per场清洗层映射, 不在此断言。 'basis': {'mast_usability': 'NOT_CHECKED (判别器不接mast参数; 勿判无塔 — 部分场实有测风塔)', 'window_months': cov.get('window_months'), 'status_source': '对风角度列 per机中位 (列源=清洗层映射, 未在此验证native/重构 — 见per场清洗config)', 'curtail_strip': {'method': '对风角与限功率正交(限电降功率不改对中, §4.3)', 'n_removed': 'N/A (正交, 无需剔)'}}, 'falsifiability': fals, 'detail': '相对对中健康=per机中位聚集; ⚠绝对偏航零点(可下发校正)仍须 LiDAR/测风塔/偏航阶跃(SOP §4.4 铁律)', 'affected_machines': affected, 'problem_nature': '性能', 'recommendation': reco, # ① 源头落盘 'discriminator': res, } def _finding_nbm(res, cov, *, title='温度 NBM 残差第二道'): """§4.6 温度 NBM 残差 → finding。flagged 非空=候选(需绝对温度+持续性二判); 空=参考(筛过非保证健康)。""" flagged = res.get('flagged') or [] affected, reco = None, None # ① 源头落盘 (候选偏热台) if not res.get('per_turbine'): verdict, ev = 'INSUFFICIENT', res.get('note', '样本不足, NBM 无法拟合') fals = '若补足样本(per台≥阈)后拟合稳定则可判' elif flagged: verdict = '候选' affected = list(flagged) # ① 候选偏热台直接落盘 reco = '现场核30秒实际温度(非仅告警计数) + 盯持续 NBM z>2 偏热台' top = res['per_turbine'][0] ev = (f"NBM 残差 fleet MAD-z 升候选 {len(flagged)} 台: {flagged}; " f"头名 {top['tid']} z={top['z']} (残差均 {top['resid_mean']}°, EWMA超 {top.get('ewma_exc')})") fals = '若候选台绝对温度未超同型分位 ∧ EWMA 持续超占比低 则为基线噪声非真发热(防 C 系低基线污染)' else: verdict = '参考' ev = f"NBM 残差第二道无 z≥{2.0} 候选台 (n={res.get('n_total')}台); 残差离群筛过" fals = '若个别台绝对温度持续超同型高分位则仍需单独绝对判(残差筛过≠绝对健康)' # §4.6c 参考通道守卫: 死 amb 台已剔出/降级 INSUFFICIENT (防"参考冻结"伪装"部件热", F0810001022 戒) if res.get('ref_dead'): ev += (' | ⚠参考(舱外温度)冻结剔出 ' + str(len(res['ref_dead'])) + ' 台(INSUFFICIENT, 非候选): ' + '; '.join(f"{x['tid']}({x['reason']})" for x in res['ref_dead'])) if res.get('ref_fallback'): ev += ' | 参考死→备选顶替: ' + '; '.join(f"{x['tid']}→{x['col']}" for x in res['ref_fallback']) return { 'title': title, 'verdict': verdict, 'claim_class': '相对排名', 'coverage': {**cov, 'per_day_evidence': ev}, 'basis': {'mast_usability': 'NONE', 'window_months': cov.get('window_months'), 'status_source': 'pooled OLS 残差(comp~P+rotor+Tmed) fleet MAD-z + EWMA 持续性', 'curtail_strip': {'method': 'NBM 回归 P 项吸收负荷(含限电降功率)共模, 限电效应回归控制非剔除(§4.3/§4.6)', 'n_removed': 'N/A (回归控制)'}}, 'falsifiability': fals, 'detail': '残差第二道(吸收功率/转速/季节共模)+绝对温度双判; 单冷 outlier 不驱动 corr(SOP §4.6)', 'affected_machines': affected, 'problem_nature': '可靠性/部件', 'recommendation': reco, # ① 源头落盘 'discriminator': res, } def _finding_ntf(res, cov, *, title='机舱风 NTF / 自由来流达成率'): """§4.5 NTF + 空间代表性 → finding。单 mast 非空间代表 → 绝对达成只报区间(INSUFFICIENT 绝对)。""" rep = res.get('spatial_representative') mast = str(cov.get('mast_usability', 'NONE')).upper() # 绝对达成(机舱风 self_ref)须双门: mast PASS ∧ 空间代表. 任一不齐 → INSUFFICIENT 绝对 # (空间代表只证机舱/自由来流比值内部一致, 与测风塔本身物理是否 PASS 无关 — 审查 HIGH-1) if rep is True and mast == 'PASS': verdict = '准定论·预警' fals = '若 NTF 随扇区极差>0.1 则为方向加权混合非机型 NTF, 绝对达成降区间' else: verdict = 'INSUFFICIENT' fals = (f'绝对达成须 mast PASS ∧ 空间代表双门 (当前 mast={mast}/空间代表={rep}) → 只报区间; ' '若两门齐(各机比∈[0.9,1.1]∧扇区极差≤0.1∧mast PASS)则可定量') ev = (f"NTF(机舱/自由来流)中位 {res.get('ntf_median')}; reads_low={res.get('reads_low')}; " f"空间代表={rep}; {res.get('caveat')}") return { 'title': title, 'verdict': verdict, 'claim_class': '绝对性能', 'coverage': {**cov, 'per_day_evidence': ev}, 'basis': {'mast_usability': cov.get('mast_usability', 'NONE'), 'window_months': cov.get('window_months'), 'status_source': '测风塔自由来流 vs 机舱风速计 NTF', 'curtail_strip': {'method': '功率曲线达成率须剔限电(§4.3 限电=性能头号混杂); 限电态由 curtail 判别器标记剔', 'n_removed': '限电样本(601态等, 见限电finding)'}}, 'falsifiability': fals, 'detail': '绝对达成率天然 self_ref(机舱风), 必自由来流修正; 单 mast 空间代表性 caveat(SOP §4.5)', 'affected_machines': None, 'problem_nature': '性能', # ① 源头落盘 (场级数据局限, 无问题机) 'recommendation': '补测风塔/独立风速量 → 绝对达成率升定论', 'discriminator': res, } # ──────────── ④部件健康 适配器 (2026-06-20 loop训练抽库 + 真数据RV-2 + 独立Critic HOLD)。 # 全部 **affected_machines 从 flagged 取, 禁 verdict-first** (Critic 逮 cooling verdict 会被整体 load_driven 淹没真衰减台)。 def _finding_matched_load(res, cov, *, title='温度 matched-load 双判 (NBM盲区兜底)'): """§4.6/3.1 matched-load 双判 → finding。flagged(同功率bin ΔT-z ∨ 多部件簇状)=候选; 空=参考(筛过非健康)。""" flagged = res.get('flagged') or [] affected, reco = None, None if not res.get('n_components'): verdict, ev = 'INSUFFICIENT', res.get('note', '无有效部件温列/功率带样本不足') fals = '若补足部件温列 + 功率带重叠样本则可判' elif flagged: verdict, affected = '候选', list(flagged) reco = '现场IR/CMS核 + 与 nbm_residual(极值)交叉对照 (matched-load 兜底中度/多部件/簇状升温)' top = res['per_turbine'].get(flagged[0], {}) ev = f"matched-load 双判升候选 {len(flagged)} 台: {flagged}; 头名 {flagged[0]} {top.get('reason')}" fals = '若候选台同功率bin ΔT 经现场核为传感器标定差/无绝对失控 则非真发热; ⚠非重叠功率带台 ΔT路径致盲(须fleet功率带重叠)' else: verdict = '参考' ev = f"matched-load 双判无升候选 (n_components={res.get('n_components')}); NBM盲区(中度/多部件/簇状)兜底筛过" fals = '若个别台同功率bin多部件持续偏热则仍需单独 matched-load 核 (筛过≠绝对健康)' if res.get('dead_channel'): # §4.6c 冻结部件温列已逐台剔出排名+簇状 (F0810001022 戒) ev += ' | ⚠§4.6c 冻结部件温列剔出: ' + '; '.join(f"{x['comp']}:{x['tids']}" for x in res['dead_channel']) return { 'title': title, 'verdict': verdict, 'claim_class': '相对排名', 'coverage': {**cov, 'per_day_evidence': ev}, 'basis': {'mast_usability': 'NONE', 'window_months': cov.get('window_months'), 'status_source': '同功率bin 部件温−fleet中位 ΔT-z + 多部件簇状同热占比', 'curtail_strip': {'method': 'matched-load 同功率bin比已控负荷(含限电降功率), bin内消(§4.6)', 'n_removed': 'N/A (同bin控制)'}}, 'falsifiability': fals, 'detail': 'NBM极值检测器盲区兜底(逮中度+4-6°C/多部件/簇状, 凉水泉2208F戒); ΔT相对fleet非OEM绝对阈, 上限候选; 真数据RV-2复现凉水泉', 'affected_machines': affected, 'problem_nature': '可靠性/部件', 'recommendation': reco, 'discriminator': res, } def _finding_cooling(res, cov, *, title='冷却系统效能'): """§4.1/3.5 冷却效能 → finding。⚠ affected 必从 flagged(C层真衰减台)取, 禁 verdict-first(整体可能 load_driven 淹没)。""" flagged = res.get('flagged') or [] affected, reco = None, None if res.get('verdict') == 'INSUFFICIENT' or not res.get('per_turbine'): verdict, ev = 'INSUFFICIENT', res.get('evidence', '样本不足') fals = '若补足同功率bin样本则可判' elif flagged: verdict, affected = '候选', list(flagged) reco = '散热器/水泵/风扇检修 + 30秒实温 vs F级155°C 核(非仅告警计数)' ev = res.get('evidence', '') fals = '若候选台同功率温偏经现场核为传感器/工况差 则非冷却衰减; 真过热须30秒实温超部件告警阈坐实' else: verdict = '参考' # ⚠Critic: flagged=[] = 同功率温检测下限内无可检出 ≠ 冷却健康定论(小样本可漏中度偏热台, A11型) ev = (res.get('evidence', '') + f" | causal={res.get('causal')}/cooling_type={res.get('cooling_type')}; " "⚠flagged=[]=检测下限内无可检出, 非'冷却健康'定论(小样本可漏中度偏热台)") fals = '若同功率bin某台温水平系统偏高(z≥2∧ΔT≥2°C)则改候选; "高温降容"若同风速高温组功率反更高则负荷驱动非衰减' if res.get('dead_channel'): # §4.6c 冻结冷却水温台已剔出排名 (F0810001022 戒) ev += ' | ⚠§4.6c 冻结冷却水温台剔出排名(INSUFFICIENT): ' + str(res['dead_channel']) return { 'title': title, 'verdict': verdict, 'claim_class': '相对排名', 'coverage': {**cov, 'per_day_evidence': ev}, 'basis': {'mast_usability': 'NONE', 'window_months': cov.get('window_months'), 'status_source': '同功率bin冷却温−fleet中位 MAD-z(C层) + 因果方向守卫(高温组功率) + 释放曲线分型', 'curtail_strip': {'method': '同功率bin比已控负荷; 因果守卫剔负荷驱动温升(非冷却衰减)', 'n_removed': 'N/A (同bin控制+因果守卫)'}}, 'falsifiability': fals, 'detail': '因果反置守卫(功率↑→温升=负荷驱动非降容, 郭家店1338→197MWh戒) + 低温冷启vs高温降容释放曲线; affected从C层flagged取(禁verdict-first, Critic); 真数据RV-2复现郭家店', 'affected_machines': affected, 'problem_nature': '可靠性/部件', 'recommendation': reco, 'discriminator': res, } def _finding_vibration(res, cov, *, title='振动 level 趋势 (粗筛)'): """§4.6/3.2 振动 level 趋势 → finding。flagged=[(测点,台)] → 候选(affected=去重台号); 退化列/无效→INSUFFICIENT。""" flagged = res.get('flagged') or [] affected, reco = None, None if str(res.get('verdict', '')).startswith('INSUFFICIENT'): verdict, ev = 'INSUFFICIENT', res.get('caveat', '退化列/样本不足/<2台') fals = '若补足有效振动列(非退化)+≥2台则可判' elif flagged: verdict = '候选' affected = sorted({t for _c, t in flagged}) reco = 'CMS高频频谱诊断定位(轴承BPFO/BPFI/齿轮啮合/叶片1P); 无CMS→建议加装振动测点' ev = '振动 level 离群/爬升候选 (分测点 fleet MAD-z / 月度斜率持续): ' + ', '.join(f'{c}:{t}' for c, t in flagged[:4]) fals = '若候选台同测点同工况健康台同样高则共模(机型/工况)非本台; 仅10min level无频谱→出不了部件根因(只候选)' else: verdict, ev = '参考', '分测点振动 level 各测点平稳 (无 z≥2/爬升); 10min level 粗筛过' fals = '若某测点 level z≥2∧月度爬升持续则改候选; level粗筛过≠CMS频谱健康' return { 'title': title, 'verdict': verdict, 'claim_class': '相对排名', 'coverage': {**cov, 'per_day_evidence': ev}, 'basis': {'mast_usability': 'NONE', 'window_months': cov.get('window_months'), 'status_source': '分测点 fleet null MAD-z(禁混池) + 月度Theil-Sen斜率+持续性', 'curtail_strip': {'method': '振动level与限功率正交(限电不直接改振动level)', 'n_removed': 'N/A (正交)'}}, 'falsifiability': fals, 'detail': '10min level 粗筛(verdict上限预警); 部件级须CMS频谱; 分测点禁混池+退化列守卫+MAD地板(SOP §4.6 R5)', 'affected_machines': affected, 'problem_nature': '可靠性/部件', 'recommendation': reco, 'discriminator': res, } def _finding_blade(res, cov, *, title='叶片不平衡 1P (CMS/秒级层)'): """§4.1/3.2b 叶片1P 重力感知门 → finding。raw 通道→INSUFFICIENT(重力/未标定); 补偿通道 flagged→候选。""" flagged = res.get('flagged') or [] affected, reco = None, None if 'INSUFFICIENT' in str(res.get('verdict', '')): verdict, ev = 'INSUFFICIENT', res.get('caveat', 'raw通道1P不可信(重力/未标定)') fals = '若用切向/已补偿通道(gravity_compensated=True)或矢量扣平衡参考窗重力相量则可跑fleet-z; raw挥舞/摆振向1P=重力/标定不一不可作不平衡(路由modal-z)' elif flagged: verdict, affected = '候选', list(flagged) reco = '现场叶片检查(质量配平/气动外形) + CMS 持续监视' ev = f"已补偿通道 1P fleet MAD-z 升候选 {len(flagged)} 台: {flagged}" fals = '若候选台 1P 经阶次跟踪/现场核落回fleet则原是泄漏/标定artifact; 极低1P台疑死/低增益通道(非健康)' else: verdict, ev = '参考', '已补偿通道 1P fleet 中位无离群' fals = '若某台 1P z≥2 则改候选' return { 'title': title, 'verdict': verdict, 'claim_class': '相对排名', 'coverage': {**cov, 'per_day_evidence': ev}, 'basis': {'mast_usability': 'NONE', 'window_months': cov.get('window_months'), 'status_source': '振动加速度阶次跟踪 1P 幅值 fleet MAD-z(切向/补偿通道); raw通道重力门拒', 'curtail_strip': {'method': 'CMS振动1P与限功率正交', 'n_removed': 'N/A (正交)'}}, 'falsifiability': fals, 'detail': '重力感知门: raw挥舞/摆振向1P=重力主导(deg2去不掉)/标定不一(跨台bimodal)→INSUFFICIENT路由modal-z; 仅补偿通道跑fleet-z(SOP §3.2b, REJECT改写+真数据RV-2)', 'affected_machines': affected, 'problem_nature': '可靠性/部件', 'recommendation': reco, 'discriminator': res, } def _finding_pitch(res, cov, *, title='变桨卡桨 (秒级层)'): """§4.1/3.6 变桨卡桨 → finding。集距死→INSUFFICIENT(不可分); flagged(运动段三叶同步+lag一致)→候选。""" flagged = res.get('flagged') or [] affected, reco = None, None if res.get('collective_pitch'): verdict, ev = 'INSUFFICIENT', res.get('caveat', '集距死: 三桨同值, per叶不可分') fals = '若秒级独立三叶桨距通道可分(非集距)则可跑卡桨筛查; 集距场叶片问题转振动1P(blade_imbalance)' elif flagged: verdict, affected = '候选', list(flagged) reco = '变桨系统检查 + 取脂送检(磨损金属Fe/Cu/水分)' top = next((r for r in res.get('per_turbine', []) if r.get('tid') == flagged[0]), {}) ev = f"运动段三叶同步极差离群+lag一致性确认 卡滞候选 {len(flagged)} 台: {flagged}; {top.get('note', '')}" fals = '若候选台经现场核三叶同步/速率正常则为噪声; SCADA只筛卡死(stick), 渐进磨损(wear)盲须油脂化验(F29/F35双null)' else: verdict, ev = '参考', f"运动段三叶同步无卡滞候选; {res.get('wear_blind_caveat', '')}" fals = '若某叶运动段持续滞后(sync离群+lag一致)则改候选; 渐进磨损 SCADA 盲(须化验)' return { 'title': title, 'verdict': verdict, 'claim_class': '相对排名', 'coverage': {**cov, 'per_day_evidence': ev}, 'basis': {'mast_usability': 'NONE', 'window_months': cov.get('window_months'), 'status_source': '运动段三叶桨距同步极差p95 + lag一致性 fleet MAD-z(秒级独立三叶)', 'curtail_strip': {'method': '卡桨与限功率正交(运动段同步性)', 'n_removed': 'N/A (正交)'}}, 'falsifiability': fals, 'detail': 'SCADA只筛卡死(stick)非渐进磨损(wear≠stick, F29/F35双null须化验); 集距死场不可分INSUFFICIENT(SOP §3.6, 真数据RV-2复现zyx集距)', 'affected_machines': affected, 'problem_nature': '可靠性/部件', 'recommendation': reco, 'discriminator': res, } # kind → (判别器函数 wrapper, 适配器). 判别器 wrapper 接 (df10, args) 出 res dict。 def _run_curtail(df10, a, n_total): return D.curtail_vs_fault_sync( df10, time_col='bin', turbine_col=a['turbine_col'], curtail_col=a['curtail_col'], freq_col=a.get('freq_col'), gen_col=a.get('gen_col'), n_total=n_total, sync_thresh=a.get('sync_thresh', 0.5)) def _run_yaw(df10, a, n_total): return D.yaw_column_usability( df10, turbine_col=a['turbine_col'], yaw_err_col=a['yaw_err_col'], spread_thresh=a.get('spread_thresh', 30.0), max_abs_thresh=a.get('max_abs_thresh', 45.0), min_per_turbine=a.get('min_per_turbine', 200)) def _run_nbm(df10, a, n_total): gen_mask = None gfc = a.get('gen_frac_col') if gfc and gfc in df10: gen_mask = df10[gfc] > a.get('gen_frac_thresh', 0.8) return D.nbm_residual( df10, comp_col=a['comp_col'], p_col=a['p_col'], rotor_col=a['rotor_col'], turbine_col=a['turbine_col'], time_col='bin', amb_col=a.get('amb_col'), gen_mask=gen_mask, z_thresh=a.get('z_thresh', 2.0), min_per_turbine=a.get('min_per_turbine', 1000), amb_fallback_cols=a.get('amb_fallback_cols'), # §4.6c 死参考备选: 舱内/塔底环境温度 amb_valid_range=tuple(a.get('amb_valid_range', (-40.0, 55.0))), # 冷浸闭环: developing 命中自动跑冷浸审(真发热vs传感器漂移物理判别, 平陆20#教训)。 # 配置缺省=None → 与旧行为等价(developing 标 not_checked 须人工); 配了 ref/rotor 则自动审。 cold_soak_ref_col=a.get('cold_soak_ref_col'), cold_soak_rotor_col=a.get('cold_soak_rotor_col')) def _run_matched_load(df10, a, n_total): gen_mask = None gfc = a.get('gen_frac_col') if gfc and gfc in df10: gen_mask = df10[gfc] > a.get('gen_frac_thresh', 0.8) comps = [c for c in a['comp_cols'] if c in df10.columns] return D.matched_load_overheat( df10, comp_cols=comps, p_col=a['p_col'], turbine_col=a['turbine_col'], time_col='bin', rated=a['rated'], gen_mask=gen_mask, dt_z_thresh=a.get('dt_z_thresh', 3.0), p_lo_frac=a.get('p_lo_frac', 0.30), p_hi_frac=a.get('p_hi_frac', 0.95), bin_kw=a.get('bin_kw', 100.0), min_cell=a.get('min_cell', 30), min_bins=a.get('min_bins', 3)) def _run_cooling(df10, a, n_total): if a['temp_col'] not in df10.columns: return {'causal': None, 'cooling_type': None, 'per_turbine': [], 'flagged': [], 'n_total': 0, 'verdict': 'INSUFFICIENT', 'evidence': f"温度列 {a['temp_col']} 不存在"} return D.cooling_efficiency( df10, temp_col=a['temp_col'], p_col=a['p_col'], turbine_col=a['turbine_col'], time_col='bin', ws_col=a.get('ws_col'), rated=a.get('rated'), z_thresh=a.get('z_thresh', 2.0), min_per_turbine=a.get('min_per_turbine', 1000), min_cell=a.get('min_cell', 200)) def _run_vibration(df10, a, n_total): vibs = [c for c in a['vib_cols'] if c in df10.columns] if not vibs: return {'per_channel': {}, 'flagged': [], 'verdict': 'INSUFFICIENT(无有效振动列)', 'ceiling': '预警', 'caveat': '无振动列'} return D.vibration_trend( df10, vib_cols=vibs, turbine_col=a['turbine_col'], time_col='bin', z_thresh=a.get('z_thresh', 2.0)) _DISPATCH = { 'curtail_sync': (_run_curtail, _finding_curtail), 'yaw_usability': (_run_yaw, _finding_yaw), 'nbm': (_run_nbm, _finding_nbm), 'matched_load': (_run_matched_load, _finding_matched_load), # 3.1 NBM盲区兜底 (10min 温度) 'cooling': (_run_cooling, _finding_cooling), # 3.5 冷却效能 (10min 温度/功率/风速) 'vibration': (_run_vibration, _finding_vibration), # 3.2 振动level趋势 (10min 振动RMS列, 有则) # 'ntf' 走单独路径 (需 nacelle/freestream 对齐 array, 非 df10 列) — 见 run_pipeline ntf 块 } # 秒级/CMS 层适配器: blade_imbalance_1p(CMS波形) / pitch_stick(秒级独立三叶) 须 10min 之外的数据层, # 不进 _DISPATCH(run_pipeline 只携 df10); 由专用秒级 runner 调判别器后用这些适配器出标准 finding(flagged-based)。 _SEC_ADAPTERS = {'blade': _finding_blade, 'pitch': _finding_pitch} def _window_days(df10, time_col='bin'): if time_col not in df10: return None t = pd.to_datetime(df10[time_col], errors='coerce') if t.notna().sum() == 0: return None # floor 到 ~1 个 10min bin (0.007 天): 单/少 bin 窗有数据必 truthy, 不被 validate_finding # 误判"缺/空"(审查 HIGH-2); 0.007 诚实=至少一个 10min 窗, 多日时取真实跨度 span = (t.max() - t.min()).total_seconds() / 86400.0 return max(round(span, 3), 0.007) def run_pipeline(raw_df, config, *, write=True, out_root=None): """ON-2 编排: raw long df → clean(CL-1) → 10min(§3.4) → discriminators(§4) → findings.json。 raw_df: 长表 (所有机台堆叠), 含 config['time_col'] / config['machine_col'] / 原始信号列。 按场加载 (编码/路径/列名是 per-farm 噪声, 不进库)。 config 必填: farm outputs//sop/ 落盘 contract 契约 yaml 路径 (或 dict); section 契约 section 名 time_col 原始时间列 (raw, 派生后变 'bin') machine_col 机号列 derive: {signal_cols, power_col?, vib_cols?, status_col?, code_sets?, min_valid?} discriminators: [{kind, args, title?}, ...] kind ∈ {curtail_sync, yaw_usability, nbm} config 可选: ntf: {nacelle_ws, freestream_ws, turbine?, sector?} — array 对齐好的机舱/自由来流(见 dogfood 脚本) window_months coverage 用 (单季→资源类降级判据) 返回 dict: {clean_report, df10(DataFrame), discriminators(list), findings(list), validation(dict)}。 write=True 落 outputs//sop/findings.json (+ clean_gate/cleaned 由 contract_gate 落)。 """ # ON-0 焊点1 (§0.8; 2026-07-29): 出数前必过前置锁定门。 # **默认 warn 不硬拦** —— 实测 4 个锁 3 个不可用 (1 YAML 坏 / 2 空壳), 现在硬接会卡死所有场; # 锁债清完后置 ON0_STRICT=1 翻硬 (判据见 SOP §0.8 焊点1 条)。新写的 per 场分析脚本 # 请**直接调 assert_frozen(farm)** (默认 strict=True) —— 那才是焊点1 的主用法。 import os as _os from src.sop.analysis_lock import assert_frozen as _af _af(config.get("farm", "?"), strict=_os.environ.get("ON0_STRICT") == "1") farm = config['farm'] section = config.get('section', 'turbine') time_col = config['time_col'] machine_col = config['machine_col'] target = Path(out_root) if out_root else P.sop(farm) # ---- stage 1: CL-1 清洗门 (sentinel/range/死列/ffill + 清洗快照指纹) clean_df, clean_report = apply_contract_gate( raw_df, config['contract'], section, farm=farm if not out_root else None, out_dir=str(target) if out_root else None, machine_col=machine_col, time_col=time_col, persist=config.get('persist', True)) # ---- stage 2: §3.4 10min 真时均派生 dv = config['derive'] df10 = derive_10min( clean_df, time_col=time_col, signal_cols=dv['signal_cols'], machine_col=machine_col, power_col=dv.get('power_col'), vib_cols=dv.get('vib_cols', ()), status_col=dv.get('status_col'), code_sets=dv.get('code_sets'), min_valid=dv.get('min_valid', 30), freq=dv.get('freq', '10min')) # ON-2 默认10min; 15min采样场(如肥城上汽900S)传 freq='15min' n_total = int(df10[machine_col].nunique()) wdays = _window_days(df10) cov_base = {'n_machines': n_total, 'window_days': wdays, 'window_months': config.get('window_months')} # ---- stage 3: §4 判别器 → finding findings, discr_dump = [], [] for spec in config.get('discriminators', []): kind = spec['kind'] run_fn, adapt_fn = _DISPATCH[kind] res = run_fn(df10, spec.get('args', {}), n_total) discr_dump.append({'kind': kind, 'result': res}) kw = {'title': spec['title']} if spec.get('title') else {} findings.append(adapt_fn(res, cov_base, **kw)) # ntf 单独 (array 入, 非 df10 列) if config.get('ntf'): nt = config['ntf'] res = D.ntf_correction(nt['nacelle_ws'], nt['freestream_ws'], turbine=nt.get('turbine'), sector=nt.get('sector')) discr_dump.append({'kind': 'ntf', 'result': res}) cov_ntf = {**cov_base, 'mast_usability': nt.get('mast_usability', 'NONE')} findings.append(_finding_ntf(res, cov_ntf, **({'title': nt['title']} if nt.get('title') else {}))) # ---- stage 4: 校验 + 落盘 (issue 不静默) val_issues = {} for i, it in enumerate(findings): iss = schemas.validate_finding(it, tag=f"finding[{i}]({it.get('title','')[:16]})") if iss: val_issues[i] = iss cons = schemas.validate_conservation(findings, n_total, window_days=wdays) perm = schemas.validate_per_machine(findings) # ① 逐台字段契约核 (affected_machines/problem_nature/recommendation) validation = {'per_finding': val_issues, 'conservation': cons, 'per_machine': perm, 'all_pass': not val_issues and not cons and not perm} doc = { '_pipeline': 'src/sop/farm_pipeline.run_pipeline (ON-2)', '_generated': datetime.now().isoformat(timespec='seconds'), '_farm': farm, '_n_machines': n_total, '_window_days': wdays, '_validation': validation, 'findings': findings, } # ON-0 焊点2 (§0.8): 产物落盘**前**再核锁哈希 —— 锁若中途被改, 这批数已不对应声明配置, 不该落盘。 # 放 write 之前是刻意的: 放之后再发现, 错的产物已经在盘上了。 from src.sop.analysis_lock import assert_unchanged as _au _au(config.get("farm", "?"), stage="findings 落盘前") if write: target.mkdir(parents=True, exist_ok=True) (target / 'findings.json').write_text( json.dumps(doc, ensure_ascii=False, indent=2, default=str), encoding='utf-8') # ---- stage 4: 柱1 生产器自动跑 (findings → cases 跨批次持续档, 2026-07-11 接线) ---- # 每场 onboard 一条命令到 cases; 之前每场手动 `python -m src.sop.cases ` = 会忘。 # config['produce_cases']=False 可关 (如只想重算 findings 不动案卷)。 cases_out = None if write and config.get('produce_cases', True): cases_out = produce_cases_stage(doc, target) return {'clean_report': clean_report, 'df10': df10, 'discriminators': discr_dump, 'findings': findings, 'validation': validation, 'cases': cases_out} def produce_cases_stage(findings_doc, out_dir): """柱1 stage 4 (可独测): findings doc → produce_cases(跨批次持续档) → validate → cases.json。 校验不过 → **响亮 raise, cases.json 不更新** (守护失败必响亮; findings 已落盘, 案卷保持上批)。 返回 {n_cases, n_prev, states}。 """ from collections import Counter from src.sop.cases import produce_cases from src.sop.schemas import validate_cases out_dir = Path(out_dir) cpath = out_dir / 'cases.json' prev = json.loads(cpath.read_text(encoding='utf-8')) if cpath.exists() else [] cases = produce_cases(findings_doc, prev, findings_doc['_generated'], str(findings_doc['_generated'])[:10]) issues = validate_cases(cases) if issues: raise RuntimeError(f'produce_cases validate FAIL ({len(issues)} issue), cases.json 未更新: ' f'{issues[:5]}') cpath.write_text(json.dumps(cases, ensure_ascii=False, indent=2), encoding='utf-8') return {'n_cases': len(cases), 'n_prev': len(prev), 'states': dict(Counter(c['state'] for c in cases))}