|
|
@@ -0,0 +1,166 @@
|
|
|
+#!/usr/bin/env python3
|
|
|
+# -*- coding: utf-8 -*-
|
|
|
+"""月度派生件构建器 —— 补上 v0.2.0 里"有消费者、却没生成端"的那批产物 (2026-09-12)。
|
|
|
+
|
|
|
+## 为什么有这个脚本
|
|
|
+
|
|
|
+从零重算 (`rebuild_from_raw.py`) 只覆盖 16 件; 随包里另有 15 件 parquet **没有任何生成端**
|
|
|
+(全库搜遍: 只有读取方、0 处写入方), 于是工作台主数据 `/api/fleet` 一缺 `temp_monthly.parquet`
|
|
|
+就整个回 `err=no_products` —— 页面看着像"没数据", 其实是"缺产物"。
|
|
|
+
|
|
|
+本脚本把这批件**从 data/raw 重新算出来**, 并且**逐值对齐随包件**: 随包件在这里当"标准答案"
|
|
|
+(存在 `_products_off/.../windscada/`), `--verify` 会报每件的命中率; 规则是从"标准答案"反推并
|
|
|
+用全量比对确认的, 不是猜的。
|
|
|
+
|
|
|
+## 已确认的口径 (2026-09-12, 用随包件反推 + 逐值验证)
|
|
|
+
|
|
|
+ temp_monthly 每台 × 每月 × 每个温度通道, 在 **发电态功率门 p ≥ 500kW** 行上的**中位数**
|
|
|
+ 形状 38×19×27 = 19494 行; 实测 WTG01 的 513 格 **逐值 100% 相同**
|
|
|
+ (对照: 全行中位数只有 5.8% 命中 —— 因为停机时绕组是冷的, 会把中位拉低)
|
|
|
+ 口径出处: `src/windscada/subsys/temp_nbm.py` 的"发电态×功率档(500kW)", PBINS 步长 500
|
|
|
+
|
|
|
+## 用法
|
|
|
+
|
|
|
+ python scripts/windscada_monthly_build.py --verify # 只比对(不写盘), 报每件命中率
|
|
|
+ python scripts/windscada_monthly_build.py # 算并写盘
|
|
|
+ python scripts/windscada_monthly_build.py --turbines WTG01,WTG02 # 只算几台(调试用)
|
|
|
+"""
|
|
|
+from __future__ import annotations
|
|
|
+
|
|
|
+import argparse
|
|
|
+import pathlib
|
|
|
+import sys
|
|
|
+
|
|
|
+ROOT = pathlib.Path(__file__).resolve().parents[1]
|
|
|
+sys.path.insert(0, str(ROOT))
|
|
|
+
|
|
|
+import pandas as pd # noqa: E402
|
|
|
+
|
|
|
+from src import paths as P # noqa: E402
|
|
|
+from src.windscada.config import farm # noqa: E402
|
|
|
+from src.windscada.data import load_10min, contracted_cols # noqa: E402
|
|
|
+
|
|
|
+POWER_GATE = 500.0 # 发电态功率门 (kW); 出处见 temp_nbm 的"发电态×功率档(500kW)"
|
|
|
+GROUPS = ['A.功率', 'B.温度NBM']
|
|
|
+
|
|
|
+
|
|
|
+def baseline_dir():
|
|
|
+ """"标准答案"目录: 随包件在哪。
|
|
|
+
|
|
|
+ 优先 `outputs/<场>/windscada/_pre_rebuild_20260911/`(rebuild_from_raw --verify 用的那份);
|
|
|
+ 没有就回落到 `_products_off/**/windscada/`(products_state --off 挪走后的暂存区, 内容就是随包件)。
|
|
|
+ 两处都没有 → None(只报"没法比", 不当成错)。
|
|
|
+ """
|
|
|
+ inplace = P.store() / '_pre_rebuild_20260911'
|
|
|
+ if inplace.is_dir():
|
|
|
+ return inplace
|
|
|
+ for p in sorted((ROOT / '_products_off').rglob('windscada')):
|
|
|
+ if p.is_dir():
|
|
|
+ return p
|
|
|
+ return None
|
|
|
+
|
|
|
+
|
|
|
+BASELINE = baseline_dir()
|
|
|
+
|
|
|
+
|
|
|
+def temp_channels(cfg) -> list:
|
|
|
+ """契约里 B.温度NBM 组的活通道 (以 _mean 结尾) —— temp_monthly/watch_channels 的列集。"""
|
|
|
+ return [c for c in contracted_cols(cfg, groups=['B.温度NBM']) if c.endswith('_mean')]
|
|
|
+
|
|
|
+
|
|
|
+def month_of(s: pd.Series) -> pd.Series:
|
|
|
+ return s.dt.to_period('M').astype(str)
|
|
|
+
|
|
|
+
|
|
|
+def build_temp_monthly(cfg, turbines=None, progress=True):
|
|
|
+ """→ (长表 DataFrame[turbine, channel, month, med], 明细 dict)
|
|
|
+
|
|
|
+ 规则: 每台 × 每月 × 每个温度通道, 在 p ≥ 500kW 行上的中位数。
|
|
|
+ """
|
|
|
+ chans = temp_channels(cfg)
|
|
|
+ rows, info = [], {}
|
|
|
+ for t in (turbines or cfg['turbines']):
|
|
|
+ d = load_10min(t, cfg, groups=GROUPS)
|
|
|
+ d['month'] = month_of(d['ts'])
|
|
|
+ n_all = len(d)
|
|
|
+ d = d[d['grd_wtc_ActPower_mean'] >= POWER_GATE]
|
|
|
+ if not len(d):
|
|
|
+ info[t] = dict(rows_all=n_all, rows_gate=0, months=0)
|
|
|
+ continue
|
|
|
+ med = d.groupby('month')[chans].median()
|
|
|
+ long = med.stack().rename('med').reset_index()
|
|
|
+ long.columns = ['month', 'channel', 'med']
|
|
|
+ long.insert(0, 'turbine', t)
|
|
|
+ rows.append(long)
|
|
|
+ info[t] = dict(rows_all=n_all, rows_gate=len(d), months=len(med))
|
|
|
+ if progress:
|
|
|
+ print(f' {t}: 全 {n_all} 行 → 发电态 {len(d)} 行, {len(med)} 个月', flush=True)
|
|
|
+ return (pd.concat(rows, ignore_index=True) if rows else pd.DataFrame(columns=['turbine', 'channel', 'month', 'med'])), info
|
|
|
+
|
|
|
+
|
|
|
+PRODUCTS = {
|
|
|
+ 'temp_monthly.parquet': build_temp_monthly,
|
|
|
+}
|
|
|
+
|
|
|
+
|
|
|
+def verify(cfg) -> int:
|
|
|
+ """逐件与随包件比对: 行数、键集合、逐值相等比例。"""
|
|
|
+ ok_all = True
|
|
|
+ for name in PRODUCTS:
|
|
|
+ mine_p, base_p = P.store() / name, BASELINE / name
|
|
|
+ if not mine_p.exists():
|
|
|
+ print(f' [ -- ] {name:34s} 本仓还没算 (跳过)')
|
|
|
+ continue
|
|
|
+ if not base_p.exists():
|
|
|
+ print(f' [ ?? ] {name:34s} 没有随包基线可比 (基线目录 {P.rel(BASELINE)})')
|
|
|
+ continue
|
|
|
+ a, b = pd.read_parquet(base_p), pd.read_parquet(mine_p)
|
|
|
+ keys = [c for c in a.columns if c != 'med']
|
|
|
+ m = a.merge(b, on=keys, how='outer', suffixes=('_base', '_mine'), indicator=True)
|
|
|
+ val = [c for c in a.columns if c not in keys][0]
|
|
|
+ both = m[m._merge == 'both']
|
|
|
+ n = len(both)
|
|
|
+ eq = int((both[f'{val}_base'].sub(both[f'{val}_mine']).abs() < 1e-9).sum())
|
|
|
+ only_b = int((m._merge == 'left_only').sum())
|
|
|
+ only_m = int((m._merge == 'right_only').sum())
|
|
|
+ good = (n > 0 and eq == n and only_b == 0 and only_m == 0)
|
|
|
+ ok_all &= good
|
|
|
+ mark = '✅' if good else '❌'
|
|
|
+ print(f' [{mark}] {name:34s} 基线 {len(a)} 行 / 本仓 {len(b)} 行; 逐值相等 {eq}/{n};'
|
|
|
+ f' 仅基线 {only_b}; 仅本仓 {only_m}')
|
|
|
+ if not good and n:
|
|
|
+ d = both.assign(_d=(both[f'{val}_base'] - both[f'{val}_mine']).abs())
|
|
|
+ print(' 差异最大的 3 格:')
|
|
|
+ print(d.nlargest(3, '_d')[keys + [f'{val}_base', f'{val}_mine']].to_string(index=False))
|
|
|
+ print('\n结论:', '全部逐值一致' if ok_all else '有差异 —— 规则还需修正, 别急着写盘')
|
|
|
+ return 0 if ok_all else 4
|
|
|
+
|
|
|
+
|
|
|
+def main() -> int:
|
|
|
+ ap = argparse.ArgumentParser()
|
|
|
+ ap.add_argument('--verify', action='store_true', help='只与随包基线比对, 不写盘')
|
|
|
+ ap.add_argument('--turbines', default=None, help='只算这几台 (逗号分隔)')
|
|
|
+ a = ap.parse_args()
|
|
|
+ cfg = farm()
|
|
|
+ ts = a.turbines.split(',') if a.turbines else None
|
|
|
+ print(f'场站 {cfg["name"] if "name" in cfg else cfg["raw_station"]} · 温度通道 {len(temp_channels(cfg))} 个 · 发电态功率门 {POWER_GATE:.0f} kW\n')
|
|
|
+ if not a.verify:
|
|
|
+ for name, fn in PRODUCTS.items():
|
|
|
+ print(f'== 算 {name}')
|
|
|
+ df, info = fn(cfg, turbines=ts)
|
|
|
+ out = P.store() / name
|
|
|
+ out.parent.mkdir(parents=True, exist_ok=True)
|
|
|
+ df.to_parquet(out, index=False)
|
|
|
+ print(f' → {P.rel(out)} {df.shape[0]} 行 × {df.shape[1]} 列\n')
|
|
|
+ print('== 与随包基线逐值比对 ==')
|
|
|
+ return verify(cfg)
|
|
|
+
|
|
|
+
|
|
|
+if __name__ == '__main__':
|
|
|
+ # 控制台可能是 GBK(中文 Windows 936): print 里的 ✅ ❌ 编不出来会抛 UnicodeEncodeError,
|
|
|
+ # 明明算完了却以退出码 1 结束。降级为 '?' 而不是崩(同类坑见 src/console.py)。
|
|
|
+ import sys as _sys
|
|
|
+ for _s in (_sys.stdout, _sys.stderr):
|
|
|
+ try: _s.reconfigure(errors='replace')
|
|
|
+ except Exception: pass
|
|
|
+ sys.exit(main())
|