#!/usr/bin/env python3 # -*- coding: utf-8 -*- r"""部件历史生成端 —— 补上"包内没有生成端"的 `m5_cms_tcm/component_history.json` (用户令 2026-09-19)。 ## 为什么要写它 用户令:「**所有的计算 均要形成 观澜的源代码**,确保 观澜 在别的电脑安装后,系统运行正常。」 `component_history.json` 是**随包快照**(族表 kind=shipped、gen=None,见 `scripts/products_reverse_audit.py::HUMAN_SNAPSHOTS`),随包件里的"在升/闭环/换件日期"含人工裁决与 厂家锚,**不是测量数据的函数**。于是另一台机器上装完、放数据、重算之后: `src/ontology/trend_ingest.py` 拿不到在升/闭环证据(图谱里两类 Evidence 全缺)、 `scripts/windscada_serve.py` 的决策链盘「振动在升/闭环」两格永远空。 本脚本把它**从已有产物算出来**(CMS 窗索引 + 工单台账),口径写在下面与产物 meta 里。 ## 产物与消费者 (字段契约) outputs/<场>/m5_cms_tcm/component_history.json meta built / windows / segments / segments_kind / 口径 / ★数据边界 components[] turbine, component, n_windows, n_replace, series, events summary['★当前在升 (末三窗/早三窗 ≥1.3)'] ← src/ontology/trend_ingest.py 逐行读 turbine, component, ratio, early, late, n_events, 工单可对 summary['★换件闭环 (检出→换件→恢复, 算法端到端验证)'] ← 同上 + scripts/windscada_serve.py 链盘 turbine, component, ratio, pre, post, replace_date, title, verdict summary[非★键] 消费端只认 ★ 前缀的键 ⇒ 非★键是给人/给诊断看的, 不被误读 ★ `component` 必须是 `src/ontology/trend_ingest.py::COMP_MAP` 里的中文名 (齿轮箱/发电机/主轴承/变桨/偏航/变流器), 否则该行被消费者跳过。本器只出**齿轮箱/发电机/主轴承**三类 —— CMS 索引里只有这三类的测点 (Main_bearing_* / Gear_* / Generator_*),变桨/偏航/变流器在本产物无振动测点,不臆造。 ## 分段 (segment) 口径 —— 两种模式, 自动选 参考口径是**导出窗**为段 (末三窗/早三窗; 健康期窗中位)。本机实况 (2026-09-19 实测): `src/windcms/config.py::FARMS['rudong']['windows']` 只解析出**一个**窗 `w0316` (`outputs/rudong/m5_cms_tcm/windows/w0316/index.parquet`,2,066,686 行 × 54 列, trigger_time 2026-03-16 17:27 → 2026-04-21 10:18)。窗数 < 3 ⇒ 自动降为 **窗内周段** (ISO 周, 2026-W12 … 2026-W17 共 6 段)。两种模式都在本文件里实现: len(windows) >= 3 segments_kind='window' 段名=窗名, 段日=该窗 trigger_time 最早日 len(windows) < 3 segments_kind='week' 段名=ISO 周 (2026-W12), 段日=该周周一 **多窗口径未在本机取数验证**(本机只有 1 窗 ⇒ 该分支只有单测/纸面正确性,不给"已验证"的说法)。 ## 数据现实与诚实边界 (不造数) * **`★当前在升` 本机必然是空表**:判据要求"末三窗/早三窗"且**多窗口径**;本机 1 窗 ⇒ 不可算。 ★ 不降门槛、不换判据去凑一行。作为补充, 出**非★**键 `窗内周段趋势 (单窗补充口径, 仅参考不作判据)`:同一条比值算式 (末三段/早三段, 要求 6 段以取到 **不相交**的早/晚两组) 在周段上跑出来是什么, 就写什么。 * **`★换件闭环` 本机为空**:闭环要求"该台该部件的更换/换件工单,且段跨度里**既有段在换件日之前、 又有段在其后**"(前段须**结束**在工单日之前、后段须**开始**在工单日之后 —— 工单所在那一段两边都不算, 否则换后数据会混进"换前",实测污染过 3 条)。工单台账 (`outputs/<场>/windscada/workorders.parquet`, 5,876 行) 实测跨 **2020-01-03 → 2026-07-08**,其中传动链件 903 行(去重后)/换件类 12 条 —— 但**全部落在段跨度 (2026-03-16~04-21) 之外** (换件工单最晚 2025-10-14),窗内 0 条 ⇒ 无"换件前段"可比 ⇒ 空表。空表本身就是结论。 * 台账有**逐字重复行**(实测 2026-01-03 环网柜那条同一台同一分钟 7 行;换件类里 WTG23 2024-09-12 发电机轴承更换 4 行)⇒ 本器按 (台,日,标题,部件) 去重后再计数,原始/去重两个数都写进 meta.台账。 * 段内样本数一并写进 series 的 `n` —— 周段只有 2~142 个标量记录 (中位 14), 中位数据此散布, 读者可自行判断。 ## 用法 python scripts/component_history_build.py # 算并写盘 + 自登记 python scripts/component_history_build.py --dry-run # 只算不写, 打印漏斗 python scripts/component_history_build.py --farm rudong # 场名 (默认 rudong) """ from __future__ import annotations import argparse import json import pathlib import re import sys import time ROOT = pathlib.Path(__file__).resolve().parents[1] sys.path.insert(0, str(ROOT)) # ★GBK 控制台/日志重定向不再因 '²' '⇒' 这类字符抛 UnicodeEncodeError (2026-09-19 远端实逮: # `baseline_38_build.py` 打印 'm/s²' 时在 cp936 下 rc=1, 整条恢复链少一件产物)。 # 只把**不可编码字符替换掉**, 不改流编码 —— 中文照旧可读。 try: sys.stdout.reconfigure(errors='replace') except Exception: pass import numpy as np # noqa: E402 import pandas as pd # noqa: E402 import pyarrow.parquet as pq # noqa: E402 from src import paths as P # noqa: E402 from src.derived_manifest import record as dm_record # noqa: E402 from src.windscada.config import farm as ws_farm # noqa: E402 (场名/台数: windscada 侧) from src.windcms import data as CD # noqa: E402 (窗发现: 唯一取用口) from src.windcms.config import KEY_SCALARS, SENSORS, SENSOR_CN # noqa: E402 from src.windcms.config import farm as cms_farm # noqa: E402 (windcms 场配置: FARMS[name]) # ── 标量集 (bounded): KEY_SCALARS 去掉 CrestFactor ──────────────────────────────── # 为什么去 CrestFactor: `src/windcms/report.py` 明写"Kurtosis/CrestFactor 判别力弱(回测实证, 不展示)"; # 且随包旧件的 fleet 键集合恰好 = {Kurtosis, Peak, Rms_HP, Rms_Vel} ∪ {iso_rms_vel (发电机), # rms_200 (主轴承)} —— 本器沿用同一集合, 便于与旧件对拍。 MEAS_SERIES = [m for m in KEY_SCALARS if m != 'CrestFactor'] # ── 部件 ↔ 测点 ──────────────────────────────────────────────────────────────── # 依据 `src/windcms/config.py`: 8 测点 = 2 主轴承 + 4 齿轮箱 + 2 发电机。 # `System Monitor` (磁盘/内存/健康度) 不是部件测点 ⇒ 排除 (它连 SENSORS 都不在)。 COMP_OF_SENSOR = { 'Main_bearing_front': '主轴承', 'Main_bearing_rear': '主轴承', 'Gear_planet': '齿轮箱', 'Gear_IMS': '齿轮箱', 'Gear_HS_rotor_side': '齿轮箱', 'Gear_HS_generator_side': '齿轮箱', 'Generator_DE': '发电机', 'Generator_NDE': '发电机', } # 段数门: 单窗口径下"末三/早三"需要 ≥3 段 (与旧件判据同名同门槛); 周段补充口径另要求 ≥6 段 # (早三/末三**不相交**, 否则末三段里含早三段的中间段, 比值被自己稀释)。 MIN_SEG_WINDOW = 3 MIN_SEG_WEEK = 6 RATIO_HI = 1.3 # 在升判据 (与旧件同名同值): 末/早 ≥ 1.3 RECOVER_HI = 0.5 # 闭环"恢复"判据: 换件后/换件前 ≤ 0.5 FLAT_LO, FLAT_HI = 0.9, 1.1 # 负面对照分桶: 持平带 # ── 工单 → (部件, 是否更换) 规则 (见模块 docstring「诚实边界」) ───────────────────── WO_COMP_KEYS = (('齿轮箱', ('齿轮箱', '齿圈', '行星', '高速轴', '中间轴', '齿轮油')), ('主轴承', ('主轴承', '主轴')), ('发电机', ('发电机',))) PART_PAT = re.compile(r'轴承|齿圈|行星|高速轴|中间轴|箱体|发电机本体') # 耗材类词: 出现即**不算换件** (油/油脂/加注/油泵/滤芯/密封/传感器 —— 换了也不该改变振动) CONSUMABLE_PAT = re.compile(r'油脂|润滑|油位|油泵|滤芯|密封|加注|油管') REPLACE_PAT = re.compile(r'更换|替换|换件') EVENT_CAP = 8 # 每 (台, 部件) 落进 events 的最多工单条数 (最近 N 条), 总数另存 n_events WANT_COLS = ['turbine', 'sensor_name', 'meas_name', 'trigger_time', 'scalar_value', 'ds_size'] # ── 工具 ─────────────────────────────────────────────────────────────────────── def _r6(v): """6 位有效数字的浮点 (0.0028182 → 0.0028182; 4.630123 → 4.63012)。""" if v is None: return None v = float(v) if not np.isfinite(v): return None return float(f'{v:.6g}') def _r3(v): if v is None: return None v = float(v) return float(f'{v:.3g}') if np.isfinite(v) else None def _py(o): """numpy 标量 → python 标量; NaN/Inf → None (JSON 严格可解析, 见 json.dumps(allow_nan=False))。""" if isinstance(o, dict): return {k: _py(v) for k, v in o.items()} if isinstance(o, (list, tuple)): return [_py(v) for v in o] if isinstance(o, (np.integer,)): return int(o) if isinstance(o, (np.floating, float)): f = float(o) return f if np.isfinite(f) else None if isinstance(o, (np.bool_,)): return bool(o) return o def load_scalars(cfg, meas): """各 CMS 窗的标量子集 —— 口径同 `src/windcms/data.py::_read_index/load_scalars` (只取 ds_size==1 的标量行, 只取 SENSORS 内的测点), 但**不落缓存**: 本器只读, 不在产物仓里另生成件。 窗的覆盖区间另取 `src.windcms.data.window_spans` (逐窗读 trigger_time; 页面同源) —— 本器只筛了 6 个标量 × 8 测点, 单看筛后行会把窗起点说晚 (w0316 全表起 03-16 17:27, 筛后起 03-17)。 """ ws = CD.windows(cfg) spans = CD.window_spans(cfg) parts, used = [], {} for w, info in ws.items(): p = pathlib.Path(info['index']) names = set(pq.ParquetFile(p).schema_arrow.names) cols = [c for c in WANT_COLS if c in names] d = pd.read_parquet(p, columns=cols) for c in ('ds_size', 'scalar_value'): if c not in d.columns: d[c] = np.nan d = d[(d.ds_size == 1) & d.meas_name.isin(meas) & d.sensor_name.isin(SENSORS)].copy() d['window'] = w parts.append(d) t = pd.to_datetime(d['trigger_time'], errors='coerce') sp = spans.get(w) or {} used[w] = dict(index=P.rel(p), rows_scalar=int(len(d)), rows_all=int(sp.get('rows') or 0), span_min=str(sp.get('t_min') or '')[:10], span_max=str(sp.get('t_max') or '')[:10], seg_min=str(t.min())[:10], seg_max=str(t.max())[:10], turbines=int(d['turbine'].nunique())) if not parts: return pd.DataFrame(columns=WANT_COLS + ['window', 't', 'val']), used d = pd.concat(parts, ignore_index=True) d['t'] = pd.to_datetime(d['trigger_time'], errors='coerce') d['val'] = pd.to_numeric(d['scalar_value'], errors='coerce') d = d[d['val'].notna() & d['t'].notna()].copy() return d, used def mark_segments(d, n_win, span_min=None): """加 `seg` 列 + 返回 (kind, {段名: 段日起日}, 段序)。见 docstring「分段口径」。 段日: 多窗口径 = 该窗全表覆盖起点 (`window_spans`, 与页面同源); 单窗口径 = 该 ISO 周的**周一** (取周一而不是"该周首条记录日", 是让段日只依赖日历, 不随哪个测点先出数而漂)。 """ span_min = span_min or {} if n_win >= 3: # 多窗口径: 段 = 导出窗 d['seg'] = d['window'].astype(str) kind = 'window' else: # 单窗口径: 段 = 窗内 ISO 周 iso = d['t'].dt.isocalendar() d['seg'] = iso['year'].astype(int).astype(str) + '-W' + iso['week'].astype(int).astype(str).str.zfill(2) kind = 'week' g = d.groupby('seg')['t'].min() order = [s for s, _ in sorted(g.items(), key=lambda kv: kv[1])] span_max = str(d['t'].max())[:10] if kind == 'week': seg_date = {s: str(pd.Timestamp.fromisocalendar(int(s[:4]), int(s[6:]), 1).date()) for s in order} seg_end = {s: str((pd.Timestamp.fromisocalendar(int(s[:4]), int(s[6:]), 1) + pd.Timedelta(days=6)).date()) for s in order} else: seg_date = {s: str(span_min.get(s) or str(g[s])[:10])[:10] for s in order} seg_end = {s: str(d[d.seg == s]['t'].max())[:10] for s in order} return kind, seg_date, order, seg_end def seg_table(d): """(台, 测点, 标量, 段) 的段内中位/样本数。""" g = d.groupby(['turbine', 'sensor_name', 'meas_name', 'seg'])['val'].agg(['median', 'size']) g = g.reset_index().rename(columns={'median': 'med', 'size': 'n'}) return g def trend_table(sg, order, k=3, min_seg=3): """(台, 测点, 标量) 的「末 K 段中位 / 早 K 段中位」比 —— 只对段数 ≥ min_seg 的组合出数。""" rows = [] idx = {s: i for i, s in enumerate(order)} for (t, sen, meas), sub in sg.groupby(['turbine', 'sensor_name', 'meas_name']): m = sub.assign(_i=sub['seg'].map(idx)).dropna(subset=['_i']).sort_values('_i') if len(m) < min_seg: continue vals = m['med'].to_numpy(dtype=float) early, late = float(np.median(vals[:k])), float(np.median(vals[-k:])) if not (early > 0): continue rows.append(dict(turbine=t, sensor=sen, component=COMP_OF_SENSOR.get(sen, ''), measure=meas, early=_r6(early), late=_r6(late), ratio=_r3(late / early), n_seg=int(len(m)), segs=list(m['seg']))) return pd.DataFrame(rows) # ── 工单台账 ─────────────────────────────────────────────────────────────────── def load_workorders(store): """工单台账 → (逐行事件表, 台账范围说明)。台账不存在 ⇒ 空表 + 如实说明 (不造事件)。""" p = pathlib.Path(store) / 'workorders.parquet' if not p.exists(): return pd.DataFrame(columns=['turbine', 'date', 'title', 'replaced', 'comps', 'raw']), \ dict(path=P.rel(p), exists=False, rows=0, note='台账不在位 ⇒ 全部 events 为空数组, 不是"无工单"'), w = pd.read_parquet(p) w['t_report'] = pd.to_datetime(w.get('t_report'), errors='coerce') txt = (w.get('故障名称', pd.Series('', index=w.index)).fillna('') + '|' + w.get('故障描述', pd.Series('', index=w.index)).fillna('') + '|' + w.get('故障位置二级', pd.Series('', index=w.index)).fillna('') + '|' + w.get('故障位置三级', pd.Series('', index=w.index)).fillna('') + '|' + w.get('排查项目', pd.Series('', index=w.index)).fillna('') + '|' + w.get('故障原因', pd.Series('', index=w.index)).fillna('') + '|' + w.get('维修对象', pd.Series('', index=w.index)).fillna('') + '|' + w.get('元器件名称', pd.Series('', index=w.index)).fillna('')) act = w.get('维修动作', pd.Series('', index=w.index)).fillna('').astype(str) obj = (w.get('故障名称', pd.Series('', index=w.index)).fillna('') + '|' + w.get('维修对象', pd.Series('', index=w.index)).fillna('') + '|' + w.get('元器件名称', pd.Series('', index=w.index)).fillna('')) recs = [] for i, s in txt.items(): comps = [c for c, keys in WO_COMP_KEYS if any(k in s for k in keys)] if not comps: continue o = str(obj.iloc[i] if hasattr(obj, 'iloc') else '') repl = bool(REPLACE_PAT.search(act.iloc[i])) and bool(PART_PAT.search(o)) and not CONSUMABLE_PAT.search(o) title = str(w.get('故障名称', pd.Series('', index=w.index)).iloc[i] or '').strip() recs.append(dict(turbine=str(w['turbine'].iloc[i]), date=str(w['t_report'].iloc[i])[:10], title=title[:60], replaced=repl, comps=comps, act=str(act.iloc[i]), obj=o[:60])) ev = pd.DataFrame(recs) # ★ 台账里有**逐字重复行** (实测: 2026-01-03 环网柜 V 柜出线电缆接头三相击穿 同一台同一分钟 7 行; # 换件类里 WTG23 2024-09-12 发电机轴承更换 4 行)。不去重会让 n_replace / n_events 虚高。 n_raw = int(len(ev)) if n_raw: ev['_k'] = ev['turbine'] + '|' + ev['date'] + '|' + ev['title'] + '|' + ev['comps'].map(lambda L: ','.join(L)) ev = ev.drop_duplicates(subset=['_k']).drop(columns=['_k']).reset_index(drop=True) tv = w['t_report'].dropna() info = dict(path=P.rel(p), exists=True, rows=int(len(w)), span=[str(tv.min())[:10], str(tv.max())[:10]] if len(tv) else ['', ''], n_transmission=int(len(ev)), n_transmission_raw=n_raw, n_replace=int(ev['replaced'].sum()) if len(ev) else 0, replace_span=[str(pd.to_datetime(ev[ev.replaced].date).min())[:10], str(pd.to_datetime(ev[ev.replaced].date).max())[:10]] if len(ev) and ev.replaced.any() else None) return ev, info def events_of(ev, turbine, comp): """该 (台, 部件) 的工单事件 (按台/部件关键字映射) → (最近 EVENT_CAP 条明细, 总条数, 其中换件条数)。 ★ 总条数与换件条数在**全量**匹配上数 (明细只是最近的 EVENT_CAP 条): 早期版本用明细数当总数, 一旦该台该部件工单超过 8 条, n_replace 就会**少报**。 """ if not len(ev) or 'comps' not in ev.columns: return [], 0, 0 m = ev[ev['comps'].map(lambda L: comp in L) & (ev['turbine'] == turbine)].sort_values('date') out = [dict(turbine=turbine, comp=comp, date=r['date'], src='工单', title=r['title'], replaced=bool(r['replaced'])) for _, r in m.tail(EVENT_CAP).iterrows()] return out, int(len(m)), int(m['replaced'].sum()) def build(cfg, farm_name='rudong', dry=False): m5 = P.m5(farm_name) store = P.store(farm_name) t0 = time.time() meas = MEAS_SERIES d, used = load_scalars(cfg, meas) if not len(d): print('[X] 一个 CMS 窗索引都没读到 —— 先放数据并跑重算链 (scripts/rebuild_all.py); 本器不写空产物。') return None, 4 kind, seg_date, order, seg_end = mark_segments(d, len(used), {w: v['span_min'] for w, v in used.items()}) sg = seg_table(d) ev, wo_info = load_workorders(store) # 台名单以场配置为准 (38 台), 数据里没有的台如实留空 series (不静默丢台) turbines = [t for t in ws_farm(farm_name).get('turbines', [])] or sorted(d['turbine'].unique()) no_data = [t for t in turbines if t not in set(d['turbine'].unique())] # 数据覆盖区间取**全表** (src.windcms.data.window_spans, 与页面同源): 只筛 6 标量 × 8 测点会把窗起点说晚一天 span = [min(v['span_min'] for v in used.values()), max(v['span_max'] for v in used.values())] print(f'部件历史生成端 · 场={farm_name}') print(f' 窗 {len(used)} 个: ' + '; '.join( f'{w}(全表 {v["span_min"]}→{v["span_max"]}, {v["rows_all"]}行; 标量 {v["rows_scalar"]}行)' for w, v in used.items())) print(f' 标量行 (ds_size==1 ∧ {len(meas)} 标量 ∧ {len(SENSORS)} 测点): {len(d)}' + (f' · 场配置 {len(turbines)} 台, 无数据的台: {no_data}' if no_data else f' · 台数 {len(turbines)}')) print(f' 分段口径: segments_kind={kind} · 段 {len(order)} 个: {order[0]}…{order[-1]} ' f'(段日起 {seg_date[order[0]]}→{seg_date[order[-1]]}; 数据覆盖 {span[0]}→{span[1]})') print(f' 段内样本 (台×测点×标量×段) 中位: {int(sg["n"].median())} · 最少 {int(sg["n"].min())} · 最多 {int(sg["n"].max())}') print(f' 工单台账: {wo_info["path"]} 行={wo_info["rows"]} 跨度={wo_info.get("span")} ' f'传动链件={wo_info.get("n_transmission")} 换件={wo_info.get("n_replace")} ' f'换件跨度={wo_info.get("replace_span")}') if len(ev): print(' 换件类工单 (全量, 判据 candidate): ' + '; '.join( f'{r["turbine"]} {"/".join(r["comps"])} {r["date"]} {r["title"][:18]}' for _, r in ev[ev['replaced']].sort_values('date').iterrows())) # ── trend: 单窗口径的末三/早三 (判据) 与 周段的末三/早三 (补充, 不相交) ── if kind == 'window': tw = trend_table(sg, order, k=3, min_seg=MIN_SEG_WINDOW) else: tw = trend_table(sg, order, k=3, min_seg=MIN_SEG_WEEK) tw = tw.sort_values('ratio', ascending=False) if len(tw) else tw wo_ok = bool(wo_info.get('exists') and wo_info.get('span') and wo_info['span'][0] and wo_info['span'][0] <= span[1] and wo_info['span'][1] >= span[0]) # ── components ── comps_out = [] for t in turbines: for comp in ('主轴承', '齿轮箱', '发电机'): sens = [s for s, c in COMP_OF_SENSOR.items() if c == comp] series = {} n_seg_c = 0 for sen in sens: for ms in meas: sub = sg[(sg.turbine == t) & (sg.sensor_name == sen) & (sg.meas_name == ms)] if not len(sub): continue m = sub.set_index('seg')['med'].reindex(order).dropna() if not len(m): continue base = float(np.median(m.to_numpy(dtype=float))) if not (base > 0): continue entries = [] for s in order: if s not in m.index: continue v = float(m.loc[s]) n = int(sub.loc[sub.seg == s, 'n'].iloc[0]) entries.append(dict(seg=s, med=_r6(v), x_self=_r3(v / base), n=n)) series[f'{sen}|{ms}'] = entries n_seg_c = max(n_seg_c, len(entries)) e_list, n_ev, n_repl = events_of(ev, t, comp) comps_out.append(dict(turbine=t, component=comp, n_windows=int(n_seg_c), n_events=n_ev, n_replace=n_repl, events_capped=bool(n_ev > len(e_list)), series=series, events=e_list)) # ── summary ── # ★判据一: 只认多窗口径 (本机必然空) —— 不降门槛凑数 win_trend = trend_table(sg, order, k=3, min_seg=MIN_SEG_WINDOW) if kind == 'window' else pd.DataFrame() rising = [] if kind == 'window' and len(win_trend): for (t, comp), sub in win_trend.groupby(['turbine', 'component']): w = sub.sort_values('ratio', ascending=False).iloc[0] if w['ratio'] >= RATIO_HI: _, n_ev, _ = events_of(ev, t, comp) rising.append(dict(turbine=t, component=comp, ratio=float(w['ratio']), early=float(w['early']), late=float(w['late']), n_events=n_ev, 工单可对=wo_ok, sensor=w['sensor'], measure=w['measure'], n_windows=int(w['n_seg']))) rising = sorted(rising, key=lambda r: -r['ratio']) # 补充诊断 (非★, 消费者忽略): 周段末三/早三 —— 逐行带段内最小样本数, 让"薄段噪声"可见 diag = [] if len(tw): nmin = sg.groupby(['turbine', 'sensor_name', 'meas_name'])['n'].min() for _, r in tw[tw['ratio'] >= RATIO_HI].iterrows(): _, n_ev, _ = events_of(ev, r['turbine'], r['component']) diag.append(dict(turbine=r['turbine'], component=r['component'], sensor=r['sensor'], measure=r['measure'], early=float(r['early']), late=float(r['late']), ratio=float(r['ratio']), n_seg=int(r['n_seg']), seg_n_min=int(nmin.get((r['turbine'], r['sensor'], r['measure']), 0)), n_events=n_ev, 工单可对=wo_ok)) diag_all = len(diag) diag = diag[:20] # ★判据二: 换件闭环 —— 需要"换件日之前有段 ∧ 之后有段" loop, loop_skip, repl_pairs = [], [], [] repl = ev[ev['replaced']] if len(ev) and 'replaced' in ev.columns else ev fusion = {} fp = m5 / 'fusion_38.csv' if fp.exists(): # 检出候选旗标 (融合/CMS 面), 只作 verdict 的佐证 fu = pd.read_csv(fp) for _, r in fu.iterrows(): fusion[str(r.get('台'))] = dict(融合=str(r.get('融合') or ''), CMS=str(r.get('CMS') or '')) for _, r in (repl.iterrows() if len(repl) else []): for comp in r['comps']: repl_pairs.append((r['turbine'], comp, r['date'])) pre_segs = [s for s in order if seg_end[s] < r['date']] # 段**结束**在工单日之前 post_segs = [s for s in order if seg_date[s] > r['date']] # 段**开始**在工单日之后 if not pre_segs or not post_segs: loop_skip.append(dict(turbine=r['turbine'], component=comp, date=r['date'], title=r['title'], why=f'段跨度 {span[0]}~{span[1]} 内' + ('无换件前段' if not pre_segs else '无换件后段'))) continue sens = [s for s, c in COMP_OF_SENSOR.items() if c == comp] cand = [] for sen in sens: for ms in meas: sub = sg[(sg.turbine == r['turbine']) & (sg.sensor_name == sen) & (sg.meas_name == ms)] if not len(sub): continue a = sub[sub.seg.isin(pre_segs)]['med'] b = sub[sub.seg.isin(post_segs)]['med'] if len(a) and len(b) and float(np.median(a)) > 0: cand.append((ms, sen, float(np.median(a)), float(np.median(b)))) if not cand: loop_skip.append(dict(turbine=r['turbine'], component=comp, date=r['date'], title=r['title'], why='换件前后段内该部件无可用标量')) continue ms, sen, pre, post = max(cand, key=lambda c: c[2]) # 取换前幅值最大者 (最不利/最有分辨力) ratio = post / pre f = fusion.get(r['turbine']) or {} hit = bool(f) and (f.get('融合') not in ('', '正常') or '红' in f.get('CMS', '') and '红0' not in f.get('CMS', '')) if ratio <= RECOVER_HI: verdict = '恢复' if hit else '回落 (非检出台, 不作端到端实证)' elif ratio < FLAT_LO: verdict = '部分回落' elif ratio <= FLAT_HI: verdict = '持平' else: verdict = '未降/上升' loop.append(dict(turbine=r['turbine'], component=comp, ratio=_r3(ratio), pre=_r6(pre), post=_r6(post), replace_date=r['date'], title=r['title'], verdict=verdict, sensor=sen, measure=ms, 检出候选=hit, 融合=(f.get('融合') or None), pre_segs=pre_segs, post_segs=post_segs)) loop = sorted(loop, key=lambda r: r['ratio']) # 负面对照: 工单后振动未降/持平 (非换件的传动链工单, 换前后各有段) # 不可算的原因**分桶计数**: 单看一个 1512 的总数会把"工单在窗外"和"窗内但无前后段"混成一团。 neg_rows, neg_noc = [], dict(工单在段跨度外=0, 窗内但无前段或后段=0, 该部件无标量=0) nonrepl = ev[~ev['replaced']] if len(ev) and 'replaced' in ev.columns else ev for _, r in (nonrepl.iterrows() if len(nonrepl) else []): for comp in r['comps']: # ★ 前/后段都**不含工单所在那一段**: 段界线按"段结束 < 工单日"与"段开始 > 工单日"判, # 否则工单落在段中间时, 该段自己会被当成"换前", 把换后的数据混进 pre (实测污染过 3 条)。 pre_segs = [s for s in order if seg_end[s] < r['date']] post_segs = [s for s in order if seg_date[s] > r['date']] if r['date'] < seg_date[order[0]] or r['date'] > seg_end[order[-1]]: neg_noc['工单在段跨度外'] += 1 continue if not pre_segs or not post_segs: neg_noc['窗内但无前段或后段'] += 1 continue sens = [s for s, c in COMP_OF_SENSOR.items() if c == comp] cand = [] for sen in sens: for ms in meas: sub = sg[(sg.turbine == r['turbine']) & (sg.sensor_name == sen) & (sg.meas_name == ms)] if not len(sub): continue a = sub[sub.seg.isin(pre_segs)]['med'] b = sub[sub.seg.isin(post_segs)]['med'] if len(a) and len(b) and float(np.median(a)) > 0: cand.append((ms, sen, float(np.median(a)), float(np.median(b)))) if not cand: neg_noc['该部件无标量'] += 1 continue ms, sen, pre, post = max(cand, key=lambda c: c[2]) neg_rows.append(dict(turbine=r['turbine'], component=comp, date=r['date'], title=r['title'], sensor=sen, measure=ms, pre=_r6(pre), post=_r6(post), ratio=_r3(post / pre))) bucket = dict(下降=sum(1 for r in neg_rows if r['ratio'] < FLAT_LO), 持平=sum(1 for r in neg_rows if FLAT_LO <= r['ratio'] <= FLAT_HI), 上升=sum(1 for r in neg_rows if r['ratio'] > FLAT_HI)) neg_noc_total = sum(neg_noc.values()) n_comp_rows = sum(1 for c in comps_out for _ in c['series']) win_txt = '; '.join(f'{w} {v["span_min"]}~{v["span_max"]}' for w, v in used.items()) n_pair_inspan = sum(1 for _, _, dt in repl_pairs if span[0] <= dt <= span[1]) 判读 = ( f'分段口径 segments_kind={kind}: 本机只解析出 {len(used)} 个导出色窗 ({win_txt}) ' f'⇒ 以窗内 ISO 周段为分段 ({len(order)} 段 {order[0]}…{order[-1]}, {span[0]}~{span[1]}); ' f'多窗 (≥3) 时自动改回多窗口径。' f'★当前在升 (末三窗/早三窗 ≥{RATIO_HI}) 本机{("有 " + str(len(rising)) + " 行") if rising else "不可算 ⇒ 空表"}: ' f'该判据要求多窗口径且 ≥{MIN_SEG_WINDOW} 窗, 本机 {len(used)} 窗。' f'★换件闭环 本机{("有 " + str(len(loop)) + " 行") if loop else "不可算 ⇒ 空表"}: 要求换件工单在段跨度内' f'且其前后各有至少一段; 台账 {wo_info.get("span")} 共 {wo_info.get("n_transmission")} 条传动链工单, ' f'其中换件类 {wo_info.get("n_replace")} 条 (跨度 {wo_info.get("replace_span")}) 展开为 {len(repl_pairs)} 个 ' f'(台,部件) 组合, 段跨度内 {n_pair_inspan} 个 ⇒ 成闭环 {len(loop)} 行, 其余 {len(loop_skip)} 个因' f'"无换件前段/后段"逐条记入 meta.未成闭环。' f'补充口径 (非判据): 周段末三/早三比值 ≥{RATIO_HI} 的候选 {diag_all} 条 (列出前 {len(diag)} 条); ' f'周段每 (台,测点,标量) 只有 {int(sg["n"].min())}~{int(sg["n"].max())} 条标量记录/段 ' f'(中位 {int(sg["n"].median())}) ⇒ 该比值只说明窗内走向, 不作在升判据 (明细逐行给出 seg_n_min)。' f'负面对照: 传动链非换件工单里换前后各有段的 {len(neg_rows)} 条可算 (下降 {bucket["下降"]} · 持平 {bucket["持平"]} ' f'· 上升 {bucket["上升"]}), 其余 {neg_noc_total} 条不可算 ({neg_noc})。') seg_kind_txt = '导出窗' if kind == 'window' else f'窗内 ISO 周 (本机窗数 {len(used)} < 3 ⇒ 降级; 多窗时自动改回)' meta = dict( built=time.strftime('%Y-%m-%d'), farm=farm_name, by='scripts/component_history_build.py', windows={w: v['span_min'] for w, v in used.items()}, windows_detail={w: dict(rows_all=v['rows_all'], rows_scalar=v['rows_scalar'], span_min=v['span_min'], span_max=v['span_max'], seg_min=v['seg_min'], seg_max=v['seg_max'], turbines=v['turbines']) for w, v in used.items()}, segments_kind=kind, segments=seg_date, n_segments=len(order), 来源声明=('本文件由 scripts/component_history_build.py **现算**写成: 输入 = ' 'outputs/<场>/m5_cms_tcm/windows/*/index.parquet (data/raw 重算产物, 只读 ds_size==1 标量行) ' '+ outputs/<场>/windscada/workorders.parquet (工单台账) ' '+ outputs/<场>/m5_cms_tcm/fusion_38.csv (仅作"检出候选"佐证)。' '不含随包旧件 (人工裁决/厂家锚) 的数值搬运; 算不出来的 (在升/闭环) 一律空表 + 边界说明, 不填估值。' 'components[].n_windows = 该部件的**段数** (segments_kind="week" 时即窗内周段数, 不是导出窗数)。'), 口径=(f'标量: CMS 窗索引 ds_size==1 行, 标量集 = KEY_SCALARS 去 CrestFactor ({", ".join(meas)}); ' f'测点 → 部件: Main_bearing_*→主轴承, Gear_*→齿轮箱, Generator_*→发电机 (System Monitor 非部件测点, 排除)。' f'段: {seg_kind_txt}。' f'series 值 = 该段该 (测点|标量) 标量**中位**; 归一 x_self = 段中位 / 该台该测点该标量**全段中位数** ' f'(随包旧件同口径, 旧件写作 x_self; 本件不出 x_fleet —— 那需要跨台同工况对照)。' f'在升判据 = 末三窗中位 / 早三窗中位 ≥ {RATIO_HI}, 需多窗口径且 ≥{MIN_SEG_WINDOW} 窗。' f'闭环判据 = 换件工单日之前有段 ∧ 之后有段; pre/post = 该部件换前/后段中位 (取换前幅值最大的测点|标量), ' f'verdict: ratio ≤ {RECOVER_HI} 且该台为检出候选 ⇒ 恢复, 否则按 部分回落/持平/未降 如实记。' f'工单口径: 部件由关键词映射 (齿轮箱/齿圈/行星/高速轴/中间轴/齿轮油 → 齿轮箱; 主轴承/主轴 → 主轴承; ' f'发电机 → 发电机); replaced = 维修动作含 更换|替换|换件 ∧ 目标文本含部件词 ' f'(轴承|齿圈|行星|高速轴|中间轴|箱体|发电机本体) ∧ 目标文本无耗材词 ' f'(油脂|润滑|油位|油泵|滤芯|密封|加注|油管) —— 换油/换传感器不算换件。'), inputs=dict(scalars=[v['index'] for v in used.values()], workorders=wo_info.get('path'), fusion=P.rel(fp) if fp.exists() else None), 未成闭环=loop_skip[:40], 未成闭环总数=len(loop_skip), 台账=wo_info, ) meta['★数据边界'] = ( f'① 本机只有 {len(used)} 个 CMS 导出色窗 ({win_txt}) ⇒ 在升 (末三窗/早三窗) 与 多窗口径的健康期分段' f'**都不可算** ⇒ summary["★当前在升…"] 为空表; 本件以窗内周段 ({len(order)} 段) 作补充诊断, ' f'只描述窗内走向, 不作判据。' f'② 换件闭环需要"换件日前后各有段": 工单台账 {wo_info.get("path")} 跨度 {wo_info.get("span")} ' f'(传动链件 {wo_info.get("n_transmission")} 条, 其中换件类 {wo_info.get("n_replace")} 条, ' f'跨度 {wo_info.get("replace_span")}), 与段跨度 {seg_date[order[0]]}~{seg_date[order[-1]]} **无交叠** ' f'⇒ summary["★换件闭环…"] 为空表。★ 这两张空表是数据边界, 不是"没有在升机组/没有闭环案例"; ' f'站点需 (a) 补导 ≥3 个 CMS 色窗 (尤其换件日前后各一窗), (b) 补 2026 年传动链换件工单, ' f'二者到位后本脚本自动出数。' f'③ 段内样本薄: 周段每 (台,测点,标量) 仅 {int(sg["n"].min())}~{int(sg["n"].max())} 条标量记录 ' f'(中位 {int(sg["n"].median())}), 中位数据此散布 —— series 的 n 逐条给出, 读者自行判断可信度。') meta['★换件断点'] = '换件前后不可混算趋势; 本表在 events 里标 replaced=true, 读者须自行分段读 series' payload = dict( meta=meta, components=_py(comps_out), summary=_py({ '★换件闭环 (检出→换件→恢复, 算法端到端验证)': loop, '★当前在升 (末三窗/早三窗 ≥1.3)': rising, '窗内周段趋势 (单窗补充口径, 仅参考不作判据)': diag, '窗内周段趋势·候选总数': diag_all, '工单后振动未降/持平 (多为油位/传感器/加注类, 非轴承更换 ⇒ 振动本不该变)': dict( 可算=len(neg_rows), 不可算=neg_noc_total, 不可算原因=neg_noc, 分桶=bucket, 口径='逐条传动链非换件工单: 换件日前后各有至少一段 ⇒ 可算; 指标取该部件换前幅值最大的 (测点|标量) 的 ' f'post/pre 比值; 分桶 下降<{FLAT_LO} · 持平 {FLAT_LO}~{FLAT_HI} · 上升>{FLAT_HI}', 明细=sorted(neg_rows, key=lambda r: -r['ratio'])[:12]), '判读': 判读}), ) txt = json.dumps(payload, ensure_ascii=False, indent=1, allow_nan=False, default=str) out = m5 / 'component_history.json' print(f'\n 段表行 {len(sg)} · 部件行 {len(comps_out)} · series 键 {n_comp_rows} · ' f'events {sum(len(c["events"]) for c in comps_out)} 条') print(f' summary: ★换件闭环 {len(loop)} 行 · ★当前在升 {len(rising)} 行 · ' f'周段候选 {diag_all} 条(列 {len(diag)}) · 负面对照可算 {len(neg_rows)}/不可算 {neg_noc_total} {neg_noc}') if neg_rows: print(' 负面对照明细: ' + '; '.join(f'{r["turbine"]} {r["component"]} {r["date"]} {r["title"][:16]} ' f'{r["sensor"]}|{r["measure"]} {r["ratio"]}×' for r in sorted(neg_rows, key=lambda r: -r['ratio']))) if diag: print(' 周段候选 top: ' + '; '.join(f'{r["turbine"]} {r["component"]} {r["sensor"]}|{r["measure"]} ' f'{r["ratio"]}× ({r["early"]}→{r["late"]}, {r["n_seg"]}段)' for r in diag[:5])) if loop_skip: print(f' 换件工单未成闭环 {len(loop_skip)} 条, 例: ' + '; '.join(f'{s["turbine"]} {s["component"]} {s["date"]} {s["why"]}' for s in loop_skip[:3])) print(f' 产物体积预估 {len(txt.encode("utf-8")) / 1024:.0f} KB ({P.rel(out)}) · 用时 {time.time() - t0:.1f}s') if dry: print('\n[dry-run] 不写盘、不登记。') return payload, 0 m5.mkdir(parents=True, exist_ok=True) out.write_text(txt, encoding='utf-8') dm_record(P.out_root(farm_name), {out.relative_to(P.out_root(farm_name)).as_posix():'scripts/component_history_build.py (CMS 窗标量 ds_size==1: 段中位/x_self 序列 + ' '工单台账事件; 在升/闭环判据与空表情形见产物 meta.★数据边界)'}, by='scripts/component_history_build.py') print(f'\n[OK] 已写 {P.rel(out)} ({out.stat().st_size / 1024:.0f} KB) · 已自登记') return payload, 0 def main() -> int: ap = argparse.ArgumentParser(description='部件历史生成端 (m5_cms_tcm/component_history.json)') ap.add_argument('--farm', default='rudong') ap.add_argument('--dry-run', action='store_true', help='只算不写盘 (打印漏斗与体积预估)') args = ap.parse_args() cfg = cms_farm(args.farm) _, rc = build(cfg, args.farm, dry=args.dry_run) return rc if __name__ == '__main__': sys.exit(main())