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