scada_slim_build.py 7.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151
  1. #!/usr/bin/env python3
  2. # -*- coding: utf-8 -*-
  3. r"""SCADA 10min **精简轴仓**(slim10min):按窗重算的公共底座(用户令 2026-09-21)。
  4. ## 为什么要有它
  5. 用户令「部件问题/发电性能 应随时间窗变化」时,实测卡点不是算不出来,而是**每换一个窗都要重读 15 GB 10min CSV**:
  6. 变桨/偏航/蓄能/温度四个判级面与七镜头曲线各自的 `build_store()` 都是"逐台 `load_10min` 全量 → 只留一个固定窗",
  7. 换窗重算一次 ≈ 2 分钟/面。而它们真正需要的列只有 **约 50 列**(温度 `_mean` 27 + 功率/风况/转速/桨距/偏航/液压)。
  8. 本器把这份"公共列子集"抽成**逐台一份的窄仓**:一次全量扫(≈3 分钟),之后任何窗的重算都只是
  9. 窄仓上的过滤+分组(秒级),并且**口径与各面原有实现逐行同源**(同一份 `load_10min` 输出、同列同值)。
  10. ## 落点与口径
  11. outputs/<场>/windscada/slim10min/<台号>.parquet (列: ts + 公共列子集)
  12. outputs/<场>/windscada/slim10min/_manifest.json (行数/列清单/缺列/构建时间)
  13. · 缺列**如实记**(`columns_missing`),不造 0、不静默丢列 —— 各面按列名取,缺了就当"不可判"。
  14. · 本件是**产物**(raw-derived):落盘后自登记进 `outputs/<场>/_derived_manifest.json`。
  15. · 口径不变式:本件 = `load_10min(台, groups=GROUPS)` 的列子集,**不裁剪行**(时间范围 = 该台 CSV 全量)。
  16. ## 用法
  17. python scripts/scada_slim_build.py # 全 38 台
  18. python scripts/scada_slim_build.py --turbines WTG01,WTG02
  19. python scripts/scada_slim_build.py --check # 只看现状(行数/列/最新时间),不写
  20. 退出码: 0 成功 · 2 参数/源不齐 · 5 落盘为空(不覆盖既有仓)
  21. """
  22. from __future__ import annotations
  23. import argparse
  24. import json
  25. import pathlib
  26. import sys
  27. import time
  28. ROOT = pathlib.Path(__file__).resolve().parents[1]
  29. sys.path.insert(0, str(ROOT))
  30. from src import paths as P # noqa: E402
  31. from src.windscada.config import farm # noqa: E402
  32. from src.windscada.data import contracted_cols, load_10min # noqa: E402
  33. GROUPS = ['A.功率', 'A.风况', 'A.转速', 'B.变桨', 'B.温度NBM', 'B.偏航', 'B.润滑液压']
  34. # 公共列子集 = 四个判级面 + 七镜头曲线实际取用的列(逐条对着源码点名,别用整组,整组 128 列太肥)
  35. COLS_CORE = [
  36. 'grd_wtc_ActPower_mean', # 功率: 五态/pbin/发电态门
  37. 'tur_wtc_PowerRef_endvalue', # 限电命令面: perf.curtail.classify 的必需列 (缺它直接 KeyError)
  38. 'tur_wtc_SecAnemo_mean', # 机舱风: 曲线 X 轴 + 低风待机门 (活风源)
  39. 'tur_wtc_GenRpm_mean', 'tur_wtc_MainSRpm_mean', # 转速: 曲线/转矩/λ
  40. 'tur_wtc_PitcPosA_mean', 'tur_wtc_PitcPosB_mean', 'tur_wtc_PitcPosC_mean', # 桨距: 曲线/三叶不平衡
  41. 'tmp_wtc_AmbieTmp_mean', # 环温: 空气密度口径
  42. # 偏航面 (yaw.build_store)
  43. 'prs_wtc_YawPress_mean', 'prs_wtc_YawPress_max', 'prs_wtc_YawPress_min',
  44. 'dot_wtc_YawLubPu_timeon', 'din_wtc_Y1LubPis_timeon', 'flg_wtc_ScYawOpe_counts',
  45. 'cnt_wtc_WindFauS_endvalue', 'cnt_wtc_WindFauT_endvalue',
  46. # 蓄能面 (hydraulic.build_store)
  47. 'prs_wtc_HydPress_mean', 'prs_wtc_HydPress_max', 'prs_wtc_HydPress_min', 'tmp_wtc_HydOilTm_mean',
  48. ]
  49. def slim_cols(cfg) -> list[str]:
  50. """公共列 = 核心列 + 温度 NBM 的全部 `_mean`(温度面按 NBM 全通道比档)。"""
  51. temps = [c for c in contracted_cols(cfg, groups=['B.温度NBM']) if c.endswith('_mean')]
  52. out = []
  53. for c in COLS_CORE + temps:
  54. if c not in out:
  55. out.append(c)
  56. return out
  57. def out_dir(farm_name: str | None = None) -> pathlib.Path:
  58. return P.out_root(farm_name or P.farm()) / 'windscada' / 'slim10min'
  59. def build(cfg, want: list[str] | None, check: bool) -> int:
  60. cols = slim_cols(cfg)
  61. d = out_dir()
  62. d.mkdir(parents=True, exist_ok=True)
  63. man_p = d / '_manifest.json'
  64. man = json.loads(man_p.read_text(encoding='utf-8')) if man_p.is_file() else {}
  65. if check:
  66. import pandas as pd
  67. files = sorted(d.glob('WTG*.parquet'))
  68. print(f'== slim10min 现状 · {P.rel(d)} ==')
  69. print(f' 件 {len(files)} · 期望列 {len(cols)} · 手册: '
  70. f'{(man.get("at") or "—")} · 构建行 {man.get("rows_total", 0):,}')
  71. for f in files[:3]:
  72. x = pd.read_parquet(f, columns=['ts'])
  73. print(f' {f.name}: {len(x):,} 行 · {x.ts.min()} → {x.ts.max()}')
  74. if len(files) > 3:
  75. print(f' … 共 {len(files)} 件')
  76. return 0
  77. turbs = want or list(cfg['turbines'])
  78. total = 0
  79. missing_all: dict[str, list[str]] = {}
  80. t0 = time.time()
  81. for i, t in enumerate(turbs, 1):
  82. df = load_10min(t, cfg, groups=GROUPS)
  83. df = df.sort_values('ts').reset_index(drop=True)
  84. miss = [c for c in cols if c not in df.columns]
  85. if miss:
  86. missing_all[t] = miss
  87. keep = ['ts'] + [c for c in cols if c in df.columns]
  88. sub = df[keep]
  89. sub.to_parquet(d / f'{t}.parquet', index=False)
  90. total += len(sub)
  91. print(f' [{i}/{len(turbs)}] {t}: {len(sub):,} 行 × {len(keep) - 1} 列 '
  92. f'({sub.ts.min()} → {sub.ts.max()}){f" · 缺 {len(miss)} 列" if miss else ""} '
  93. f'已用 {time.time() - t0:.0f}s', flush=True)
  94. if not total:
  95. print('[X] 一行都没产出 —— 不覆盖既有仓')
  96. return 5
  97. man = dict(at=time.strftime('%Y-%m-%d %H:%M:%S'), turbines=len(turbs), rows_total=total,
  98. columns=cols, columns_n=len(cols), groups=GROUPS,
  99. columns_missing=missing_all,
  100. note='公共列子集(判级/曲线按窗重算的底座);不裁剪行;缺列如实记,不造数')
  101. man_p.write_text(json.dumps(man, ensure_ascii=False, indent=1), encoding='utf-8')
  102. print(f'\n合计 {total:,} 行 · 列 {len(cols)} · 用时 {time.time() - t0:.0f}s → {P.rel(d)}')
  103. if missing_all:
  104. print(f' ⚠ 有缺列的台 {len(missing_all)} 个(例 {list(missing_all)[:3]})—— 已写进 _manifest.json')
  105. try:
  106. from src.derived_manifest import record as _rec
  107. # ★键必须是**相对产物仓**的路径(2026-09-21 实逮: 先写成 basename `WTG01.parquet` ⇒
  108. # 反向审计按 `windscada/slim10min/*` 一条也匹配不上,39 件全被记成"未归类")。
  109. _root = P.out_root(P.farm())
  110. rels = {d.joinpath(p.name).relative_to(_root).as_posix():
  111. 'scada_slim_build.py(10min 公共列子集,供判级/曲线按窗重算)'
  112. for p in sorted(d.glob('WTG*.parquet'))}
  113. rels[d.joinpath('_manifest.json').relative_to(_root).as_posix()] = 'scada_slim_build.py(窄仓手册)'
  114. _rec(_root, rels, by='scada_slim_build.py')
  115. print(f' 自登记 {len(rels)} 件 → {P.rel(_root / "_derived_manifest.json")}')
  116. except Exception as e: # 登记失败不该把数据构建判成失败, 但要响亮
  117. print(f' ⚠ 自登记失败: {type(e).__name__}: {e}')
  118. return 0
  119. def main() -> int:
  120. ap = argparse.ArgumentParser(description='SCADA 10min 精简轴仓(按窗重算底座)')
  121. ap.add_argument('--turbines', default='', help='只做这几台(逗号分隔)')
  122. ap.add_argument('--check', action='store_true', help='只看现状, 不写')
  123. a = ap.parse_args()
  124. cfg = farm()
  125. want = [x.strip() for x in a.turbines.split(',') if x.strip()] or None
  126. return build(cfg, want, a.check)
  127. if __name__ == '__main__':
  128. raise SystemExit(main())