windscada_alarms_ingest.py 7.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156
  1. #!/usr/bin/env python3
  2. # -*- coding: utf-8 -*-
  3. r"""报警事件摄入: data/raw/<场站>/故障报警/*.xls → <store>/alarms.parquet (2026-09-11)
  4. ## 为什么有这个脚本
  5. 维护页「数据层 · 报警事件」一行原先写着这个命令, 但 v0.2.0 包里**没有这个文件** ——
  6. 声明与实现不符, 客户照页面敲命令只会报 not found; 而 alarms.parquet 又是约十个下游产物的
  7. 输入 (perf/curtail · perf/faults · subsys/yaw · subsys/hydraulic · subsys/fusion · perf/reliability …),
  8. 所以"从 data/raw 重算/更新"这条链一直断在这里。
  9. ## 源件真身 (2026-09-11 实测)
  10. `故障报警\*.xls` **不是** BIFF 表格, 而是 SpreadsheetML(XML), 每行一个报警事件:
  11. <row ID="…" TimeOn="2025-03-31 21:07:59.000" TimeOff="…" DurationSec="240"
  12. AlarmGroupText="Turbine" Name="WTG01" Alarmcode="3130" AlarmText="变桨润滑" … />
  13. 9 个季度/月度导出件 (2025Q1…2026.7) 的 `<row>` 数与随包 alarms.parquet 按 file 分组**逐件相同**
  14. (2773/4770/7830/4782/10446/1699/3288/2252/1371 = 39211), 故列映射按随包件锁定为:
  15. file=源文件名 · turbine=Name · code=Alarmcode · text=AlarmText · group=AlarmGroupText
  16. t_on=TimeOn · t_off=TimeOff · dur_s=DurationSec · dur_true_s=(TimeOff−TimeOn).total_seconds()
  17. (2025Q1 全量比对: 2773/2773 行、8 个字段**逐值一致**。)
  18. ## 两条纪律
  19. 1. **累计快照不叠加**: 目录里还有 `2025年全年故障记录.xls`(20155 行) 与 `2026年至今.xls`(19002 行),
  20. 它们分别是 Q1–Q4 并集 (2773+4770+7830+4782 = 20155, 精确相等) 与 2026 各件的并集 —— 属**累计快照**。
  21. 直接 concat 会让同一事件按文件数翻倍 (2026-08-31 长停台账虚增 35 倍是同款成因)。本脚本按
  22. "键 (Name, Alarmcode, TimeOn) 是否有新增"判定: 一行新键都不贡献的件 = 快照, 跳过并打印证据。
  23. 2. **逐源件幂等**: 重跑时只替换本次涉及的那些 file 的行, 其余行原样保留 —— 这样"放新文件 → 重跑"
  24. 不会把已有月份洗掉, 也不会凭空复制。
  25. 用法:
  26. python scripts/windscada_alarms_ingest.py --dry-run # 只报计划与证据
  27. python scripts/windscada_alarms_ingest.py # 写 <store>/alarms.parquet
  28. python scripts/windscada_alarms_ingest.py --farm rudong --raw D:\\别的现场
  29. """
  30. from __future__ import annotations
  31. import argparse
  32. import pathlib
  33. import re
  34. import sys
  35. ROOT = pathlib.Path(__file__).resolve().parents[1]
  36. sys.path.insert(0, str(ROOT))
  37. import pandas as pd # noqa: E402
  38. COLS = ["file", "turbine", "code", "text", "group", "t_on", "t_off", "dur_s", "dur_true_s"]
  39. KEY = ["turbine", "code", "t_on"]
  40. ATTR = re.compile(r'(\w+)="([^"]*)"')
  41. def parse_event_file(p: pathlib.Path) -> pd.DataFrame | None:
  42. """XML 报警导出 → DataFrame(COLS); 不是事件导出 (缺 TimeOn/Alarmcode) 返回 None。"""
  43. txt = p.read_text(encoding="utf-8-sig", errors="replace")
  44. rows = []
  45. for m in re.finditer(r"<row\b([^>]*)/?>", txt):
  46. a = dict(ATTR.findall(m.group(1)))
  47. if "TimeOn" not in a or "Alarmcode" not in a:
  48. continue
  49. t_on = pd.to_datetime(a.get("TimeOn") or None, errors="coerce")
  50. t_off = pd.to_datetime(a.get("TimeOff") or None, errors="coerce")
  51. dur_s = pd.to_numeric(a.get("DurationSec") or None, errors="coerce")
  52. rows.append(dict(file=p.name, turbine=a.get("Name"), code=a.get("Alarmcode"),
  53. text=a.get("AlarmText"), group=a.get("AlarmGroupText"),
  54. t_on=t_on, t_off=t_off, dur_s=dur_s,
  55. dur_true_s=(t_off - t_on).total_seconds() if pd.notna(t_off) and pd.notna(t_on) else float("nan")))
  56. if not rows:
  57. return None
  58. d = pd.DataFrame(rows, columns=COLS)
  59. d["code"] = d["code"].astype("string")
  60. for c in ("file", "turbine", "text", "group"):
  61. d[c] = d[c].astype("string")
  62. return d
  63. def main() -> int:
  64. ap = argparse.ArgumentParser()
  65. ap.add_argument("--farm", default=None)
  66. ap.add_argument("--raw", default=None, help="覆盖原始件根目录 (默认取场配置 data/raw)")
  67. ap.add_argument("--dry-run", action="store_true")
  68. a = ap.parse_args()
  69. from src.windscada.config import farm
  70. cfg = farm(a.farm)
  71. src = pathlib.Path(a.raw) if a.raw else pathlib.Path(str(cfg["src_alarm"]))
  72. store = pathlib.Path(cfg["store"])
  73. out = store / "alarms.parquet"
  74. print(f"源: {src}\n出: {out}\n")
  75. files = sorted(src.glob("*.xls"))
  76. if not files:
  77. raise SystemExit(f"源目录没有 .xls: {src}")
  78. parsed, skipped = {}, []
  79. for p in files:
  80. d = parse_event_file(p)
  81. if d is None:
  82. skipped.append((p.name, "XML 里没有 TimeOn/Alarmcode 行 —— 不是报警事件导出 (统计/报表件)"))
  83. continue
  84. parsed[p.name] = d
  85. # 累计快照判定: 行数少的先入, 一行新键都不贡献的件视为快照
  86. order = sorted(parsed, key=lambda n: len(parsed[n]))
  87. kept, seen, snapshots = {}, set(), []
  88. for n in order:
  89. d = parsed[n]
  90. keys = set(map(tuple, d[KEY].astype(str).to_numpy()))
  91. new = keys - seen
  92. if not new:
  93. snapshots.append((n, len(d), len(keys)))
  94. continue
  95. kept[n] = d
  96. seen |= keys
  97. new = pd.concat(kept.values(), ignore_index=True) if kept else pd.DataFrame(columns=COLS)
  98. print("== 源件 ==")
  99. for n in sorted(parsed, key=lambda x: -len(parsed[x])):
  100. tag = "快照·跳过" if any(n == s[0] for s in snapshots) else "摄入"
  101. print(f" [{tag}] {n:22s} {len(parsed[n]):6d} 行")
  102. for n, why in skipped:
  103. print(f" [非事件件·跳过] {n:22s} {why}")
  104. if snapshots:
  105. print("\n== 累计快照证据 (精确等于既有件的并集, 叠加即翻倍) ==")
  106. for n, rows, keys in snapshots:
  107. print(f" {n}: {rows} 行 / {keys} 个唯一键, 新增键 0")
  108. # 逐源件幂等: 只替换本次涉及的 file
  109. old = pd.read_parquet(out) if out.exists() else pd.DataFrame(columns=COLS)
  110. for c in ("file", "turbine", "code", "text", "group"):
  111. if c in old:
  112. old[c] = old[c].astype("string")
  113. touch = set(kept)
  114. keep_old = old[~old["file"].isin(touch)] if len(old) else old
  115. merged = pd.concat([keep_old, new], ignore_index=True) if len(keep_old) or len(new) else old
  116. merged = merged.sort_values(["t_on", "turbine", "code"], na_position="last").reset_index(drop=True)
  117. print(f"\n== 结果 ==\n 旧 {len(old)} 行 → 新 {len(merged)} 行 "
  118. f"(本次涉及 {len(touch)} 个源件: {len(new)} 行; 保留未涉及的 {len(keep_old)} 行)")
  119. print(f" 覆盖: {merged.t_on.min()} ~ {merged.t_on.max()} 码 {merged.code.nunique()} 个 台 {merged.turbine.nunique()} 台")
  120. if a.dry_run:
  121. print("\n(dry-run, 未写盘)")
  122. return 0
  123. store.mkdir(parents=True, exist_ok=True)
  124. merged.to_parquet(out, index=False)
  125. print(f"\n已写 {out}")
  126. return 0
  127. if __name__ == "__main__":
  128. sys.exit(main())