# -*- coding: utf-8 -*- """windscada L0 数据层: 契约驱动加载 + 校验门 (统一技术栈, 双翼共用). 铁律: 契约外列不加载; 校验门不过不进分析; 稀发数字量按类判活.""" import pandas as pd, pathlib, yaml from src.windscada.config import farm # config 属后端模块(P4),暂留原路径 def contract(cfg=None): cfg = cfg or farm() return yaml.safe_load(pathlib.Path(cfg['contract']).read_text(encoding='utf-8')) def contracted_cols(cfg=None, groups=None): c = contract(cfg) out = [] for g, chans in c['groups'].items(): if groups and g not in groups: continue out += [name for name, meta in chans.items() if meta['liveness'] == '活'] return sorted(set(out)) def gate(cfg=None): """G 门: 契约存在 + 活列占比 + 死嫌清单. 不过门 → 下游拒跑 (响亮).""" c = contract(cfg) total = sum(len(v) for v in c['groups'].values()) alive = sum(1 for v in c['groups'].values() for m in v.values() if m['liveness'] == '活') ok = alive / max(total, 1) >= 0.8 return dict(ok=ok, alive=alive, total=total, dead=c.get('dead_suspect', []), note=f"契约 {c['version']}: 活 {alive}/{total}; 未契约列 {c['n_cols_uncontracted']} 已登记不假装覆盖") def load_10min(turbine, cfg=None, groups=None): """一台的 10min 数据(**CSV 与 MDB 两种形态都认**,用户令 2026-09-19)。 取数走 `scada_source`:`/<台>.csv` 在位就用 CSV;该台没有 CSV 才回落 `data/raw/<场站>/scada_mdb/<年>年/<月>月/*.mdb`(现场年度归档形态)。用的哪种记在 返回表的 `attrs['scada_source']`(并在首次使用时打一行,不静默换源)。 """ cfg = cfg or farm() g = gate(cfg) if not g['ok']: raise RuntimeError(f"L0 校验门不过: {g['note']}") c = contract(cfg) want = [c['ts_col']] + contracted_cols(cfg, groups) from app_ETL.app_ETL_guanlan import scada_source as SS d = SS.load('10min', turbine, columns=want, cfg=cfg) src = str(d.attrs.get('scada_source') or '?') _note_source(src) # 契约外列不加载(两形态统一)—— CSV 侧原先在 read_csv 时用 usecols 做到, MDB 侧在这里收口 keep = [x for x in d.columns if x in want] d = d.loc[:, keep].copy() if not keep: raise RuntimeError(f'{turbine}: 源件里没有契约列 (想要 {want[:4]}…; 实际列 {list(d.columns)[:6]}…)') ts = c['ts_col'] d[ts] = pd.to_datetime(d[ts], errors='coerce') out = d.dropna(subset=[ts]).rename(columns={ts: 'ts'}) out.attrs['scada_source'] = src return out _NOTE_DONE: set = set() def _note_source(src: str) -> None: """首次用到某种形态时说明一次(让人知道这次重算是从 CSV 还是从 MDB 取的)。""" if src in _NOTE_DONE: return _NOTE_DONE.add(src) print(f' [i] SCADA 源形态: {src}' + ('(/<台>.csv)' if src == 'csv' else '(scada_mdb 年度归档库)'), flush=True) def source_report(cfg=None) -> dict: """本机 SCADA 收资形态盘点(页面/自检用): 各形态件数 + 逐台会走哪种。""" from app_ETL.app_ETL_guanlan import scada_source as SS cfg = cfg or farm() out = {} for cls in ('10min', '1min'): s = SS.sources(cls, cfg) turbs = SS.turbines(cls, cfg) kinds = {} for t in turbs: k = SS.source_of(cls, t, cfg) or 'none' kinds[k] = kinds.get(k, 0) + 1 out[cls] = dict(n_csv=s['n_csv'], n_mdb=s['n_mdb'], n_zip=s['n_zip'], n_turbines=len(turbs), by_source=kinds, csv_dir=str(s['csv_dir']), mdb_root=str(s['mdb_root'])) return out