fleet_views.py 34 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591
  1. # -*- coding: utf-8 -*-
  2. r"""整场取数视图(P12):`fleet_view` 及其**传递闭包**内的全部本地函数与模块级全局,机械搬到算法服务。
  3. ★来源:`app_backEnd/app_backEnd_guanlan/serve.py`(原样搬运,未改写逻辑)。
  4. 搬运范围由 AST 计算(种子 `fleet_view` → 递归展开被调本地函数 → 带上被引用的模块级全局),避免漏项。
  5. """
  6. from __future__ import annotations
  7. import collections
  8. import datetime as _dt
  9. import json
  10. import pathlib
  11. import re
  12. import threading
  13. import time
  14. from typing import Any
  15. from app_common.app_common_guanlan.api import paths as _P
  16. from src.windscada.config import farm
  17. from src import paths as _P # 路径唯一真源 (与 cwd 无关)
  18. from src.windscada import i18n
  19. from src.windscada import taxonomy
  20. from src.windscada.subsys import temp_nbm, hydraulic, yaw as yawmod, pitch as pitchmod
  21. import hashlib, time
  22. import numpy as np, pandas as pd
  23. import sys, json, pathlib, re, threading, urllib.parse, collections
  24. import threading as _thr
  25. # ── 原 serve.py 的模块级全局(逐字搬) ──
  26. CFG = farm(); ST = pathlib.Path(CFG['store'])
  27. LOCK = threading.Lock()
  28. CFG = farm(); ST = pathlib.Path(CFG['store'])
  29. TS = {'报警': 0, '危险': 0, '良好': 1, '不可判': 2, '优秀': 3}
  30. _CACHE = {}
  31. _HEAVY_LOCK = _thr.Lock()
  32. _STAMP = {'val': None, 'at': 0.0}
  33. _STAMP_TTL = 2.0 # 秒: 指纹有效期 (stat ~150 个文件 ≈ 1~3 ms, 不值得每请求都做)
  34. _WIN_BUSY: dict = {'sysmx': set(), 'curves': set(), 'm9': set()}
  35. _WIN_CACHE: dict = {'sysmx': {}, 'curves': {}, 'm9': {}}
  36. _WIN_ERR: dict = {'sysmx': {}, 'curves': {}, 'm9': {}}
  37. _WIN_LOCK = _thr.Lock()
  38. def _product_files():
  39. """页面取数依赖的产物文件清单 (顺序稳定: 供指纹; 只列**产物**, 不含 reference/ 随包契约)。
  40. 覆盖 _load() 直读的件 + 它经 taxonomy/temp_nbm/hydraulic/fusion 间接读的件:
  41. windscada/*.parquet|csv|json · ontology/*.json · pitch/*.parquet
  42. m5_cms_tcm/{handoff_vibration_v2,component_history,baseline_38}.json · windcms/报告_CMS*.md
  43. """
  44. try:
  45. from src.windcms.config import cms_out as _cms_out # 与写侧同一解析口 (WINDCMS_OUT)
  46. _cms_dir = _cms_out()
  47. except Exception:
  48. _cms_dir = _P.cms()
  49. pats = ((ST, ('*.parquet', '*.csv', '*.json')),
  50. (_P.ont(), ('*.json',)),
  51. (_P.pitch(), ('*.parquet',)),
  52. (_P.m5(), ('handoff_vibration_v2.json', 'component_history.json', 'baseline_38.json')),
  53. (_cms_dir, ('报告_CMS振动状态评估报告_*.md',)))
  54. files = []
  55. for d, ps in pats:
  56. for pat in ps:
  57. files.extend(sorted(d.glob(pat)))
  58. return files
  59. def products_stamp(force=False):
  60. """产物指纹 (sha1 前 16 位)。force=True 时忽略 TTL 立即重算。"""
  61. now = time.time()
  62. if not force and _STAMP['val'] is not None and (now - _STAMP['at']) < _STAMP_TTL:
  63. return _STAMP['val']
  64. h = hashlib.sha1()
  65. for p in _product_files():
  66. try:
  67. st = p.stat()
  68. h.update(f'{_P.rel(p)}|{st.st_mtime_ns}|{st.st_size}\n'.encode('utf-8'))
  69. except OSError:
  70. h.update(f'{_P.rel(p)}|MISSING\n'.encode('utf-8'))
  71. _STAMP.update(val=h.hexdigest()[:16], at=now)
  72. return _STAMP['val']
  73. def _load():
  74. with LOCK:
  75. stamp = products_stamp()
  76. if _CACHE.get('__loaded'):
  77. if _CACHE.get('__stamp') == stamp:
  78. return
  79. # 产物变了 (典型: 刚跑完重算) → 重载。**先建后换**: 下面任何一步抛异常都不会破坏旧缓存,
  80. # 请求照旧能用旧数 (降级但不空白), 同时日志留痕。
  81. print(f'[reload] 产物指纹变化 {_CACHE.get("__stamp")} → {stamp}, 重载', flush=True)
  82. tmp = {}
  83. try:
  84. for key, name in (('tm', 'temp_monthly.parquet'), ('al', 'alarms.parquet'),
  85. ('lm', 'loss_monthly.parquet'), ('bins', 'powercurve_bins.parquet'),
  86. ('pcd', 'powercurve_dev.parquet')):
  87. f = ST / name
  88. if not f.exists():
  89. raise ProductsMissing(name, f)
  90. tmp[key] = pd.read_parquet(f)
  91. tmp['al']['month'] = tmp['al']['t_on'].dt.to_period('M').astype(str)
  92. tmp['pcd'] = tmp['pcd'].set_index('turbine')
  93. for key, name in (('duty', 'duty_monthly.parquet'),):
  94. f = ST / name
  95. tmp[key] = pd.read_parquet(f) if f.exists() else None
  96. try:
  97. tmp['sysmx'] = taxonomy.system_matrix()
  98. tmp['treg'] = temp_nbm.registry()
  99. tmp['treg'] = tmp['treg'][0] if isinstance(tmp['treg'], tuple) else tmp['treg']
  100. tmp['hyd'], _ = hydraulic.registry()
  101. tmp['hyd'] = tmp['hyd'].set_index('turbine')
  102. except FileNotFoundError as e: # 这些派生件同样在产物仓里; 缺了就按"无产物"处理
  103. raise ProductsMissing(getattr(e, 'filename', '派生产物'), getattr(e, 'filename', ST))
  104. zp = _P.pitch() / 'pitch_zero_monthly.parquet'
  105. tmp['zero'] = pd.read_parquet(zp) if zp.exists() else None
  106. tmp['__stamp'] = stamp
  107. tmp['__loaded'] = True
  108. _CACHE.update(tmp) # 原子提交: 失败时不留下半截缓存
  109. except Exception:
  110. if _CACHE.get('__loaded'):
  111. print('[reload] 重载失败 → 继续用上一份缓存 (页面不会空白, 但数是旧的; 看上面的异常)', flush=True)
  112. raise
  113. def _month_end(m):
  114. """'YYYY-MM' → 该月最后一天 'YYYY-MM-DD'。"""
  115. import calendar
  116. y, mo = int(m[:4]), int(m[-2:])
  117. return f'{m}-{calendar.monthrange(y, mo)[1]:02d}'
  118. def months_of(win):
  119. """窗 → 覆盖到的月份列表(月度类数据按此过滤)。"""
  120. _load()
  121. all_m = sorted(_CACHE['tm'].month.unique())
  122. if win == '全程': return all_m
  123. if win == '2025H2': return [m for m in all_m if '2025-07' <= m <= '2025-12']
  124. if win == '2026H1': return [m for m in all_m if '2026-01' <= m <= '2026-06']
  125. if win == '2026年': return [m for m in all_m if m >= '2026-01']
  126. import re as _re
  127. if _re.fullmatch(r'\d{4}-\d{2}', win): # 单月窗 (2026-08-27 用户令: 按月份过滤和选择)
  128. return [m for m in all_m if m == win]
  129. m2 = _re.fullmatch(r'(\d{4}-\d{2})~(\d{4}-\d{2})', win)
  130. if m2: # 月区间窗 "2025-04~2025-10"
  131. return [m for m in all_m if m2.group(1) <= m <= m2.group(2)]
  132. m3 = _re.fullmatch(r'(\d{4}-\d{2}-\d{2})~(\d{4}-\d{2}-\d{2})', win)
  133. if m3: # ★自定义起止(含)→ 相交的月 (用户令 2026-09-21)
  134. a, b = sorted((m3.group(1), m3.group(2)))
  135. return [m for m in all_m if _month_end(m) >= a and f'{m}-01' <= b]
  136. n = 1 if win == '近30日' else 3
  137. return all_m[-n:]
  138. def win_range(win):
  139. """窗 → (起, 止) 日期串,**含两端**(用户令 2026-09-21:时间窗支持自定义起止日期)。
  140. 预设/逐月/月区间窗一律落到"该时间窗覆盖月的月首/月末",于是与 `months_of()` 同口径;
  141. 日区间窗原样返回。日粒度件(停机事件/报警/日粒度判据)用它做**含端**过滤。
  142. """
  143. import re as _re
  144. if win:
  145. if _re.fullmatch(r'\d{4}-\d{2}-\d{2}~\d{4}-\d{2}-\d{2}', win):
  146. a, b = win.split('~')
  147. return (a, b) if a <= b else (b, a)
  148. if _re.fullmatch(r'\d{4}-\d{2}', win):
  149. return (f'{win}-01', _month_end(win))
  150. m2 = _re.fullmatch(r'(\d{4}-\d{2})~(\d{4}-\d{2})', win)
  151. if m2:
  152. return (f'{m2.group(1)}-01', _month_end(m2.group(2)))
  153. if win == '2025H2':
  154. return ('2025-07-01', '2025-12-31')
  155. if win == '2026H1':
  156. return ('2026-01-01', '2026-06-30')
  157. ms = months_of(win)
  158. if not ms:
  159. return ('0001-01-01', '9999-12-31')
  160. return (f'{ms[0]}-01', _month_end(ms[-1]))
  161. def span_of(win):
  162. """窗 → 计算层用的 (起, 止) 元组(含);判级/曲线等按窗重算的口子都吃这个。"""
  163. return win_range(win)
  164. def _mem_mb():
  165. """本进程可用内存(MB);拿不到就返回 None(不因为这些诊断信息把重算搞挂)。"""
  166. try:
  167. import ctypes
  168. class _MS(ctypes.Structure):
  169. _fields_ = [('dwLength', ctypes.c_ulong), ('dwMemoryLoad', ctypes.c_ulong),
  170. ('ullTotalPhys', ctypes.c_ulonglong), ('ullAvailPhys', ctypes.c_ulonglong),
  171. ('ullTotalPageFile', ctypes.c_ulonglong), ('ullAvailPageFile', ctypes.c_ulonglong),
  172. ('ullTotalVirtual', ctypes.c_ulonglong), ('ullAvailVirtual', ctypes.c_ulonglong),
  173. ('ullAvailExtendedVirtual', ctypes.c_ulonglong)]
  174. st = _MS()
  175. st.dwLength = ctypes.sizeof(_MS)
  176. ctypes.windll.kernel32.GlobalMemoryStatusEx(ctypes.byref(st))
  177. return int(st.ullAvailPhys // (1024 * 1024))
  178. except Exception:
  179. return None
  180. def _win_key(win):
  181. a, b = span_of(win)
  182. return f'{a}~{b}'
  183. def _win_get(kind, win, fn):
  184. """→ 值 或 None(None = 正在算/刚起算)。失败**记名**(`_WIN_ERR`)而不是静默当"没数据"。"""
  185. key = _win_key(win)
  186. with _WIN_LOCK:
  187. if key in _WIN_CACHE[kind]:
  188. return _WIN_CACHE[kind][key]
  189. if key not in _WIN_BUSY[kind]:
  190. _WIN_BUSY[kind].add(key)
  191. def _run():
  192. try:
  193. with _HEAVY_LOCK: # 重活串行: 同一时刻只算一份
  194. m0 = _mem_mb()
  195. print(f'[win] {kind} 按窗重算开算 {key}(可用内存 {m0} MB)', flush=True)
  196. v = fn()
  197. print(f'[win] {kind} 按窗重算完成 {key}(可用内存 {_mem_mb()} MB)', flush=True)
  198. except Exception as e: # 守护失败必响亮
  199. _WIN_ERR[kind][key] = f'{type(e).__name__}: {e}'[:200]
  200. print(f'[win] {kind} 按窗重算失败 {key}: {_WIN_ERR[kind][key]}', flush=True)
  201. v = None
  202. with _WIN_LOCK:
  203. if v is not None:
  204. _WIN_CACHE[kind][key] = v
  205. _WIN_BUSY[kind].discard(key)
  206. _thr.Thread(target=_run, name=f'win-{kind}-{key}', daemon=True).start()
  207. print(f'[win] {kind} 按窗重算启动 {key}(首次约数十秒,页面会先出 pending)', flush=True)
  208. return None
  209. def sysmx_of(win):
  210. """→ (判级矩阵, pending)。`pending=True` 时返回的是"另一口径"的旧矩阵,页面必须如实标注。"""
  211. from src.windscada import taxonomy
  212. m = _win_get('sysmx', win, lambda: taxonomy.system_matrix(CFG, span=span_of(win)))
  213. if m is None:
  214. return _CACHE['sysmx'], True
  215. return m, False
  216. def m9_of(win):
  217. """→ (控制参数一致性, pending)。按窗重算(走窄仓 ≈5 s/窗);干净窗直接用正式产物。"""
  218. from src.windscada.perf import control as _cm
  219. a, b = span_of(win)
  220. if (a, b) == (_cm.WIN[0], '2025-12-31'):
  221. return _cm.registry(CFG), False
  222. v = _win_get('m9', win, lambda: _cm.registry(CFG, span=(a, b)))
  223. if v is None:
  224. return _cm.registry(CFG), True # 未就位: 先给上一口径, 页面按 pending 标注
  225. return v, False
  226. def fleet_view(win):
  227. _load()
  228. ms = months_of(win)
  229. # ★2026-09-21 用户令「判级也按所选时间窗重算」: 判级矩阵按窗算(后台+缓存,首次 pending)。
  230. sysmx, sysmx_pending = sysmx_of(win)
  231. # ① 系统分类问题 (全系统)
  232. systems, sysdist = {}, {}
  233. for s in taxonomy.SYSTEMS:
  234. rows = [dict(t=t, st=sysmx[t][s]['状态'], why=i18n.humanize(sysmx[t][s]['依据']))
  235. for t in CFG['turbines'] if sysmx[t][s]['状态'] in ('报警', '不可判')] # 良好不出 (用户令: 只显示有问题的)
  236. rows.sort(key=lambda r: TS.get(r['st'], 9))
  237. systems[s] = rows
  238. # 全场分布 (系统入口卡用: rows 只含问题台, 不能当全场分母 — 会显示"共12台")
  239. from collections import Counter as _C
  240. sysdist[s] = dict(_C(sysmx[t][s]['状态'] for t in CFG['turbines']))
  241. # 良好台清单单列 (系统详情页第三档)
  242. systems[s + '·良好'] = [dict(t=t, st='良好', why=i18n.humanize(sysmx[t][s]['依据']))
  243. for t in CFG['turbines'] if sysmx[t][s]['状态'] == '良好']
  244. # ①b 关注清单: ≥2 系统报警 (逐台独立台页)
  245. watch = []
  246. for t in CFG['turbines']:
  247. al_sys = [x for x in taxonomy.SYSTEMS if sysmx[t][x]['状态'] == '报警']
  248. if len(al_sys) >= 2:
  249. watch.append(dict(t=t, n=len(al_sys), syss=al_sys,
  250. # ★不截断: 截到 70 字会把依据切成半句, 而模板反解要完整形态
  251. # ⇒ 英文侧整条落回中文 (2026-09-03 实逮, 总览 7 卡 4 张中文)。
  252. # 显示长度归前端 CSS 管。
  253. why=';'.join(f"{x}: {i18n.humanize(sysmx[t][x]['依据'])}" for x in al_sys)))
  254. watch.sort(key=lambda r: -r['n'])
  255. n_alarm_t = sum(1 for t in CFG['turbines'] if any(sysmx[t][x]['状态'] == '报警' for x in taxonomy.SYSTEMS))
  256. # ② 故障统计梳理 (真时间窗)
  257. al = _CACHE['al']; a = al[al.month.isin(ms)]
  258. top_txt = a.groupby(['code', 'text']).agg(n=('code', 'size'), dur_h=('dur_s', lambda x: x.sum() / 3600)).reset_index()
  259. pareto_n = top_txt.sort_values('n', ascending=False).head(12)
  260. pareto_d = top_txt.sort_values('dur_h', ascending=False).head(12)
  261. monthly = a.groupby('month').size().reindex(ms).fillna(0)
  262. per_t = a.groupby('turbine').size().sort_values(ascending=False).head(10)
  263. # 经典故障分析: 条数与时长必须配对看 (2026-08-28)。两个独立排行榜分不出
  264. # "高频短时(信号抖动或重复触发)" 与 "低频长时(硬故障)" — 处置方向相反, 混在一起会误派工。
  265. # 第三维=影响台数, 分离"单台刷屏"与"全场批次共性"。
  266. # 残月识别 (2026-08-28 审核逮): 2026-07 只有 6 天数据(覆盖 19.4%), 却被当整月
  267. # 并进"2026年"的占比与月度趋势 → 末柱"下降"是假象。占比类分母也因此不完整。
  268. import calendar as _cal
  269. _lmc = _CACHE['lm']; _lmc = _lmc[_lmc.month.astype(str).isin(ms)]
  270. _cov = {}
  271. for _mm in ms:
  272. _sub = _lmc[_lmc.month.astype(str) == _mm]
  273. if not len(_sub):
  274. _cov[_mm] = 0.0; continue
  275. _y, _mo = int(_mm[:4]), int(_mm[-2:])
  276. _cal_h = _cal.monthrange(_y, _mo)[1] * 24
  277. _cov[_mm] = round(float(_sub['rows_'].sum() / 6 / max(_sub['turbine'].nunique(), 1) / _cal_h), 3)
  278. _nt = a.groupby(['code', 'text'])['turbine'].nunique().rename('nt')
  279. _med = a.groupby(['code', 'text'])['dur_s'].median().rename('med_s')
  280. qd = top_txt.set_index(['code', 'text']).join([_nt, _med]).reset_index()
  281. qd = qd.sort_values('n', ascending=False).head(30)
  282. faults = dict(
  283. pareto_n=[dict(k=f"{r.code} {i18n.alarm_label(r.code, r.text)[:18]}", v=int(r.n)) for _, r in pareto_n.iterrows()],
  284. pareto_d=[dict(k=f"{r.code} {i18n.alarm_label(r.code, r.text)[:18]}", v=round(float(r.dur_h), 1)) for _, r in pareto_d.iterrows()],
  285. monthly=dict(months=ms, vals=[int(v) for v in monthly],
  286. cov=[_cov.get(x, 0.0) for x in ms],
  287. per_day=[round(float(monthly[x]) / max(_cov.get(x, 0.0) * _cal.monthrange(int(x[:4]), int(x[-2:]))[1], 1e-9), 1)
  288. if _cov.get(x, 0) > 0.02 else None for x in ms]),
  289. cov=_cov, cov_min=min(_cov.values()) if _cov else 1.0,
  290. partial=[x for x in ms if _cov.get(x, 1) < 0.5],
  291. quad=[dict(code=str(r.code), k=i18n.alarm_label(r.code, r.text)[:20], n=int(r.n),
  292. h=round(float(r.dur_h), 1), nt=int(r.nt), med=round(float(r.med_s), 1))
  293. for _, r in qd.iterrows()],
  294. n_codes=int(len(top_txt)),
  295. per_t=[dict(k=k, v=int(v)) for k, v in per_t.items()], total=int(len(a)))
  296. # ③ SOP 控制策略
  297. lm = _CACHE['lm']; l = lm[lm.month.astype(str).isin(ms)]
  298. hrs = l.groupby('state')['rows_'].sum() / 6
  299. loss = l.groupby('state')['loss'].sum() / 1000
  300. hrs_tot = max(hrs.sum(), 1)
  301. stop_h = hrs.get('停机', 0) + hrs.get('停机(调度令)', 0)
  302. disp_h = hrs.get('停机(调度令)', 0)
  303. lw_h = hrs.get('低风待机', 0)
  304. # 口径与 availability.summary 单源一致: 1 − 停机/(总 − 调度令 − 低风待机)
  305. avail = 1 - (stop_h - disp_h) / max(hrs_tot - disp_h - lw_h, 1)
  306. pcd = _CACHE['pcd']
  307. # 能量账闭合 (2026-08-28 经典分析): 四行工况表看不出"电量流向"。
  308. # 理论可发 = 实际上网 + Σ各态损失; 过物理上限核 (不可超 38台×4.0MW×窗时长)。
  309. _act = l.groupby('state')['act'].sum() / 1000
  310. _wf = [dict(k=str(k), loss=round(float(loss.get(k, 0)), 0), act=round(float(_act.get(k, 0)), 0),
  311. h=round(float(hrs.get(k, 0)), 0)) for k in hrs.index]
  312. _A = float(_act.sum()); _L = float(loss.sum())
  313. _cap = len(CFG['turbines']) * 4.0 * float(hrs_tot) / max(len(CFG['turbines']), 1)
  314. # 逐月能量 (2026-08-28 报表端口): 汇报要环比与趋势, 全时间窗聚合给不出。
  315. # 残月覆盖率一并带出 — 报表里拿残月和整月比环比会读反 (故障月度图已踩过一次)。
  316. _em = []
  317. for _mm in ms:
  318. _sub = l[l.month.astype(str) == _mm]
  319. if not len(_sub):
  320. continue
  321. _a1 = float(_sub['act'].sum()) / 1000
  322. _l1 = float(_sub['loss'].sum()) / 1000
  323. _em.append(dict(m=_mm, act=round(_a1, 0), loss=round(_l1, 0), theo=round(_a1 + _l1, 0),
  324. loss_pct=round(_l1 / max(_a1 + _l1, 1) * 100, 1),
  325. cov=_cov.get(_mm, 1.0),
  326. eflh=round(_a1 / max(len(CFG['turbines']) * 4.0, 1), 0)))
  327. energy = dict(monthly=_em, act=round(_A, 0), loss=round(_L, 0), theo=round(_A + _L, 0),
  328. loss_pct=round(_L / max(_A + _L, 1) * 100, 1),
  329. eflh=round(_A / max(len(CFG['turbines']) * 4.0, 1), 0),
  330. cap=round(_cap, 0), cap_ok=bool(_A + _L < _cap),
  331. items=sorted(_wf, key=lambda r: -r['loss']))
  332. control = dict(
  333. energy=energy,
  334. states=[dict(k=k, h=round(float(v), 0), pct=round(float(v / hrs_tot * 100), 1), loss=round(float(loss.get(k, 0)), 0)) for k, v in hrs.items()],
  335. avail=round(float(avail * 100), 1),
  336. curve=dict(sigma=round(float(pcd.dev_w.std() * 100), 2),
  337. cands=[dict(t=i, dev=round(float(r.dev_w * 100), 2), 判=r['判别']) for i, r in pcd.iterrows() if r['判别'] != '—'],
  338. note='固定判别窗 2025-H2'),
  339. anchors='限电命令面占比 2025≈0.197 → 2026≈0.769 (专项锚)')
  340. m8 = None
  341. try:
  342. from src.windscada.perf import faults as fmod
  343. if (ST / 'stop_events.parquet').exists():
  344. m8 = dict(mtbf=fmod.mtbf_summary(ms), stop_pareto=fmod.stop_pareto(ms), seasonal=fmod.seasonal())
  345. except Exception as e:
  346. m8 = dict(err=str(e)[:120])
  347. m9 = None
  348. try:
  349. from src.windscada.perf import control as cmod
  350. if (ST / 'control_profile.parquet').exists():
  351. m9, m9_pending = m9_of(win) # ★按窗重算(用户令 2026-09-21)
  352. except Exception as e:
  353. m9 = dict(err=str(e)[:120])
  354. m9_pending = False
  355. try:
  356. from src.windscada.perf import reliability as relmod
  357. rel = relmod.overview(CFG, span=span_of(win))
  358. except Exception as e:
  359. rel = dict(err=str(e)[:120])
  360. try:
  361. from src.windscada.subsys import fusion as fusmod
  362. ftab, fmeta = fusmod.fusion_table(CFG, all_turbines=True)
  363. mrows, mheads = fusmod.matrix(CFG, ms) # 整体时间窗: 矩阵现窗/记忆轨迹随全局窗选择器
  364. cap = fusmod.capability()
  365. ev, ev_bounds = fusmod.drivetrain_events(ms, CFG)
  366. lvcnt = {}
  367. for r in mrows:
  368. lv = r['振动']['level']; lvcnt[lv] = lvcnt.get(lv, 0) + 1
  369. fkpi = dict(定论=lvcnt.get('bad', 0), 预警=lvcnt.get('warn', 0), 候选观察=lvcnt.get('note', 0),
  370. 销案正常=lvcnt.get('ok', 0), 未列=lvcnt.get('unlisted', 0),
  371. 收录台=len(CFG['turbines']) - lvcnt.get('unlisted', 0), 全场=len(CFG['turbines']),
  372. 机制链=sum(1 for r in mrows if '机制链' in r['润滑'].get('note', '')),
  373. 系数=dict(wear=cap['classes'].get('wear_progressive', {}).get('coef'),
  374. thermal=cap['classes'].get('thermal_acute', {}).get('coef'),
  375. unknown=cap['classes'].get('unknown', {}).get('coef')))
  376. # 矩阵补 报警/工单 两列 (2026-08-28 用户令) — 报警=时间窗内真窗计数+族分布; 工单=历史台账(窗见列头)
  377. # 报警/工单两列只统计传动链三部件 (2026-08-28 用户令: 只说发电机/齿轮箱/主轴) —
  378. # 本矩阵每行是"该台传动链状态", 全场报警计数会把变桨/偏航/主控的量混进来 (316 条里绝大多数与传动链无关)
  379. _DT_PAT = '齿轮|齿箱|润滑油|油冷|滤芯|发电机|定子|绕组|滑环|主轴承|主轴|轴承'
  380. _a = _CACHE['al']; _aw = _a[_a.month.isin(ms)]
  381. _aw = _aw[_aw.text.str.contains(_DT_PAT, na=False)]
  382. _acnt = _aw.groupby('turbine').size()
  383. _arank = _acnt.rank(ascending=False, method='min')
  384. _FAM = [('齿轮箱', '齿轮|齿箱|润滑油|油冷|滤芯'), ('发电机', '发电机|定子|绕组|滑环'),
  385. ('主轴承', '主轴承|主轴')]
  386. try:
  387. from src.windscada.subsys import workorder as _wo
  388. _wod = _wo.load(CFG)
  389. except Exception:
  390. _wod = None
  391. for _r in mrows:
  392. _t = _r['turbine']; _g = _aw[_aw.turbine == _t]
  393. _fam, _taken = [], set()
  394. _CODEFAM = {'3225': '变桨叶片'} # 文本为空的码走码号归族 (3225=变桨液压, 现场西门子资料对译表)
  395. for _nm, _pat in _FAM:
  396. _byc = _g.code.astype(str).map(_CODEFAM) == _nm
  397. _sub = _g[~_g.index.isin(_taken) & (_g.text.str.contains(_pat, na=False) | _byc)]
  398. _taken |= set(_sub.index)
  399. if len(_sub): _fam.append(dict(k=_nm, n=int(len(_sub))))
  400. _oth = len(_g) - sum(f['n'] for f in _fam)
  401. if _oth > 0: _fam.append(dict(k='其他', n=int(_oth)))
  402. # 刷屏闸: 单码占比过半 → 条数不代表"问题多"而是一个码在抖 (08# 4609/4918=93.7% 码3225 集中2026-01)
  403. _burst = None
  404. if len(_g):
  405. _bc = _g.groupby('code').size().sort_values(ascending=False)
  406. _share = float(_bc.iloc[0]) / len(_g)
  407. if _share >= 0.5:
  408. _sub2 = _g[_g.code == _bc.index[0]]
  409. _burst = dict(code=str(_bc.index[0]), share=round(_share, 3), n=int(_bc.iloc[0]),
  410. months=sorted({str(m)[:7] for m in _sub2.t_on.dt.to_period('M').astype(str)}),
  411. med_s=float(_sub2.dur_s.median()) if 'dur_s' in _sub2 else None)
  412. # 明细 (用户令"浮窗说明之前发生问题"): top 码 + 首末时间
  413. _top = []
  414. if len(_g):
  415. for (_c2, _tx), _sub3 in _g.groupby(['code', 'text']):
  416. _top.append(dict(code=str(_c2), text=str(_tx)[:34], n=int(len(_sub3)),
  417. first=str(_sub3.t_on.min())[:10], last=str(_sub3.t_on.max())[:10]))
  418. _top.sort(key=lambda x: -x['n'])
  419. _r['报警'] = dict(n=int(_acnt.get(_t, 0)), rank=int(_arank.get(_t, 0)) if _t in _arank else None,
  420. tot=len(CFG['turbines']), fam=_fam, burst=_burst, items=_top[:6])
  421. if _wod is not None and len(_wod):
  422. _w = _wod[_wod.turbine == _t]
  423. _wtxt = (_w.get('故障名称', '').astype(str) + ' ' + _w.get('故障位置二级', '').astype(str)
  424. + ' ' + _w.get('维修对象', '').astype(str) + ' ' + _w.get('元器件名称', '').astype(str))
  425. _w = _w[_wtxt.str.contains(_DT_PAT, na=False)]
  426. _acts = _w['维修动作'].replace('', pd.NA).dropna().value_counts() if len(_w) else None
  427. _wi = []
  428. for _, _wr in _w.sort_values('t_report', ascending=False).head(6).iterrows():
  429. _wi.append(dict(date=str(_wr.get('t_report'))[:10], name=str(_wr.get('故障名称') or '')[:30],
  430. act=str(_wr.get('维修动作') or ''), part=str(_wr.get('维修对象') or _wr.get('元器件名称') or '')[:14]))
  431. _r['工单'] = dict(n=int(len(_w)),
  432. acts=[dict(k=str(k), n=int(v)) for k, v in (_acts.head(3).items() if _acts is not None else [])],
  433. items=_wi,
  434. last=(lambda v: str(v.max())[:10] if len(v) else None)(_w.t_report.dropna()[_w.t_report.dropna().dt.year > 2000]) if len(_w) else None)
  435. else:
  436. _r['工单'] = dict(n=0, acts=[], last=None)
  437. _wyr = ''
  438. if _wod is not None and len(_wod):
  439. _tv = _wod.t_report.dropna()
  440. _tv = _tv[_tv.dt.year > 2000] # 空日期落 1970 epoch, 剔除后再报窗 (否则列头写"1970~")
  441. if len(_tv): _wyr = f"{_tv.min():%Y}~{_tv.max():%Y}"
  442. mheads = dict(mheads, 报警=f"传动链三部件 · {ms[0]}~{ms[-1]}" if ms else '传动链三部件',
  443. 工单=f"传动链三部件 · 台账 {_wyr}" if _wyr else '台账不可用')
  444. # 决策链进度盘 (2026-08-28 用户问"这个到底能怎么用"): 原来只有 29# 一台的静态六格,
  445. # 是展示牌不是工作面 — 打开它做不了任何决定。改为对全部需跟踪台报"走到第几步、卡在哪、下一步做什么"。
  446. # 判定规则单源在此, 前端只渲染不判断。
  447. _CLOOP = {}
  448. try:
  449. _ch = json.loads((ST.parent / 'm5_cms_tcm' / 'component_history.json').read_text(encoding='utf-8'))
  450. for _k, _v in _ch.get('summary', {}).items():
  451. if _k.startswith('★换件闭环') and isinstance(_v, list):
  452. for _r in _v:
  453. _CLOOP[_r['turbine']] = _r
  454. if _k.startswith('★当前在升') and isinstance(_v, list):
  455. for _r in _v:
  456. _CLOOP.setdefault(_r['turbine'], {}).update(rising=_r)
  457. except Exception:
  458. pass
  459. _HIT = {'bad', 'warn', 'note'}
  460. # ★一台可能有多行 (29#/17# 各 2 行: 齿轮箱 + 主轴承)。用 {turbine: row} 直接建字典会被后一行覆盖,
  461. # 实测 29# 因此取到"齿轮箱二级行星内齿圈·监视"那行, 而它的决策链讲的是"主轴承+集中润滑泵·定论"
  462. # ⇒ 浮窗里部件、温度通道、油样全是另一个部件的 (memory: 画图须确认各维度来自同一个体)。
  463. _SEV = {'bad': 0, 'warn': 1, 'note': 2, 'stale': 3, 'ok': 4, 'unlisted': 5}
  464. _tbi = {}
  465. for _row in ftab.to_dict('records'):
  466. _k = _row['turbine']
  467. if _k not in _tbi or _SEV.get(_row.get('级'), 9) < _SEV.get(_tbi[_k].get('级'), 9):
  468. _tbi[_k] = _row
  469. board = []
  470. for _r in mrows:
  471. t = _r['turbine']
  472. srcs = [k for k in ('振动', '温度', '润滑', '油液') if (_r.get(k) or {}).get('level') in _HIT]
  473. lvl = (_r.get('振动') or {}).get('level')
  474. tb = _tbi.get(t, {})
  475. mech = ('机制链' in ((_r.get('润滑') or {}).get('note') or '')
  476. or '机制链' in (tb.get('振动结论') or ''))
  477. n_wo = (_r.get('工单') or {}).get('n') or 0
  478. cl = _CLOOP.get(t) or {}
  479. closed = bool(cl.get('verdict') == '恢复')
  480. rising = cl.get('rising')
  481. if not (lvl in _HIT or (srcs and lvl != 'ok')):
  482. continue
  483. # 六步: done=已达 / open=未达 / unknown=数据不足不可判
  484. steps = [
  485. dict(k='证据', st='done' if srcs else 'open', v=f"{len(srcs)} 源", d='/'.join(srcs)),
  486. dict(k='机制', st='done' if mech else 'open',
  487. v='已定性' if mech else '未定性', d='根因定到部件与失效模式'),
  488. dict(k='判级', st='done' if lvl in _HIT else 'open',
  489. v=tb.get('证据状态') or lvl or '—', d=f"设备状态 {tb.get('设备状态') or '—'}"),
  490. dict(k='排期', st='done' if t == 'WTG29' else 'open',
  491. v='情景A 2026-09' if t == 'WTG29' else '未排',
  492. d='进检修排程沙盘 (风险×损失双轴)'),
  493. # ★这一列一律"不可判", 不是偷懒: 工单台账止 2024-11 而判级时间窗在 2026,
  494. # 台账里的历史工单不可能是针对本次问题的 ⇒ 用 n>0 判"已派工"会造出假的完成态。
  495. dict(k='动作', st='unknown',
  496. v=(f"台账 {n_wo} 单" if n_wo else '台账无'),
  497. d='工单台账 2020~2024-11, 25-26 现场未提供 — 历史单不对应本次问题, 须向现场核实'),
  498. # ★历史闭环不占本格: 05#/20# 2024~2025 换过齿轮箱且已验证恢复, 但本轮又有新的候选级证据,
  499. # 把历史闭环填进第六格会出现"验收已达而机制未定"的自相矛盾行。历史闭环走 hist_loop 单独标记。
  500. dict(k='验收', st='open', v='未验',
  501. d='本轮动作执行后复测同一判据是否回落' + (
  502. f" (该台 {cl.get('replace_date')} 有过一次闭环: {cl.get('component')} 降至 {cl.get('ratio')}×)" if closed else '')),
  503. ]
  504. _first = next((x for x in steps if x['st'] != 'done'), None)
  505. stuck = _first['k'] if _first else None
  506. stuck_kind = _first['st'] if _first else None
  507. NEXT = {'证据': '补测: 四源全为常态, 先确认监测面有效',
  508. '机制': '定性: 现场/取样把根因定到部件与失效模式 — 机制不清则排期与动作都无依据',
  509. '判级': '送审: 按 RV-1 触发相应审级',
  510. '排期': '排程: 进沙盘比情景, 出建议窗口',
  511. '动作': '核实: 向现场调取 2025~2026 工单 — 台账止于 2024-11, 系统无法判定是否已派工',
  512. '验收': '复测: 动作已执行, 复测同一判据是否回落'}
  513. board.append(dict(t=t, lvl=lvl, ostate=tb.get('设备状态'), estate=tb.get('证据状态'),
  514. part=tb.get('部件'), steps=steps, stuck=stuck, stuck_kind=stuck_kind,
  515. next=NEXT.get(stuck, ''), done=sum(1 for x in steps if x['st'] == 'done'),
  516. rising=rising, hist_loop=(cl if closed else None), srcs=srcs))
  517. _ORD = {'bad': 0, 'warn': 1, 'note': 2}
  518. board.sort(key=lambda r: (_ORD.get(r['lvl'], 3), -r['done']))
  519. _cl_rows = [v for v in _CLOOP.values() if v.get('verdict') == '恢复']
  520. fus = dict(链盘=dict(rows=board, closed=_cl_rows,
  521. stuck=dict(collections.Counter(r['stuck'] for r in board if r['stuck'])),
  522. wo_bound='工单台账 2020~2024-11 (现场未提供 25-26)'),
  523. 表=ftab.to_dict('records'),
  524. 矩阵=mrows, 列窗=mheads, kpi=fkpi,
  525. 事件=ev, 事件边界=ev_bounds,
  526. 能力=dict(classes=cap['classes'], blind=[str(b) for b in cap['blind']], loop=cap.get('loop', {})),
  527. 窗=dict(振动=fmeta.get('振动窗', '—'), SCADA=fmeta.get('SCADA窗', ''), 油样=fmeta.get('油样窗', '')),
  528. 色标=fmeta.get('色标', {}), 数据时点=fmeta.get('数据时点', ''), handoff日期=fmeta.get('handoff日期', ''),
  529. 证据窗末=fmeta.get('证据窗末', ''),
  530. 纪律=str(fmeta.get('纪律', '')), 盲区=[str(b) for b in fmeta.get('盲区', [])],
  531. open_items=[str(o) for o in fmeta.get('open_items', [])], gap=fmeta.get('gap'))
  532. except Exception as e:
  533. fus = dict(err=str(e)[:160])
  534. kpi = dict(报警台=n_alarm_t, 关注台=len(watch), 全场=len(CFG['turbines']))
  535. return dict(win=win, months=ms, all_months=sorted(_CACHE['tm'].month.unique()),
  536. systems=systems, sysdist=sysdist, watch=watch, kpi=kpi, rel=rel, fus=fus, faults=faults, control=control, m8=m8, m9=m9,
  537. # ★2026-09-21: 判级轴现在**按所选窗重算**;首次是后台算,这一份仍是旧口径 ⇒ 如实标 pending
  538. sysmx_pending=bool(sysmx_pending),
  539. m9_pending=bool(m9_pending),
  540. win_pending=bool(sysmx_pending or m9_pending),
  541. sysmx_span=list(span_of(win)),
  542. note=('判级轴(变桨/偏航/蓄能/温度)按所选时间窗重算;曲线按所选时间窗重算;'
  543. '故障统计/五态/温度月轨迹/停机台账=所选窗真窗'
  544. + ('(⚠ 判级矩阵正在按所选时间窗重算,下面系统卡暂为上一份口径,稍后自动刷新)' if sysmx_pending else '')))