windscada_monthly_build.py 9.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206
  1. #!/usr/bin/env python3
  2. # -*- coding: utf-8 -*-
  3. """月度派生件构建器 —— 补上 v0.2.0 里"有消费者、却没生成端"的那批产物 (2026-09-12)。
  4. ## 为什么有这个脚本
  5. 从零重算 (`rebuild_from_raw.py`) 只覆盖 16 件; 随包里另有 15 件 parquet **没有任何生成端**
  6. (全库搜遍: 只有读取方、0 处写入方), 于是工作台主数据 `/api/fleet` 一缺 `temp_monthly.parquet`
  7. 就整个回 `err=no_products` —— 页面看着像"没数据", 其实是"缺产物"。
  8. 本脚本把这批件**从 data/raw 重新算出来**, 并且**逐值对齐随包件**: 随包件在这里当"标准答案"
  9. (存在 `_products_off/.../windscada/`), `--verify` 会报每件的命中率; 规则是从"标准答案"反推并
  10. 用全量比对确认的, 不是猜的。
  11. ## 已确认的口径 (2026-09-12, 用随包件反推 + 逐值验证)
  12. temp_monthly 每台 × 每月 × 每个温度通道, 在 **发电态功率门 p ≥ 500kW** 行上的**中位数**
  13. 形状 38×19×27 = 19494 行; 实测 WTG01 的 513 格 **逐值 100% 相同**
  14. (对照: 全行中位数只有 5.8% 命中 —— 因为停机时绕组是冷的, 会把中位拉低)
  15. 口径出处: `src/windscada/subsys/temp_nbm.py` 的"发电态×功率档(500kW)", PBINS 步长 500
  16. ## 用法
  17. python scripts/windscada_monthly_build.py --verify # 只比对(不写盘), 报每件命中率
  18. python scripts/windscada_monthly_build.py # 算并写盘
  19. python scripts/windscada_monthly_build.py --turbines WTG01,WTG02 # 只算几台(调试用)
  20. """
  21. from __future__ import annotations
  22. try:
  23. from app_common.app_common_guanlan.api import install_root as _install_root
  24. except ImportError: # 理论不可达;包结构异常时回退到按位置上跳
  25. from pathlib import Path as _P
  26. def _install_root(_f): return _P(_f).resolve().parents[3]
  27. import argparse
  28. import pathlib
  29. import sys
  30. ROOT = _install_root(__file__) # 模块化后按标记找安装根(原 parents[1] 已不成立)
  31. sys.path.insert(0, str(ROOT))
  32. import pandas as pd # noqa: E402
  33. from src import paths as P # noqa: E402
  34. from src.windscada.config import farm # noqa: E402
  35. from src.windscada.data import load_10min, contracted_cols # noqa: E402
  36. POWER_GATE = 500.0 # 发电态功率门 (kW); 出处见 temp_nbm 的"发电态×功率档(500kW)"
  37. GROUPS = ['A.功率', 'B.温度NBM']
  38. def baseline_dir():
  39. """"标准答案"目录: 随包件在哪 (整目录视角; 单件比对用 baseline_for)。"""
  40. inplace = P.store() / '_pre_rebuild_20260911'
  41. if inplace.is_dir():
  42. return inplace
  43. for p in sorted((ROOT / '_products_off').rglob('windscada')):
  44. if p.is_dir():
  45. return p
  46. return None
  47. def baseline_for(name: str):
  48. """这一件产物的"标准答案"在哪。
  49. 注意两个暂存处**内容不同**, 必须逐件找:
  50. · `outputs/<场>/windscada/_pre_rebuild_20260911/` = **台账验收基线**(3 件: alarms/workorders/oil_samples);
  51. · `_products_off*/**/windscada/` = 随包整份产物(31 件 parquet, 含月度派生件)。
  52. 原先只认第一个目录, 于是 temp_monthly 会被报成"没有随包基线可比", 而结论却打印"全部逐值一致"
  53. —— 2026-09-12 实测逮到(假结论比没结论更糟)。
  54. """
  55. cands = []
  56. bd = baseline_dir()
  57. if bd:
  58. cands.append(bd / name)
  59. for p in sorted((ROOT).glob('_products_off*/**/windscada')):
  60. cands.append(p / name)
  61. for c in cands:
  62. if c.is_file():
  63. return c
  64. return None
  65. BASELINE = baseline_dir()
  66. def temp_channels(cfg) -> list:
  67. """契约里 B.温度NBM 组的活通道 (以 _mean 结尾) —— temp_monthly/watch_channels 的列集。"""
  68. return [c for c in contracted_cols(cfg, groups=['B.温度NBM']) if c.endswith('_mean')]
  69. def month_of(s: pd.Series) -> pd.Series:
  70. return s.dt.to_period('M').astype(str)
  71. def build_temp_monthly(cfg, turbines=None, progress=True):
  72. """→ (长表 DataFrame[turbine, channel, month, med], 明细 dict)
  73. 规则: 每台 × 每月 × 每个温度通道, 在 p ≥ 500kW 行上的中位数。
  74. ★**不做"样本太少就不出月"的门** (用户令 2026-09-20 裁决: 保持现状、如实反映源件):
  75. 实测 2026-09 只有 1 行(补充件 `WTG01-B2.csv` 的最后一个时间戳 `2026-09-01 00:00`),
  76. 该月照样出一行 —— 页面时间窗按月过滤时它就"一格数据"。理由是**产物=当前 data/raw 的函数**:
  77. 源件里有这个时间戳, 就该看得见; 加"最少样本"门会让"数据到没到"这件事变成猜的。
  78. 要少这一格, 现场把那个时间戳裁掉或让现场别导出月初那一行即可, 不在本器里静默过滤。
  79. """
  80. chans = temp_channels(cfg)
  81. rows, info = [], {}
  82. for t in (turbines or cfg['turbines']):
  83. d = load_10min(t, cfg, groups=GROUPS)
  84. d['month'] = month_of(d['ts'])
  85. n_all = len(d)
  86. d = d[d['grd_wtc_ActPower_mean'] >= POWER_GATE]
  87. if not len(d):
  88. info[t] = dict(rows_all=n_all, rows_gate=0, months=0)
  89. continue
  90. med = d.groupby('month')[chans].median()
  91. long = med.stack().rename('med').reset_index()
  92. long.columns = ['month', 'channel', 'med']
  93. long.insert(0, 'turbine', t)
  94. rows.append(long)
  95. info[t] = dict(rows_all=n_all, rows_gate=len(d), months=len(med))
  96. if progress:
  97. print(f' {t}: 全 {n_all} 行 → 发电态 {len(d)} 行, {len(med)} 个月', flush=True)
  98. return (pd.concat(rows, ignore_index=True) if rows else pd.DataFrame(columns=['turbine', 'channel', 'month', 'med'])), info
  99. PRODUCTS = {
  100. 'temp_monthly.parquet': build_temp_monthly,
  101. }
  102. def verify(cfg) -> int:
  103. """逐件与随包件比对: 行数、键集合、逐值相等比例。
  104. 三种结果分开报, 不许混:
  105. ✅ 逐值一致 ❌ 有差异(规则需修正) ?? 找不到标准答案(**不能**当成一致)
  106. """
  107. ok_all, compared, no_base = True, 0, []
  108. for name in PRODUCTS:
  109. mine_p = P.store() / name
  110. if not mine_p.exists():
  111. print(f' [ -- ] {name:34s} 本仓还没算 (跳过)')
  112. continue
  113. base_p = baseline_for(name)
  114. if base_p is None:
  115. no_base.append(name)
  116. print(f' [ ?? ] {name:34s} 找不到随包件可对比 (台账基线只有 3 件; 随包整份产物在 _products_off*/ 下)')
  117. continue
  118. compared += 1
  119. a, b = pd.read_parquet(base_p), pd.read_parquet(mine_p)
  120. keys = [c for c in a.columns if c != 'med']
  121. m = a.merge(b, on=keys, how='outer', suffixes=('_base', '_mine'), indicator=True)
  122. val = [c for c in a.columns if c not in keys][0]
  123. both = m[m._merge == 'both']
  124. n = len(both)
  125. eq = int((both[f'{val}_base'].sub(both[f'{val}_mine']).abs() < 1e-9).sum())
  126. only_b = int((m._merge == 'left_only').sum())
  127. only_m = int((m._merge == 'right_only').sum())
  128. good = (n > 0 and eq == n and only_b == 0 and only_m == 0)
  129. ok_all &= good
  130. mark = '✅' if good else '❌'
  131. print(f' [{mark}] {name:34s} 基线 {len(a)} 行 / 本仓 {len(b)} 行; 逐值相等 {eq}/{n};'
  132. f' 仅基线 {only_b}; 仅本仓 {only_m} [标准答案: {P.rel(base_p)}]')
  133. if not good and n:
  134. d = both.assign(_d=(both[f'{val}_base'] - both[f'{val}_mine']).abs())
  135. print(' 差异最大的 3 格:')
  136. print(d.nlargest(3, '_d')[keys + [f'{val}_base', f'{val}_mine']].to_string(index=False))
  137. if not compared:
  138. print(f'\n结论: **无法比对**(没有一件找到随包件) —— 这不是"通过"。缺标准答案: {no_base}')
  139. return 5
  140. if no_base:
  141. print(f'\n注意: 另有 {len(no_base)} 件找不到标准答案, 未比对: {no_base}')
  142. print('\n结论:', '已比对的部分全部逐值一致' if ok_all else '有差异 —— 规则还需修正, 别急着写盘')
  143. return 0 if ok_all else 4
  144. def main() -> int:
  145. ap = argparse.ArgumentParser()
  146. ap.add_argument('--verify', action='store_true', help='只与随包基线比对, 不写盘')
  147. ap.add_argument('--turbines', default=None, help='只算这几台 (逗号分隔)')
  148. a = ap.parse_args()
  149. cfg = farm()
  150. ts = a.turbines.split(',') if a.turbines else None
  151. print(f'场站 {cfg["name"] if "name" in cfg else cfg["raw_station"]} · 温度通道 {len(temp_channels(cfg))} 个 · 发电态功率门 {POWER_GATE:.0f} kW\n')
  152. if not a.verify:
  153. for name, fn in PRODUCTS.items():
  154. print(f'== 算 {name}')
  155. df, info = fn(cfg, turbines=ts)
  156. out = P.store() / name
  157. out.parent.mkdir(parents=True, exist_ok=True)
  158. df.to_parquet(out, index=False)
  159. print(f' → {P.rel(out)} {df.shape[0]} 行 × {df.shape[1]} 列\n')
  160. print('== 与随包基线逐值比对 ==')
  161. return verify(cfg)
  162. if __name__ == '__main__':
  163. # 控制台可能是 GBK(中文 Windows 936): print 里的 ✅ ❌ 编不出来会抛 UnicodeEncodeError,
  164. # 明明算完了却以退出码 1 结束。降级为 '?' 而不是崩(同类坑见 src/console.py)。
  165. import sys as _sys
  166. for _s in (_sys.stdout, _sys.stderr):
  167. try: _s.reconfigure(errors='replace')
  168. except Exception: pass
  169. sys.exit(main())