component_history_build.py 41 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641
  1. #!/usr/bin/env python3
  2. # -*- coding: utf-8 -*-
  3. r"""部件历史生成端 —— 补上"包内没有生成端"的 `m5_cms_tcm/component_history.json` (用户令 2026-09-19)。
  4. ## 为什么要写它
  5. 用户令:「**所有的计算 均要形成 观澜的源代码**,确保 观澜 在别的电脑安装后,系统运行正常。」
  6. `component_history.json` 是**随包快照**(族表 kind=shipped、gen=None,见
  7. `scripts/products_reverse_audit.py::HUMAN_SNAPSHOTS`),随包件里的"在升/闭环/换件日期"含人工裁决与
  8. 厂家锚,**不是测量数据的函数**。于是另一台机器上装完、放数据、重算之后:
  9. `src/ontology/trend_ingest.py` 拿不到在升/闭环证据(图谱里两类 Evidence 全缺)、
  10. `scripts/windscada_serve.py` 的决策链盘「振动在升/闭环」两格永远空。
  11. 本脚本把它**从已有产物算出来**(CMS 窗索引 + 工单台账),口径写在下面与产物 meta 里。
  12. ## 产物与消费者 (字段契约)
  13. outputs/<场>/m5_cms_tcm/component_history.json
  14. meta built / windows / segments / segments_kind / 口径 / ★数据边界
  15. components[] turbine, component, n_windows, n_replace, series, events
  16. summary['★当前在升 (末三窗/早三窗 ≥1.3)'] ← src/ontology/trend_ingest.py 逐行读
  17. turbine, component, ratio, early, late, n_events, 工单可对
  18. summary['★换件闭环 (检出→换件→恢复, 算法端到端验证)'] ← 同上 + scripts/windscada_serve.py 链盘
  19. turbine, component, ratio, pre, post, replace_date, title, verdict
  20. summary[非★键] 消费端只认 ★ 前缀的键 ⇒ 非★键是给人/给诊断看的, 不被误读
  21. ★ `component` 必须是 `src/ontology/trend_ingest.py::COMP_MAP` 里的中文名 (齿轮箱/发电机/主轴承/变桨/偏航/变流器),
  22. 否则该行被消费者跳过。本器只出**齿轮箱/发电机/主轴承**三类 —— CMS 索引里只有这三类的测点
  23. (Main_bearing_* / Gear_* / Generator_*),变桨/偏航/变流器在本产物无振动测点,不臆造。
  24. ## 分段 (segment) 口径 —— 两种模式, 自动选
  25. 参考口径是**导出窗**为段 (末三窗/早三窗; 健康期窗中位)。本机实况 (2026-09-19 实测):
  26. `src/windcms/config.py::FARMS['rudong']['windows']` 只解析出**一个**窗 `w0316`
  27. (`outputs/rudong/m5_cms_tcm/windows/w0316/index.parquet`,2,066,686 行 × 54 列,
  28. trigger_time 2026-03-16 17:27 → 2026-04-21 10:18)。窗数 < 3 ⇒ 自动降为
  29. **窗内周段** (ISO 周, 2026-W12 … 2026-W17 共 6 段)。两种模式都在本文件里实现:
  30. len(windows) >= 3 segments_kind='window' 段名=窗名, 段日=该窗 trigger_time 最早日
  31. len(windows) < 3 segments_kind='week' 段名=ISO 周 (2026-W12), 段日=该周周一
  32. **多窗口径未在本机取数验证**(本机只有 1 窗 ⇒ 该分支只有单测/纸面正确性,不给"已验证"的说法)。
  33. ## 数据现实与诚实边界 (不造数)
  34. * **`★当前在升` 本机必然是空表**:判据要求"末三窗/早三窗"且**多窗口径**;本机 1 窗 ⇒ 不可算。
  35. ★ 不降门槛、不换判据去凑一行。作为补充, 出**非★**键
  36. `窗内周段趋势 (单窗补充口径, 仅参考不作判据)`:同一条比值算式 (末三段/早三段, 要求 6 段以取到
  37. **不相交**的早/晚两组) 在周段上跑出来是什么, 就写什么。
  38. * **`★换件闭环` 本机为空**:闭环要求"该台该部件的更换/换件工单,且段跨度里**既有段在换件日之前、
  39. 又有段在其后**"(前段须**结束**在工单日之前、后段须**开始**在工单日之后 —— 工单所在那一段两边都不算,
  40. 否则换后数据会混进"换前",实测污染过 3 条)。工单台账 (`outputs/<场>/windscada/workorders.parquet`,
  41. 5,876 行) 实测跨 **2020-01-03 → 2026-07-08**,其中传动链件 903 行(去重后)/换件类 12 条 ——
  42. 但**全部落在段跨度 (2026-03-16~04-21) 之外** (换件工单最晚 2025-10-14),窗内 0 条 ⇒ 无"换件前段"可比
  43. ⇒ 空表。空表本身就是结论。
  44. * 台账有**逐字重复行**(实测 2026-01-03 环网柜那条同一台同一分钟 7 行;换件类里 WTG23 2024-09-12
  45. 发电机轴承更换 4 行)⇒ 本器按 (台,日,标题,部件) 去重后再计数,原始/去重两个数都写进 meta.台账。
  46. * 段内样本数一并写进 series 的 `n` —— 周段只有 2~142 个标量记录 (中位 14), 中位数据此散布, 读者可自行判断。
  47. ## 用法
  48. python scripts/component_history_build.py # 算并写盘 + 自登记
  49. python scripts/component_history_build.py --dry-run # 只算不写, 打印漏斗
  50. python scripts/component_history_build.py --farm rudong # 场名 (默认 rudong)
  51. """
  52. from __future__ import annotations
  53. try:
  54. from app_common.app_common_guanlan.api import install_root as _install_root
  55. except ImportError: # 理论不可达;包结构异常时回退到按位置上跳
  56. from pathlib import Path as _P
  57. def _install_root(_f): return _P(_f).resolve().parents[3]
  58. import argparse
  59. import json
  60. import pathlib
  61. import re
  62. import sys
  63. import time
  64. ROOT = _install_root(__file__) # 模块化后按标记找安装根(原 parents[1] 已不成立)
  65. sys.path.insert(0, str(ROOT))
  66. # ★GBK 控制台/日志重定向不再因 '²' '⇒' 这类字符抛 UnicodeEncodeError (2026-09-19 远端实逮:
  67. # `baseline_38_build.py` 打印 'm/s²' 时在 cp936 下 rc=1, 整条恢复链少一件产物)。
  68. # 只把**不可编码字符替换掉**, 不改流编码 —— 中文照旧可读。
  69. try:
  70. sys.stdout.reconfigure(errors='replace')
  71. except Exception:
  72. pass
  73. import numpy as np # noqa: E402
  74. import pandas as pd # noqa: E402
  75. import pyarrow.parquet as pq # noqa: E402
  76. from src import paths as P # noqa: E402
  77. from src.derived_manifest import record as dm_record # noqa: E402
  78. from src.windscada.config import farm as ws_farm # noqa: E402 (场名/台数: windscada 侧)
  79. from src.windcms import data as CD # noqa: E402 (窗发现: 唯一取用口)
  80. from src.windcms.config import KEY_SCALARS, SENSORS, SENSOR_CN # noqa: E402
  81. from src.windcms.config import farm as cms_farm # noqa: E402 (windcms 场配置: FARMS[name])
  82. # ── 标量集 (bounded): KEY_SCALARS 去掉 CrestFactor ────────────────────────────────
  83. # 为什么去 CrestFactor: `src/windcms/report.py` 明写"Kurtosis/CrestFactor 判别力弱(回测实证, 不展示)";
  84. # 且随包旧件的 fleet 键集合恰好 = {Kurtosis, Peak, Rms_HP, Rms_Vel} ∪ {iso_rms_vel (发电机),
  85. # rms_200 (主轴承)} —— 本器沿用同一集合, 便于与旧件对拍。
  86. MEAS_SERIES = [m for m in KEY_SCALARS if m != 'CrestFactor']
  87. # ── 部件 ↔ 测点 ────────────────────────────────────────────────────────────────
  88. # 依据 `src/windcms/config.py`: 8 测点 = 2 主轴承 + 4 齿轮箱 + 2 发电机。
  89. # `System Monitor` (磁盘/内存/健康度) 不是部件测点 ⇒ 排除 (它连 SENSORS 都不在)。
  90. COMP_OF_SENSOR = {
  91. 'Main_bearing_front': '主轴承', 'Main_bearing_rear': '主轴承',
  92. 'Gear_planet': '齿轮箱', 'Gear_IMS': '齿轮箱',
  93. 'Gear_HS_rotor_side': '齿轮箱', 'Gear_HS_generator_side': '齿轮箱',
  94. 'Generator_DE': '发电机', 'Generator_NDE': '发电机',
  95. }
  96. # 段数门: 单窗口径下"末三/早三"需要 ≥3 段 (与旧件判据同名同门槛); 周段补充口径另要求 ≥6 段
  97. # (早三/末三**不相交**, 否则末三段里含早三段的中间段, 比值被自己稀释)。
  98. MIN_SEG_WINDOW = 3
  99. MIN_SEG_WEEK = 6
  100. RATIO_HI = 1.3 # 在升判据 (与旧件同名同值): 末/早 ≥ 1.3
  101. RECOVER_HI = 0.5 # 闭环"恢复"判据: 换件后/换件前 ≤ 0.5
  102. FLAT_LO, FLAT_HI = 0.9, 1.1 # 负面对照分桶: 持平带
  103. # ── 工单 → (部件, 是否更换) 规则 (见模块 docstring「诚实边界」) ─────────────────────
  104. WO_COMP_KEYS = (('齿轮箱', ('齿轮箱', '齿圈', '行星', '高速轴', '中间轴', '齿轮油')),
  105. ('主轴承', ('主轴承', '主轴')),
  106. ('发电机', ('发电机',)))
  107. PART_PAT = re.compile(r'轴承|齿圈|行星|高速轴|中间轴|箱体|发电机本体')
  108. # 耗材类词: 出现即**不算换件** (油/油脂/加注/油泵/滤芯/密封/传感器 —— 换了也不该改变振动)
  109. CONSUMABLE_PAT = re.compile(r'油脂|润滑|油位|油泵|滤芯|密封|加注|油管')
  110. REPLACE_PAT = re.compile(r'更换|替换|换件')
  111. EVENT_CAP = 8 # 每 (台, 部件) 落进 events 的最多工单条数 (最近 N 条), 总数另存 n_events
  112. WANT_COLS = ['turbine', 'sensor_name', 'meas_name', 'trigger_time', 'scalar_value', 'ds_size']
  113. # ── 工具 ───────────────────────────────────────────────────────────────────────
  114. def _r6(v):
  115. """6 位有效数字的浮点 (0.0028182 → 0.0028182; 4.630123 → 4.63012)。"""
  116. if v is None:
  117. return None
  118. v = float(v)
  119. if not np.isfinite(v):
  120. return None
  121. return float(f'{v:.6g}')
  122. def _r3(v):
  123. if v is None:
  124. return None
  125. v = float(v)
  126. return float(f'{v:.3g}') if np.isfinite(v) else None
  127. def _py(o):
  128. """numpy 标量 → python 标量; NaN/Inf → None (JSON 严格可解析, 见 json.dumps(allow_nan=False))。"""
  129. if isinstance(o, dict):
  130. return {k: _py(v) for k, v in o.items()}
  131. if isinstance(o, (list, tuple)):
  132. return [_py(v) for v in o]
  133. if isinstance(o, (np.integer,)):
  134. return int(o)
  135. if isinstance(o, (np.floating, float)):
  136. f = float(o)
  137. return f if np.isfinite(f) else None
  138. if isinstance(o, (np.bool_,)):
  139. return bool(o)
  140. return o
  141. def load_scalars(cfg, meas):
  142. """各 CMS 窗的标量子集 —— 口径同 `src/windcms/data.py::_read_index/load_scalars`
  143. (只取 ds_size==1 的标量行, 只取 SENSORS 内的测点), 但**不落缓存**: 本器只读, 不在产物仓里另生成件。
  144. 窗的覆盖区间另取 `src.windcms.data.window_spans` (逐窗读 trigger_time; 页面同源) ——
  145. 本器只筛了 6 个标量 × 8 测点, 单看筛后行会把窗起点说晚 (w0316 全表起 03-16 17:27, 筛后起 03-17)。
  146. """
  147. ws = CD.windows(cfg)
  148. spans = CD.window_spans(cfg)
  149. parts, used = [], {}
  150. for w, info in ws.items():
  151. p = pathlib.Path(info['index'])
  152. names = set(pq.ParquetFile(p).schema_arrow.names)
  153. cols = [c for c in WANT_COLS if c in names]
  154. d = pd.read_parquet(p, columns=cols)
  155. for c in ('ds_size', 'scalar_value'):
  156. if c not in d.columns:
  157. d[c] = np.nan
  158. d = d[(d.ds_size == 1) & d.meas_name.isin(meas) & d.sensor_name.isin(SENSORS)].copy()
  159. d['window'] = w
  160. parts.append(d)
  161. t = pd.to_datetime(d['trigger_time'], errors='coerce')
  162. sp = spans.get(w) or {}
  163. used[w] = dict(index=P.rel(p), rows_scalar=int(len(d)), rows_all=int(sp.get('rows') or 0),
  164. span_min=str(sp.get('t_min') or '')[:10], span_max=str(sp.get('t_max') or '')[:10],
  165. seg_min=str(t.min())[:10], seg_max=str(t.max())[:10],
  166. turbines=int(d['turbine'].nunique()))
  167. if not parts:
  168. return pd.DataFrame(columns=WANT_COLS + ['window', 't', 'val']), used
  169. d = pd.concat(parts, ignore_index=True)
  170. d['t'] = pd.to_datetime(d['trigger_time'], errors='coerce')
  171. d['val'] = pd.to_numeric(d['scalar_value'], errors='coerce')
  172. d = d[d['val'].notna() & d['t'].notna()].copy()
  173. return d, used
  174. def mark_segments(d, n_win, span_min=None):
  175. """加 `seg` 列 + 返回 (kind, {段名: 段日起日}, 段序)。见 docstring「分段口径」。
  176. 段日: 多窗口径 = 该窗全表覆盖起点 (`window_spans`, 与页面同源); 单窗口径 = 该 ISO 周的**周一**
  177. (取周一而不是"该周首条记录日", 是让段日只依赖日历, 不随哪个测点先出数而漂)。
  178. """
  179. span_min = span_min or {}
  180. if n_win >= 3: # 多窗口径: 段 = 导出窗
  181. d['seg'] = d['window'].astype(str)
  182. kind = 'window'
  183. else: # 单窗口径: 段 = 窗内 ISO 周
  184. iso = d['t'].dt.isocalendar()
  185. d['seg'] = iso['year'].astype(int).astype(str) + '-W' + iso['week'].astype(int).astype(str).str.zfill(2)
  186. kind = 'week'
  187. g = d.groupby('seg')['t'].min()
  188. order = [s for s, _ in sorted(g.items(), key=lambda kv: kv[1])]
  189. span_max = str(d['t'].max())[:10]
  190. if kind == 'week':
  191. seg_date = {s: str(pd.Timestamp.fromisocalendar(int(s[:4]), int(s[6:]), 1).date()) for s in order}
  192. seg_end = {s: str((pd.Timestamp.fromisocalendar(int(s[:4]), int(s[6:]), 1) + pd.Timedelta(days=6)).date())
  193. for s in order}
  194. else:
  195. seg_date = {s: str(span_min.get(s) or str(g[s])[:10])[:10] for s in order}
  196. seg_end = {s: str(d[d.seg == s]['t'].max())[:10] for s in order}
  197. return kind, seg_date, order, seg_end
  198. def seg_table(d):
  199. """(台, 测点, 标量, 段) 的段内中位/样本数。"""
  200. g = d.groupby(['turbine', 'sensor_name', 'meas_name', 'seg'])['val'].agg(['median', 'size'])
  201. g = g.reset_index().rename(columns={'median': 'med', 'size': 'n'})
  202. return g
  203. def trend_table(sg, order, k=3, min_seg=3):
  204. """(台, 测点, 标量) 的「末 K 段中位 / 早 K 段中位」比 —— 只对段数 ≥ min_seg 的组合出数。"""
  205. rows = []
  206. idx = {s: i for i, s in enumerate(order)}
  207. for (t, sen, meas), sub in sg.groupby(['turbine', 'sensor_name', 'meas_name']):
  208. m = sub.assign(_i=sub['seg'].map(idx)).dropna(subset=['_i']).sort_values('_i')
  209. if len(m) < min_seg:
  210. continue
  211. vals = m['med'].to_numpy(dtype=float)
  212. early, late = float(np.median(vals[:k])), float(np.median(vals[-k:]))
  213. if not (early > 0):
  214. continue
  215. rows.append(dict(turbine=t, sensor=sen, component=COMP_OF_SENSOR.get(sen, ''), measure=meas,
  216. early=_r6(early), late=_r6(late), ratio=_r3(late / early),
  217. n_seg=int(len(m)), segs=list(m['seg'])))
  218. return pd.DataFrame(rows)
  219. # ── 工单台账 ───────────────────────────────────────────────────────────────────
  220. def load_workorders(store):
  221. """工单台账 → (逐行事件表, 台账范围说明)。台账不存在 ⇒ 空表 + 如实说明 (不造事件)。"""
  222. p = pathlib.Path(store) / 'workorders.parquet'
  223. if not p.exists():
  224. return pd.DataFrame(columns=['turbine', 'date', 'title', 'replaced', 'comps', 'raw']), \
  225. dict(path=P.rel(p), exists=False, rows=0, note='台账不在位 ⇒ 全部 events 为空数组, 不是"无工单"'),
  226. w = pd.read_parquet(p)
  227. w['t_report'] = pd.to_datetime(w.get('t_report'), errors='coerce')
  228. txt = (w.get('故障名称', pd.Series('', index=w.index)).fillna('') + '|'
  229. + w.get('故障描述', pd.Series('', index=w.index)).fillna('') + '|'
  230. + w.get('故障位置二级', pd.Series('', index=w.index)).fillna('') + '|'
  231. + w.get('故障位置三级', pd.Series('', index=w.index)).fillna('') + '|'
  232. + w.get('排查项目', pd.Series('', index=w.index)).fillna('') + '|'
  233. + w.get('故障原因', pd.Series('', index=w.index)).fillna('') + '|'
  234. + w.get('维修对象', pd.Series('', index=w.index)).fillna('') + '|'
  235. + w.get('元器件名称', pd.Series('', index=w.index)).fillna(''))
  236. act = w.get('维修动作', pd.Series('', index=w.index)).fillna('').astype(str)
  237. obj = (w.get('故障名称', pd.Series('', index=w.index)).fillna('') + '|'
  238. + w.get('维修对象', pd.Series('', index=w.index)).fillna('') + '|'
  239. + w.get('元器件名称', pd.Series('', index=w.index)).fillna(''))
  240. recs = []
  241. for i, s in txt.items():
  242. comps = [c for c, keys in WO_COMP_KEYS if any(k in s for k in keys)]
  243. if not comps:
  244. continue
  245. o = str(obj.iloc[i] if hasattr(obj, 'iloc') else '')
  246. repl = bool(REPLACE_PAT.search(act.iloc[i])) and bool(PART_PAT.search(o)) and not CONSUMABLE_PAT.search(o)
  247. title = str(w.get('故障名称', pd.Series('', index=w.index)).iloc[i] or '').strip()
  248. recs.append(dict(turbine=str(w['turbine'].iloc[i]), date=str(w['t_report'].iloc[i])[:10],
  249. title=title[:60], replaced=repl, comps=comps,
  250. act=str(act.iloc[i]), obj=o[:60]))
  251. ev = pd.DataFrame(recs)
  252. # ★ 台账里有**逐字重复行** (实测: 2026-01-03 环网柜 V 柜出线电缆接头三相击穿 同一台同一分钟 7 行;
  253. # 换件类里 WTG23 2024-09-12 发电机轴承更换 4 行)。不去重会让 n_replace / n_events 虚高。
  254. n_raw = int(len(ev))
  255. if n_raw:
  256. ev['_k'] = ev['turbine'] + '|' + ev['date'] + '|' + ev['title'] + '|' + ev['comps'].map(lambda L: ','.join(L))
  257. ev = ev.drop_duplicates(subset=['_k']).drop(columns=['_k']).reset_index(drop=True)
  258. tv = w['t_report'].dropna()
  259. info = dict(path=P.rel(p), exists=True, rows=int(len(w)),
  260. span=[str(tv.min())[:10], str(tv.max())[:10]] if len(tv) else ['', ''],
  261. n_transmission=int(len(ev)), n_transmission_raw=n_raw,
  262. n_replace=int(ev['replaced'].sum()) if len(ev) else 0,
  263. replace_span=[str(pd.to_datetime(ev[ev.replaced].date).min())[:10],
  264. str(pd.to_datetime(ev[ev.replaced].date).max())[:10]] if len(ev) and ev.replaced.any() else None)
  265. return ev, info
  266. def events_of(ev, turbine, comp):
  267. """该 (台, 部件) 的工单事件 (按台/部件关键字映射) → (最近 EVENT_CAP 条明细, 总条数, 其中换件条数)。
  268. ★ 总条数与换件条数在**全量**匹配上数 (明细只是最近的 EVENT_CAP 条): 早期版本用明细数当总数,
  269. 一旦该台该部件工单超过 8 条, n_replace 就会**少报**。
  270. """
  271. if not len(ev) or 'comps' not in ev.columns:
  272. return [], 0, 0
  273. m = ev[ev['comps'].map(lambda L: comp in L) & (ev['turbine'] == turbine)].sort_values('date')
  274. out = [dict(turbine=turbine, comp=comp, date=r['date'], src='工单', title=r['title'], replaced=bool(r['replaced']))
  275. for _, r in m.tail(EVENT_CAP).iterrows()]
  276. return out, int(len(m)), int(m['replaced'].sum())
  277. def build(cfg, farm_name='rudong', dry=False):
  278. m5 = P.m5(farm_name)
  279. store = P.store(farm_name)
  280. t0 = time.time()
  281. meas = MEAS_SERIES
  282. d, used = load_scalars(cfg, meas)
  283. if not len(d):
  284. print('[X] 一个 CMS 窗索引都没读到 —— 先放数据并跑重算链 (scripts/rebuild_all.py); 本器不写空产物。')
  285. return None, 4
  286. kind, seg_date, order, seg_end = mark_segments(d, len(used), {w: v['span_min'] for w, v in used.items()})
  287. sg = seg_table(d)
  288. ev, wo_info = load_workorders(store)
  289. # 台名单以场配置为准 (38 台), 数据里没有的台如实留空 series (不静默丢台)
  290. turbines = [t for t in ws_farm(farm_name).get('turbines', [])] or sorted(d['turbine'].unique())
  291. no_data = [t for t in turbines if t not in set(d['turbine'].unique())]
  292. # 数据覆盖区间取**全表** (src.windcms.data.window_spans, 与页面同源): 只筛 6 标量 × 8 测点会把窗起点说晚一天
  293. span = [min(v['span_min'] for v in used.values()), max(v['span_max'] for v in used.values())]
  294. print(f'部件历史生成端 · 场={farm_name}')
  295. print(f' 窗 {len(used)} 个: ' + '; '.join(
  296. f'{w}(全表 {v["span_min"]}→{v["span_max"]}, {v["rows_all"]}行; 标量 {v["rows_scalar"]}行)'
  297. for w, v in used.items()))
  298. print(f' 标量行 (ds_size==1 ∧ {len(meas)} 标量 ∧ {len(SENSORS)} 测点): {len(d)}'
  299. + (f' · 场配置 {len(turbines)} 台, 无数据的台: {no_data}' if no_data else f' · 台数 {len(turbines)}'))
  300. print(f' 分段口径: segments_kind={kind} · 段 {len(order)} 个: {order[0]}…{order[-1]} '
  301. f'(段日起 {seg_date[order[0]]}→{seg_date[order[-1]]}; 数据覆盖 {span[0]}→{span[1]})')
  302. print(f' 段内样本 (台×测点×标量×段) 中位: {int(sg["n"].median())} · 最少 {int(sg["n"].min())} · 最多 {int(sg["n"].max())}')
  303. print(f' 工单台账: {wo_info["path"]} 行={wo_info["rows"]} 跨度={wo_info.get("span")} '
  304. f'传动链件={wo_info.get("n_transmission")} 换件={wo_info.get("n_replace")} '
  305. f'换件跨度={wo_info.get("replace_span")}')
  306. if len(ev):
  307. print(' 换件类工单 (全量, 判据 candidate): ' + '; '.join(
  308. f'{r["turbine"]} {"/".join(r["comps"])} {r["date"]} {r["title"][:18]}'
  309. for _, r in ev[ev['replaced']].sort_values('date').iterrows()))
  310. # ── trend: 单窗口径的末三/早三 (判据) 与 周段的末三/早三 (补充, 不相交) ──
  311. if kind == 'window':
  312. tw = trend_table(sg, order, k=3, min_seg=MIN_SEG_WINDOW)
  313. else:
  314. tw = trend_table(sg, order, k=3, min_seg=MIN_SEG_WEEK)
  315. tw = tw.sort_values('ratio', ascending=False) if len(tw) else tw
  316. wo_ok = bool(wo_info.get('exists') and wo_info.get('span') and wo_info['span'][0]
  317. and wo_info['span'][0] <= span[1] and wo_info['span'][1] >= span[0])
  318. # ── components ──
  319. comps_out = []
  320. for t in turbines:
  321. for comp in ('主轴承', '齿轮箱', '发电机'):
  322. sens = [s for s, c in COMP_OF_SENSOR.items() if c == comp]
  323. series = {}
  324. n_seg_c = 0
  325. for sen in sens:
  326. for ms in meas:
  327. sub = sg[(sg.turbine == t) & (sg.sensor_name == sen) & (sg.meas_name == ms)]
  328. if not len(sub):
  329. continue
  330. m = sub.set_index('seg')['med'].reindex(order).dropna()
  331. if not len(m):
  332. continue
  333. base = float(np.median(m.to_numpy(dtype=float)))
  334. if not (base > 0):
  335. continue
  336. entries = []
  337. for s in order:
  338. if s not in m.index:
  339. continue
  340. v = float(m.loc[s])
  341. n = int(sub.loc[sub.seg == s, 'n'].iloc[0])
  342. entries.append(dict(seg=s, med=_r6(v), x_self=_r3(v / base), n=n))
  343. series[f'{sen}|{ms}'] = entries
  344. n_seg_c = max(n_seg_c, len(entries))
  345. e_list, n_ev, n_repl = events_of(ev, t, comp)
  346. comps_out.append(dict(turbine=t, component=comp, n_windows=int(n_seg_c),
  347. n_events=n_ev, n_replace=n_repl,
  348. events_capped=bool(n_ev > len(e_list)),
  349. series=series, events=e_list))
  350. # ── summary ──
  351. # ★判据一: 只认多窗口径 (本机必然空) —— 不降门槛凑数
  352. win_trend = trend_table(sg, order, k=3, min_seg=MIN_SEG_WINDOW) if kind == 'window' else pd.DataFrame()
  353. rising = []
  354. if kind == 'window' and len(win_trend):
  355. for (t, comp), sub in win_trend.groupby(['turbine', 'component']):
  356. w = sub.sort_values('ratio', ascending=False).iloc[0]
  357. if w['ratio'] >= RATIO_HI:
  358. _, n_ev, _ = events_of(ev, t, comp)
  359. rising.append(dict(turbine=t, component=comp, ratio=float(w['ratio']), early=float(w['early']),
  360. late=float(w['late']), n_events=n_ev, 工单可对=wo_ok,
  361. sensor=w['sensor'], measure=w['measure'], n_windows=int(w['n_seg'])))
  362. rising = sorted(rising, key=lambda r: -r['ratio'])
  363. # 补充诊断 (非★, 消费者忽略): 周段末三/早三 —— 逐行带段内最小样本数, 让"薄段噪声"可见
  364. diag = []
  365. if len(tw):
  366. nmin = sg.groupby(['turbine', 'sensor_name', 'meas_name'])['n'].min()
  367. for _, r in tw[tw['ratio'] >= RATIO_HI].iterrows():
  368. _, n_ev, _ = events_of(ev, r['turbine'], r['component'])
  369. diag.append(dict(turbine=r['turbine'], component=r['component'], sensor=r['sensor'],
  370. measure=r['measure'], early=float(r['early']), late=float(r['late']),
  371. ratio=float(r['ratio']), n_seg=int(r['n_seg']),
  372. seg_n_min=int(nmin.get((r['turbine'], r['sensor'], r['measure']), 0)),
  373. n_events=n_ev, 工单可对=wo_ok))
  374. diag_all = len(diag)
  375. diag = diag[:20]
  376. # ★判据二: 换件闭环 —— 需要"换件日之前有段 ∧ 之后有段"
  377. loop, loop_skip, repl_pairs = [], [], []
  378. repl = ev[ev['replaced']] if len(ev) and 'replaced' in ev.columns else ev
  379. fusion = {}
  380. fp = m5 / 'fusion_38.csv'
  381. if fp.exists(): # 检出候选旗标 (融合/CMS 面), 只作 verdict 的佐证
  382. fu = pd.read_csv(fp)
  383. for _, r in fu.iterrows():
  384. fusion[str(r.get('台'))] = dict(融合=str(r.get('融合') or ''), CMS=str(r.get('CMS') or ''))
  385. for _, r in (repl.iterrows() if len(repl) else []):
  386. for comp in r['comps']:
  387. repl_pairs.append((r['turbine'], comp, r['date']))
  388. pre_segs = [s for s in order if seg_end[s] < r['date']] # 段**结束**在工单日之前
  389. post_segs = [s for s in order if seg_date[s] > r['date']] # 段**开始**在工单日之后
  390. if not pre_segs or not post_segs:
  391. loop_skip.append(dict(turbine=r['turbine'], component=comp, date=r['date'], title=r['title'],
  392. why=f'段跨度 {span[0]}~{span[1]} 内'
  393. + ('无换件前段' if not pre_segs else '无换件后段')))
  394. continue
  395. sens = [s for s, c in COMP_OF_SENSOR.items() if c == comp]
  396. cand = []
  397. for sen in sens:
  398. for ms in meas:
  399. sub = sg[(sg.turbine == r['turbine']) & (sg.sensor_name == sen) & (sg.meas_name == ms)]
  400. if not len(sub):
  401. continue
  402. a = sub[sub.seg.isin(pre_segs)]['med']
  403. b = sub[sub.seg.isin(post_segs)]['med']
  404. if len(a) and len(b) and float(np.median(a)) > 0:
  405. cand.append((ms, sen, float(np.median(a)), float(np.median(b))))
  406. if not cand:
  407. loop_skip.append(dict(turbine=r['turbine'], component=comp, date=r['date'], title=r['title'],
  408. why='换件前后段内该部件无可用标量'))
  409. continue
  410. ms, sen, pre, post = max(cand, key=lambda c: c[2]) # 取换前幅值最大者 (最不利/最有分辨力)
  411. ratio = post / pre
  412. f = fusion.get(r['turbine']) or {}
  413. hit = bool(f) and (f.get('融合') not in ('', '正常') or '红' in f.get('CMS', '') and '红0' not in f.get('CMS', ''))
  414. if ratio <= RECOVER_HI:
  415. verdict = '恢复' if hit else '回落 (非检出台, 不作端到端实证)'
  416. elif ratio < FLAT_LO:
  417. verdict = '部分回落'
  418. elif ratio <= FLAT_HI:
  419. verdict = '持平'
  420. else:
  421. verdict = '未降/上升'
  422. loop.append(dict(turbine=r['turbine'], component=comp, ratio=_r3(ratio), pre=_r6(pre), post=_r6(post),
  423. replace_date=r['date'], title=r['title'], verdict=verdict,
  424. sensor=sen, measure=ms, 检出候选=hit,
  425. 融合=(f.get('融合') or None), pre_segs=pre_segs, post_segs=post_segs))
  426. loop = sorted(loop, key=lambda r: r['ratio'])
  427. # 负面对照: 工单后振动未降/持平 (非换件的传动链工单, 换前后各有段)
  428. # 不可算的原因**分桶计数**: 单看一个 1512 的总数会把"工单在窗外"和"窗内但无前后段"混成一团。
  429. neg_rows, neg_noc = [], dict(工单在段跨度外=0, 窗内但无前段或后段=0, 该部件无标量=0)
  430. nonrepl = ev[~ev['replaced']] if len(ev) and 'replaced' in ev.columns else ev
  431. for _, r in (nonrepl.iterrows() if len(nonrepl) else []):
  432. for comp in r['comps']:
  433. # ★ 前/后段都**不含工单所在那一段**: 段界线按"段结束 < 工单日"与"段开始 > 工单日"判,
  434. # 否则工单落在段中间时, 该段自己会被当成"换前", 把换后的数据混进 pre (实测污染过 3 条)。
  435. pre_segs = [s for s in order if seg_end[s] < r['date']]
  436. post_segs = [s for s in order if seg_date[s] > r['date']]
  437. if r['date'] < seg_date[order[0]] or r['date'] > seg_end[order[-1]]:
  438. neg_noc['工单在段跨度外'] += 1
  439. continue
  440. if not pre_segs or not post_segs:
  441. neg_noc['窗内但无前段或后段'] += 1
  442. continue
  443. sens = [s for s, c in COMP_OF_SENSOR.items() if c == comp]
  444. cand = []
  445. for sen in sens:
  446. for ms in meas:
  447. sub = sg[(sg.turbine == r['turbine']) & (sg.sensor_name == sen) & (sg.meas_name == ms)]
  448. if not len(sub):
  449. continue
  450. a = sub[sub.seg.isin(pre_segs)]['med']
  451. b = sub[sub.seg.isin(post_segs)]['med']
  452. if len(a) and len(b) and float(np.median(a)) > 0:
  453. cand.append((ms, sen, float(np.median(a)), float(np.median(b))))
  454. if not cand:
  455. neg_noc['该部件无标量'] += 1
  456. continue
  457. ms, sen, pre, post = max(cand, key=lambda c: c[2])
  458. neg_rows.append(dict(turbine=r['turbine'], component=comp, date=r['date'], title=r['title'],
  459. sensor=sen, measure=ms, pre=_r6(pre), post=_r6(post), ratio=_r3(post / pre)))
  460. bucket = dict(下降=sum(1 for r in neg_rows if r['ratio'] < FLAT_LO),
  461. 持平=sum(1 for r in neg_rows if FLAT_LO <= r['ratio'] <= FLAT_HI),
  462. 上升=sum(1 for r in neg_rows if r['ratio'] > FLAT_HI))
  463. neg_noc_total = sum(neg_noc.values())
  464. n_comp_rows = sum(1 for c in comps_out for _ in c['series'])
  465. win_txt = '; '.join(f'{w} {v["span_min"]}~{v["span_max"]}' for w, v in used.items())
  466. n_pair_inspan = sum(1 for _, _, dt in repl_pairs if span[0] <= dt <= span[1])
  467. 判读 = (
  468. f'分段口径 segments_kind={kind}: 本机只解析出 {len(used)} 个导出色窗 ({win_txt}) '
  469. f'⇒ 以窗内 ISO 周段为分段 ({len(order)} 段 {order[0]}…{order[-1]}, {span[0]}~{span[1]}); '
  470. f'多窗 (≥3) 时自动改回多窗口径。'
  471. f'★当前在升 (末三窗/早三窗 ≥{RATIO_HI}) 本机{("有 " + str(len(rising)) + " 行") if rising else "不可算 ⇒ 空表"}: '
  472. f'该判据要求多窗口径且 ≥{MIN_SEG_WINDOW} 窗, 本机 {len(used)} 窗。'
  473. f'★换件闭环 本机{("有 " + str(len(loop)) + " 行") if loop else "不可算 ⇒ 空表"}: 要求换件工单在段跨度内'
  474. f'且其前后各有至少一段; 台账 {wo_info.get("span")} 共 {wo_info.get("n_transmission")} 条传动链工单, '
  475. f'其中换件类 {wo_info.get("n_replace")} 条 (跨度 {wo_info.get("replace_span")}) 展开为 {len(repl_pairs)} 个 '
  476. f'(台,部件) 组合, 段跨度内 {n_pair_inspan} 个 ⇒ 成闭环 {len(loop)} 行, 其余 {len(loop_skip)} 个因'
  477. f'"无换件前段/后段"逐条记入 meta.未成闭环。'
  478. f'补充口径 (非判据): 周段末三/早三比值 ≥{RATIO_HI} 的候选 {diag_all} 条 (列出前 {len(diag)} 条); '
  479. f'周段每 (台,测点,标量) 只有 {int(sg["n"].min())}~{int(sg["n"].max())} 条标量记录/段 '
  480. f'(中位 {int(sg["n"].median())}) ⇒ 该比值只说明窗内走向, 不作在升判据 (明细逐行给出 seg_n_min)。'
  481. f'负面对照: 传动链非换件工单里换前后各有段的 {len(neg_rows)} 条可算 (下降 {bucket["下降"]} · 持平 {bucket["持平"]} '
  482. f'· 上升 {bucket["上升"]}), 其余 {neg_noc_total} 条不可算 ({neg_noc})。')
  483. seg_kind_txt = '导出窗' if kind == 'window' else f'窗内 ISO 周 (本机窗数 {len(used)} < 3 ⇒ 降级; 多窗时自动改回)'
  484. meta = dict(
  485. built=time.strftime('%Y-%m-%d'), farm=farm_name, by='scripts/component_history_build.py',
  486. windows={w: v['span_min'] for w, v in used.items()},
  487. windows_detail={w: dict(rows_all=v['rows_all'], rows_scalar=v['rows_scalar'],
  488. span_min=v['span_min'], span_max=v['span_max'],
  489. seg_min=v['seg_min'], seg_max=v['seg_max'], turbines=v['turbines'])
  490. for w, v in used.items()},
  491. segments_kind=kind,
  492. segments=seg_date,
  493. n_segments=len(order),
  494. 来源声明=('本文件由 scripts/component_history_build.py **现算**写成: 输入 = '
  495. 'outputs/<场>/m5_cms_tcm/windows/*/index.parquet (data/raw 重算产物, 只读 ds_size==1 标量行) '
  496. '+ outputs/<场>/windscada/workorders.parquet (工单台账) '
  497. '+ outputs/<场>/m5_cms_tcm/fusion_38.csv (仅作"检出候选"佐证)。'
  498. '不含随包旧件 (人工裁决/厂家锚) 的数值搬运; 算不出来的 (在升/闭环) 一律空表 + 边界说明, 不填估值。'
  499. 'components[].n_windows = 该部件的**段数** (segments_kind="week" 时即窗内周段数, 不是导出窗数)。'),
  500. 口径=(f'标量: CMS 窗索引 ds_size==1 行, 标量集 = KEY_SCALARS 去 CrestFactor ({", ".join(meas)}); '
  501. f'测点 → 部件: Main_bearing_*→主轴承, Gear_*→齿轮箱, Generator_*→发电机 (System Monitor 非部件测点, 排除)。'
  502. f'段: {seg_kind_txt}。'
  503. f'series 值 = 该段该 (测点|标量) 标量**中位**; 归一 x_self = 段中位 / 该台该测点该标量**全段中位数** '
  504. f'(随包旧件同口径, 旧件写作 x_self; 本件不出 x_fleet —— 那需要跨台同工况对照)。'
  505. f'在升判据 = 末三窗中位 / 早三窗中位 ≥ {RATIO_HI}, 需多窗口径且 ≥{MIN_SEG_WINDOW} 窗。'
  506. f'闭环判据 = 换件工单日之前有段 ∧ 之后有段; pre/post = 该部件换前/后段中位 (取换前幅值最大的测点|标量), '
  507. f'verdict: ratio ≤ {RECOVER_HI} 且该台为检出候选 ⇒ 恢复, 否则按 部分回落/持平/未降 如实记。'
  508. f'工单口径: 部件由关键词映射 (齿轮箱/齿圈/行星/高速轴/中间轴/齿轮油 → 齿轮箱; 主轴承/主轴 → 主轴承; '
  509. f'发电机 → 发电机); replaced = 维修动作含 更换|替换|换件 ∧ 目标文本含部件词 '
  510. f'(轴承|齿圈|行星|高速轴|中间轴|箱体|发电机本体) ∧ 目标文本无耗材词 '
  511. f'(油脂|润滑|油位|油泵|滤芯|密封|加注|油管) —— 换油/换传感器不算换件。'),
  512. inputs=dict(scalars=[v['index'] for v in used.values()], workorders=wo_info.get('path'),
  513. fusion=P.rel(fp) if fp.exists() else None),
  514. 未成闭环=loop_skip[:40], 未成闭环总数=len(loop_skip),
  515. 台账=wo_info,
  516. )
  517. meta['★数据边界'] = (
  518. f'① 本机只有 {len(used)} 个 CMS 导出色窗 ({win_txt}) ⇒ 在升 (末三窗/早三窗) 与 多窗口径的健康期分段'
  519. f'**都不可算** ⇒ summary["★当前在升…"] 为空表; 本件以窗内周段 ({len(order)} 段) 作补充诊断, '
  520. f'只描述窗内走向, 不作判据。'
  521. f'② 换件闭环需要"换件日前后各有段": 工单台账 {wo_info.get("path")} 跨度 {wo_info.get("span")} '
  522. f'(传动链件 {wo_info.get("n_transmission")} 条, 其中换件类 {wo_info.get("n_replace")} 条, '
  523. f'跨度 {wo_info.get("replace_span")}), 与段跨度 {seg_date[order[0]]}~{seg_date[order[-1]]} **无交叠** '
  524. f'⇒ summary["★换件闭环…"] 为空表。★ 这两张空表是数据边界, 不是"没有在升机组/没有闭环案例"; '
  525. f'站点需 (a) 补导 ≥3 个 CMS 色窗 (尤其换件日前后各一窗), (b) 补 2026 年传动链换件工单, '
  526. f'二者到位后本脚本自动出数。'
  527. f'③ 段内样本薄: 周段每 (台,测点,标量) 仅 {int(sg["n"].min())}~{int(sg["n"].max())} 条标量记录 '
  528. f'(中位 {int(sg["n"].median())}), 中位数据此散布 —— series 的 n 逐条给出, 读者自行判断可信度。')
  529. meta['★换件断点'] = '换件前后不可混算趋势; 本表在 events 里标 replaced=true, 读者须自行分段读 series'
  530. payload = dict(
  531. meta=meta,
  532. components=_py(comps_out),
  533. summary=_py({
  534. '★换件闭环 (检出→换件→恢复, 算法端到端验证)': loop,
  535. '★当前在升 (末三窗/早三窗 ≥1.3)': rising,
  536. '窗内周段趋势 (单窗补充口径, 仅参考不作判据)': diag,
  537. '窗内周段趋势·候选总数': diag_all,
  538. '工单后振动未降/持平 (多为油位/传感器/加注类, 非轴承更换 ⇒ 振动本不该变)': dict(
  539. 可算=len(neg_rows), 不可算=neg_noc_total, 不可算原因=neg_noc, 分桶=bucket,
  540. 口径='逐条传动链非换件工单: 换件日前后各有至少一段 ⇒ 可算; 指标取该部件换前幅值最大的 (测点|标量) 的 '
  541. f'post/pre 比值; 分桶 下降<{FLAT_LO} · 持平 {FLAT_LO}~{FLAT_HI} · 上升>{FLAT_HI}',
  542. 明细=sorted(neg_rows, key=lambda r: -r['ratio'])[:12]),
  543. '判读': 判读}),
  544. )
  545. txt = json.dumps(payload, ensure_ascii=False, indent=1, allow_nan=False, default=str)
  546. out = m5 / 'component_history.json'
  547. print(f'\n 段表行 {len(sg)} · 部件行 {len(comps_out)} · series 键 {n_comp_rows} · '
  548. f'events {sum(len(c["events"]) for c in comps_out)} 条')
  549. print(f' summary: ★换件闭环 {len(loop)} 行 · ★当前在升 {len(rising)} 行 · '
  550. f'周段候选 {diag_all} 条(列 {len(diag)}) · 负面对照可算 {len(neg_rows)}/不可算 {neg_noc_total} {neg_noc}')
  551. if neg_rows:
  552. print(' 负面对照明细: ' + '; '.join(f'{r["turbine"]} {r["component"]} {r["date"]} {r["title"][:16]} '
  553. f'{r["sensor"]}|{r["measure"]} {r["ratio"]}×' for r in
  554. sorted(neg_rows, key=lambda r: -r['ratio'])))
  555. if diag:
  556. print(' 周段候选 top: ' + '; '.join(f'{r["turbine"]} {r["component"]} {r["sensor"]}|{r["measure"]} '
  557. f'{r["ratio"]}× ({r["early"]}→{r["late"]}, {r["n_seg"]}段)' for r in diag[:5]))
  558. if loop_skip:
  559. print(f' 换件工单未成闭环 {len(loop_skip)} 条, 例: '
  560. + '; '.join(f'{s["turbine"]} {s["component"]} {s["date"]} {s["why"]}' for s in loop_skip[:3]))
  561. print(f' 产物体积预估 {len(txt.encode("utf-8")) / 1024:.0f} KB ({P.rel(out)}) · 用时 {time.time() - t0:.1f}s')
  562. if dry:
  563. print('\n[dry-run] 不写盘、不登记。')
  564. return payload, 0
  565. m5.mkdir(parents=True, exist_ok=True)
  566. out.write_text(txt, encoding='utf-8')
  567. 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 序列 + '
  568. '工单台账事件; 在升/闭环判据与空表情形见产物 meta.★数据边界)'},
  569. by='scripts/component_history_build.py')
  570. print(f'\n[OK] 已写 {P.rel(out)} ({out.stat().st_size / 1024:.0f} KB) · 已自登记')
  571. return payload, 0
  572. def main() -> int:
  573. ap = argparse.ArgumentParser(description='部件历史生成端 (m5_cms_tcm/component_history.json)')
  574. ap.add_argument('--farm', default='rudong')
  575. ap.add_argument('--dry-run', action='store_true', help='只算不写盘 (打印漏斗与体积预估)')
  576. args = ap.parse_args()
  577. cfg = cms_farm(args.farm)
  578. _, rc = build(cfg, args.farm, dry=args.dry_run)
  579. return rc
  580. if __name__ == '__main__':
  581. sys.exit(main())