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