windscada_alarms_ingest.py 10 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183
  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. ★2026-09-20 改口径: 原来按"逐件累加、一行新键都不贡献就跳过"的**贪心**判快照 —— 实测它依赖处理顺序,
  23. **删掉一个源件反而让总行数从 39,211 涨到 56,593**(原本被跳过的大快照件因为少了一个键而被整件摄入,
  24. 把它覆盖的键又写了一遍) ⇒ 与用户令「重算要按 <data/raw> 的最新变化」不符(既会重复计数, 又不是源集的函数)。
  25. 现在按**键集合**判: 全件并集去重 (键 = Name/Alarmcode/TimeOn), 同一键只留一行, 归属取"最窄的那个源件"
  26. (行数最少者优先, 同名再按文件名) —— 结果只由"盘上现在有哪些件"决定, 加件只会加键、删件只会去键。
  27. 2. **产物 = 当前 data/raw 的函数** (用户令 2026-09-20): 只保留**仍在盘**的源件带来的行; 已不在盘的行如实
  28. 丢掉并大声报出来(哪几个源件、多少行、涉及哪些月), 让人一眼分清"误删"还是"现场撤数"。
  29. 用法:
  30. python scripts/windscada_alarms_ingest.py --dry-run # 只报计划与证据
  31. python scripts/windscada_alarms_ingest.py # 写 <store>/alarms.parquet
  32. python scripts/windscada_alarms_ingest.py --farm rudong --raw D:\\别的现场
  33. """
  34. from __future__ import annotations
  35. import argparse
  36. import pathlib
  37. import re
  38. import sys
  39. ROOT = pathlib.Path(__file__).resolve().parents[1]
  40. sys.path.insert(0, str(ROOT))
  41. import pandas as pd # noqa: E402
  42. COLS = ["file", "turbine", "code", "text", "group", "t_on", "t_off", "dur_s", "dur_true_s"]
  43. KEY = ["turbine", "code", "t_on"]
  44. ATTR = re.compile(r'(\w+)="([^"]*)"')
  45. def parse_event_file(p: pathlib.Path) -> pd.DataFrame | None:
  46. """XML 报警导出 → DataFrame(COLS); 不是事件导出 (缺 TimeOn/Alarmcode) 返回 None。"""
  47. txt = p.read_text(encoding="utf-8-sig", errors="replace")
  48. rows = []
  49. for m in re.finditer(r"<row\b([^>]*)/?>", txt):
  50. a = dict(ATTR.findall(m.group(1)))
  51. if "TimeOn" not in a or "Alarmcode" not in a:
  52. continue
  53. t_on = pd.to_datetime(a.get("TimeOn") or None, errors="coerce")
  54. t_off = pd.to_datetime(a.get("TimeOff") or None, errors="coerce")
  55. dur_s = pd.to_numeric(a.get("DurationSec") or None, errors="coerce")
  56. rows.append(dict(file=p.name, turbine=a.get("Name"), code=a.get("Alarmcode"),
  57. text=a.get("AlarmText"), group=a.get("AlarmGroupText"),
  58. t_on=t_on, t_off=t_off, dur_s=dur_s,
  59. dur_true_s=(t_off - t_on).total_seconds() if pd.notna(t_off) and pd.notna(t_on) else float("nan")))
  60. if not rows:
  61. return None
  62. d = pd.DataFrame(rows, columns=COLS)
  63. d["code"] = d["code"].astype("string")
  64. for c in ("file", "turbine", "text", "group"):
  65. d[c] = d[c].astype("string")
  66. return d
  67. def main() -> int:
  68. ap = argparse.ArgumentParser()
  69. ap.add_argument("--farm", default=None)
  70. ap.add_argument("--raw", default=None, help="覆盖原始件根目录 (默认取场配置 data/raw)")
  71. ap.add_argument("--dry-run", action="store_true")
  72. a = ap.parse_args()
  73. from src.windscada.config import farm
  74. cfg = farm(a.farm)
  75. src = pathlib.Path(a.raw) if a.raw else pathlib.Path(str(cfg["src_alarm"]))
  76. store = pathlib.Path(cfg["store"])
  77. out = store / "alarms.parquet"
  78. print(f"源: {src}\n出: {out}\n")
  79. files = sorted(src.glob("*.xls"))
  80. if not files:
  81. raise SystemExit(f"源目录没有 .xls: {src}")
  82. parsed, skipped = {}, []
  83. for p in files:
  84. d = parse_event_file(p)
  85. if d is None:
  86. skipped.append((p.name, "XML 里没有 TimeOn/Alarmcode 行 —— 不是报警事件导出 (统计/报表件)"))
  87. continue
  88. parsed[p.name] = d
  89. # ★口径 (2026-09-20 用户令「重算要按 data/raw 最新变化」): **全件并集去重**, 结果只由盘上现有哪些件决定。
  90. # 同一键只留一行; 归属(file)取"最窄的源件"= 行数最少者优先, 同名再按文件名 —— 这样:
  91. # · 加件: 只可能新增键(或把某键的归属改到更窄的件上), 不会重复计数;
  92. # · 删件: 只可能去掉"只有它有"的键, 其余键的归属自动落到覆盖它的快照件上(并在下面报出来)。
  93. # 旧做法按"累加顺序 + 新键为零就跳过"判快照, 实测**删一个件反而多出 1.7 万行**(重复计数)。
  94. frames = []
  95. for n, d in parsed.items():
  96. dd = d.copy()
  97. dd["file"] = n
  98. frames.append(dd)
  99. allrows = pd.concat(frames, ignore_index=True) if frames else pd.DataFrame(columns=COLS)
  100. if len(allrows):
  101. allrows["_n"] = allrows["file"].map({n: len(d) for n, d in parsed.items()})
  102. allrows["_ord"] = list(range(len(allrows))) # 稳定序: 同件内保持原行序
  103. key = allrows[KEY].astype(str).agg("|".join, axis=1)
  104. allrows = (allrows.assign(_k=key)
  105. .sort_values(["_k", "_n", "file", "_ord"], kind="stable")
  106. .drop_duplicates("_k", keep="first")
  107. .drop(columns=["_n", "_ord", "_k"])
  108. .reset_index(drop=True))
  109. new = allrows[COLS] if len(allrows) else pd.DataFrame(columns=COLS)
  110. dup_rows = sum(len(d) for d in parsed.values()) - len(new)
  111. # 各源件的"贡献行数"与"被覆盖(并入更窄件)"情况 —— 让人看清快照件为什么不再是它的功劳
  112. contrib = new["file"].value_counts().to_dict() if len(new) else {}
  113. print("== 源件 (按贡献行数) ==")
  114. for n in sorted(parsed, key=lambda x: -contrib.get(x, 0)):
  115. tag = "摄入" if contrib.get(n) else "快照·全被覆盖"
  116. print(f" [{tag:10s}] {n:22s} {len(parsed[n]):6d} 行 → 贡献 {contrib.get(n, 0):6d} 行")
  117. for n, why in skipped:
  118. print(f" [非事件件·跳过] {n:22s} {why}")
  119. if dup_rows:
  120. print(f"\n== 去重证据: 逐件行数合计 {sum(len(d) for d in parsed.values())} ⇒ 键并集 {len(new)} "
  121. f"(去掉 {dup_rows} 行重复键; 键 = {'/'.join(KEY)}) ==")
  122. # 逐源件幂等: 只替换本次涉及的 file
  123. old = pd.read_parquet(out) if out.exists() else pd.DataFrame(columns=COLS)
  124. for c in ("file", "turbine", "code", "text", "group"):
  125. if c in old:
  126. old[c] = old[c].astype("string")
  127. # ★产物 = **当前源件的函数** (用户令 2026-09-20「重算要按 data/raw 的最新变化」):
  128. # 不再与旧产物合并(旧做法会把"源件已不在盘"的行一直留着), 旧产物只用来报"这回少了什么"。
  129. merged = new.sort_values(["t_on", "turbine", "code"], na_position="last").reset_index(drop=True)
  130. present = {p.name for p in files}
  131. gone_files = sorted(set(old["file"].dropna().unique()) - present) if len(old) else []
  132. old_keys = set(old[KEY].astype(str).agg("|".join, axis=1)) if len(old) else set()
  133. new_keys = set(merged[KEY].astype(str).agg("|".join, axis=1)) if len(merged) else set()
  134. lost = old_keys - new_keys
  135. gained = new_keys - old_keys
  136. print(f"\n== 结果 ==\n 旧 {len(old)} 行 → 新 {len(merged)} 行 "
  137. f"(键 {len(old_keys)} → {len(new_keys)}: 新增 {len(gained)} · 消失 {len(lost)})")
  138. if gone_files:
  139. print(f" [!] 有 {len(gone_files)} 个源件**已不在盘**: " + "、".join(gone_files[:6])
  140. + (" …" if len(gone_files) > 6 else "")
  141. + " —— 它们曾贡献的行照实不再产出(口径: 产物=当前 data/raw 的函数)")
  142. if lost:
  143. _lk = pd.DataFrame([k.split("|", 2) for k in lost], columns=KEY)
  144. _m = pd.to_datetime(_lk["t_on"], errors="coerce")
  145. print(f" 消失的键涉及月份: {_m.dt.strftime('%Y-%m').dropna().value_counts().sort_index().to_dict()}")
  146. print(" 若是**误删**: 把源件放回重跑本脚本即恢复; 若是现场撤数, 这就是正确结果。")
  147. print(f" 覆盖: {merged.t_on.min() if len(merged) else '—'} ~ {merged.t_on.max() if len(merged) else '—'}"
  148. f" 码 {merged.code.nunique() if len(merged) else 0} 个 台 {merged.turbine.nunique() if len(merged) else 0} 台")
  149. if a.dry_run:
  150. print("\n(dry-run, 未写盘)")
  151. return 0
  152. store.mkdir(parents=True, exist_ok=True)
  153. merged.to_parquet(out, index=False)
  154. print(f"\n已写 {out}")
  155. return 0
  156. if __name__ == "__main__":
  157. sys.exit(main())