# -*- 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