#!/usr/bin/env python3 # -*- coding: utf-8 -*- r"""SCADA 10min **精简轴仓**(slim10min):按窗重算的公共底座(用户令 2026-09-21)。 ## 为什么要有它 用户令「部件问题/发电性能 应随时间窗变化」时,实测卡点不是算不出来,而是**每换一个窗都要重读 15 GB 10min CSV**: 变桨/偏航/蓄能/温度四个判级面与七镜头曲线各自的 `build_store()` 都是"逐台 `load_10min` 全量 → 只留一个固定窗", 换窗重算一次 ≈ 2 分钟/面。而它们真正需要的列只有 **约 50 列**(温度 `_mean` 27 + 功率/风况/转速/桨距/偏航/液压)。 本器把这份"公共列子集"抽成**逐台一份的窄仓**:一次全量扫(≈3 分钟),之后任何窗的重算都只是 窄仓上的过滤+分组(秒级),并且**口径与各面原有实现逐行同源**(同一份 `load_10min` 输出、同列同值)。 ## 落点与口径 outputs/<场>/windscada/slim10min/<台号>.parquet (列: ts + 公共列子集) outputs/<场>/windscada/slim10min/_manifest.json (行数/列清单/缺列/构建时间) · 缺列**如实记**(`columns_missing`),不造 0、不静默丢列 —— 各面按列名取,缺了就当"不可判"。 · 本件是**产物**(raw-derived):落盘后自登记进 `outputs/<场>/_derived_manifest.json`。 · 口径不变式:本件 = `load_10min(台, groups=GROUPS)` 的列子集,**不裁剪行**(时间范围 = 该台 CSV 全量)。 ## 用法 python scripts/scada_slim_build.py # 全 38 台 python scripts/scada_slim_build.py --turbines WTG01,WTG02 python scripts/scada_slim_build.py --check # 只看现状(行数/列/最新时间),不写 退出码: 0 成功 · 2 参数/源不齐 · 5 落盘为空(不覆盖既有仓) """ from __future__ import annotations try: from app_common.app_common_guanlan.api import install_root as _install_root except ImportError: # 理论不可达;包结构异常时回退到按位置上跳 from pathlib import Path as _P def _install_root(_f): return _P(_f).resolve().parents[3] import argparse import json import pathlib import sys import time ROOT = _install_root(__file__) # 模块化后按标记找安装根(原 parents[1] 已不成立) sys.path.insert(0, str(ROOT)) from src import paths as P # noqa: E402 from src.windscada.config import farm # noqa: E402 from src.windscada.data import contracted_cols, load_10min # noqa: E402 GROUPS = ['A.功率', 'A.风况', 'A.转速', 'B.变桨', 'B.温度NBM', 'B.偏航', 'B.润滑液压'] # 公共列子集 = 四个判级面 + 七镜头曲线实际取用的列(逐条对着源码点名,别用整组,整组 128 列太肥) COLS_CORE = [ 'grd_wtc_ActPower_mean', # 功率: 五态/pbin/发电态门 'grd_wtc_ActPower_max', # 封顶参数 M9: p_cap = 99.5 分位 'tur_wtc_GenRpm_max', # 封顶参数 M9: w_cap 'tur_wtc_PowerRef_endvalue', # 限电命令面: perf.curtail.classify 的必需列 (缺它直接 KeyError) 'tur_wtc_SecAnemo_mean', # 机舱风: 曲线 X 轴 + 低风待机门 (活风源) 'tur_wtc_GenRpm_mean', 'tur_wtc_MainSRpm_mean', # 转速: 曲线/转矩/λ 'tur_wtc_PitcPosA_mean', 'tur_wtc_PitcPosB_mean', 'tur_wtc_PitcPosC_mean', # 桨距: 曲线/三叶不平衡 'tmp_wtc_AmbieTmp_mean', # 环温: 空气密度口径 # 偏航面 (yaw.build_store) 'prs_wtc_YawPress_mean', 'prs_wtc_YawPress_max', 'prs_wtc_YawPress_min', 'dot_wtc_YawLubPu_timeon', 'din_wtc_Y1LubPis_timeon', 'flg_wtc_ScYawOpe_counts', 'cnt_wtc_WindFauS_endvalue', 'cnt_wtc_WindFauT_endvalue', # 蓄能面 (hydraulic.build_store) 'prs_wtc_HydPress_mean', 'prs_wtc_HydPress_max', 'prs_wtc_HydPress_min', 'tmp_wtc_HydOilTm_mean', ] def slim_cols(cfg) -> list[str]: """公共列 = 核心列 + 温度 NBM 的全部 `_mean`(温度面按 NBM 全通道比档)。""" temps = [c for c in contracted_cols(cfg, groups=['B.温度NBM']) if c.endswith('_mean')] out = [] for c in COLS_CORE + temps: if c not in out: out.append(c) return out def out_dir(farm_name: str | None = None) -> pathlib.Path: return P.out_root(farm_name or P.farm()) / 'windscada' / 'slim10min' def build(cfg, want: list[str] | None, check: bool) -> int: cols = slim_cols(cfg) d = out_dir() d.mkdir(parents=True, exist_ok=True) man_p = d / '_manifest.json' man = json.loads(man_p.read_text(encoding='utf-8')) if man_p.is_file() else {} if check: import pandas as pd files = sorted(d.glob('WTG*.parquet')) print(f'== slim10min 现状 · {P.rel(d)} ==') print(f' 件 {len(files)} · 期望列 {len(cols)} · 手册: ' f'{(man.get("at") or "—")} · 构建行 {man.get("rows_total", 0):,}') for f in files[:3]: x = pd.read_parquet(f, columns=['ts']) print(f' {f.name}: {len(x):,} 行 · {x.ts.min()} → {x.ts.max()}') if len(files) > 3: print(f' … 共 {len(files)} 件') return 0 turbs = want or list(cfg['turbines']) total = 0 missing_all: dict[str, list[str]] = {} t0 = time.time() for i, t in enumerate(turbs, 1): df = load_10min(t, cfg, groups=GROUPS) df = df.sort_values('ts').reset_index(drop=True) miss = [c for c in cols if c not in df.columns] if miss: missing_all[t] = miss keep = ['ts'] + [c for c in cols if c in df.columns] sub = df[keep] sub.to_parquet(d / f'{t}.parquet', index=False) total += len(sub) print(f' [{i}/{len(turbs)}] {t}: {len(sub):,} 行 × {len(keep) - 1} 列 ' f'({sub.ts.min()} → {sub.ts.max()}){f" · 缺 {len(miss)} 列" if miss else ""} ' f'已用 {time.time() - t0:.0f}s', flush=True) if not total: print('[X] 一行都没产出 —— 不覆盖既有仓') return 5 man = dict(at=time.strftime('%Y-%m-%d %H:%M:%S'), turbines=len(turbs), rows_total=total, columns=cols, columns_n=len(cols), groups=GROUPS, columns_missing=missing_all, note='公共列子集(判级/曲线按窗重算的底座);不裁剪行;缺列如实记,不造数') man_p.write_text(json.dumps(man, ensure_ascii=False, indent=1), encoding='utf-8') print(f'\n合计 {total:,} 行 · 列 {len(cols)} · 用时 {time.time() - t0:.0f}s → {P.rel(d)}') if missing_all: print(f' ⚠ 有缺列的台 {len(missing_all)} 个(例 {list(missing_all)[:3]})—— 已写进 _manifest.json') try: from src.derived_manifest import record as _rec # ★键必须是**相对产物仓**的路径(2026-09-21 实逮: 先写成 basename `WTG01.parquet` ⇒ # 反向审计按 `windscada/slim10min/*` 一条也匹配不上,39 件全被记成"未归类")。 _root = P.out_root(P.farm()) rels = {d.joinpath(p.name).relative_to(_root).as_posix(): 'scada_slim_build.py(10min 公共列子集,供判级/曲线按窗重算)' for p in sorted(d.glob('WTG*.parquet'))} rels[d.joinpath('_manifest.json').relative_to(_root).as_posix()] = 'scada_slim_build.py(窄仓手册)' _rec(_root, rels, by='scada_slim_build.py') print(f' 自登记 {len(rels)} 件 → {P.rel(_root / "_derived_manifest.json")}') except Exception as e: # 登记失败不该把数据构建判成失败, 但要响亮 print(f' ⚠ 自登记失败: {type(e).__name__}: {e}') return 0 def main() -> int: ap = argparse.ArgumentParser(description='SCADA 10min 精简轴仓(按窗重算底座)') ap.add_argument('--turbines', default='', help='只做这几台(逗号分隔)') ap.add_argument('--check', action='store_true', help='只看现状, 不写') a = ap.parse_args() cfg = farm() want = [x.strip() for x in a.turbines.split(',') if x.strip()] or None return build(cfg, want, a.check) if __name__ == '__main__': raise SystemExit(main())