slim.py 2.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869
  1. # -*- coding: utf-8 -*-
  2. r"""slim10min 窄仓读取口 —— **判级/曲线按窗重算的底座**(用户令 2026-09-21)。
  3. 窄仓由 `scripts/scada_slim_build.py` 构建(`outputs/<场>/windscada/slim10min/<台号>.parquet`,
  4. 列 = 现成 `load_10min` 的公共列子集,**不裁剪行**)。本模块只做三件事:
  5. have(cfg) / manifest(cfg) 现状与手册(缺件、缺列都如实回)
  6. load(t, cfg, span=(a, b), columns=[...]) 取一台、按窗(**含两端**)过滤
  7. loads(turbs, cfg, span, columns) 批量取(各面 build_store 按窗重算用)
  8. 为什么不让各面自己读 10min:换一个窗就要重读 15 GB CSV(实测 ≈2 分钟/面)。窄仓一次建好(≈2 分钟),
  9. 之后任何窗都只是窄仓上的过滤+分组(秒级),而口径与各面原实现**同源同列**。
  10. """
  11. from __future__ import annotations
  12. import json
  13. import pathlib
  14. import pandas as pd
  15. from src.windscada.config import farm # config 属后端模块(P4),暂留原路径
  16. def dir_of(cfg=None) -> pathlib.Path:
  17. """窄仓目录 = 产物仓下的 `slim10min/`(与 `scripts/scada_slim_build.py` 同一落点)。"""
  18. cfg = cfg or farm()
  19. return pathlib.Path(cfg['store']) / 'slim10min'
  20. def path_of(turbine: str, cfg=None) -> pathlib.Path:
  21. return dir_of(cfg) / f'{turbine}.parquet'
  22. def manifest(cfg=None) -> dict:
  23. p = dir_of(cfg) / '_manifest.json'
  24. return json.loads(p.read_text(encoding='utf-8')) if p.is_file() else {}
  25. def have(cfg=None) -> bool:
  26. return path_of('WTG01', cfg).is_file()
  27. def _mask(ts: pd.Series, span) -> pd.Series:
  28. """窗(含两端):span = ('YYYY-MM-DD', 'YYYY-MM-DD');末日晚 23:59:59 也算在内。"""
  29. a, b = span
  30. t = pd.to_datetime(ts, errors='coerce')
  31. return (t >= pd.Timestamp(a)) & (t <= pd.Timestamp(b) + pd.Timedelta(days=1) - pd.Timedelta(seconds=1))
  32. def load(turbine: str, cfg=None, span=None, columns=None) -> pd.DataFrame:
  33. """取一台窄仓:`span` 给了就按窗过滤(含两端),`columns` 是列白名单(缺列不报错)。"""
  34. p = path_of(turbine, cfg)
  35. if not p.is_file():
  36. raise FileNotFoundError(f'slim10min 缺 {p} —— 先跑 python scripts/scada_slim_build.py')
  37. want = None
  38. if columns:
  39. want = ['ts'] + [c for c in columns if c != 'ts']
  40. df = pd.read_parquet(p, columns=want)
  41. if span:
  42. df = df[_mask(df['ts'], span)].reset_index(drop=True)
  43. return df
  44. def loads(turbines, cfg=None, span=None, columns=None) -> dict[str, pd.DataFrame]:
  45. """批量取(逐台一份,缺台**响亮**报错,不静默少台)。"""
  46. out = {}
  47. for t in turbines:
  48. out[t] = load(t, cfg, span=span, columns=columns)
  49. return out