data.py 3.8 KB

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