| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384 |
- # -*- 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`:`<src_10min>/<台>.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}'
- + ('(<src_10min>/<台>.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
|