| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183 |
- #!/usr/bin/env python3
- # -*- coding: utf-8 -*-
- r"""报警事件摄入: data/raw/<场站>/故障报警/*.xls → <store>/alarms.parquet (2026-09-11)
- ## 为什么有这个脚本
- 维护页「数据层 · 报警事件」一行原先写着这个命令, 但 v0.2.0 包里**没有这个文件** ——
- 声明与实现不符, 客户照页面敲命令只会报 not found; 而 alarms.parquet 又是约十个下游产物的
- 输入 (perf/curtail · perf/faults · subsys/yaw · subsys/hydraulic · subsys/fusion · perf/reliability …),
- 所以"从 data/raw 重算/更新"这条链一直断在这里。
- ## 源件真身 (2026-09-11 实测)
- `故障报警\*.xls` **不是** BIFF 表格, 而是 SpreadsheetML(XML), 每行一个报警事件:
- <row ID="…" TimeOn="2025-03-31 21:07:59.000" TimeOff="…" DurationSec="240"
- AlarmGroupText="Turbine" Name="WTG01" Alarmcode="3130" AlarmText="变桨润滑" … />
- 9 个季度/月度导出件 (2025Q1…2026.7) 的 `<row>` 数与随包 alarms.parquet 按 file 分组**逐件相同**
- (2773/4770/7830/4782/10446/1699/3288/2252/1371 = 39211), 故列映射按随包件锁定为:
- file=源文件名 · turbine=Name · code=Alarmcode · text=AlarmText · group=AlarmGroupText
- t_on=TimeOn · t_off=TimeOff · dur_s=DurationSec · dur_true_s=(TimeOff−TimeOn).total_seconds()
- (2025Q1 全量比对: 2773/2773 行、8 个字段**逐值一致**。)
- ## 两条纪律
- 1. **累计快照不叠加**: 目录里还有 `2025年全年故障记录.xls`(20155 行) 与 `2026年至今.xls`(19002 行),
- 它们分别是 Q1–Q4 并集 (2773+4770+7830+4782 = 20155, 精确相等) 与 2026 各件的并集 —— 属**累计快照**。
- 直接 concat 会让同一事件按文件数翻倍 (2026-08-31 长停台账虚增 35 倍是同款成因)。
- ★2026-09-20 改口径: 原来按"逐件累加、一行新键都不贡献就跳过"的**贪心**判快照 —— 实测它依赖处理顺序,
- **删掉一个源件反而让总行数从 39,211 涨到 56,593**(原本被跳过的大快照件因为少了一个键而被整件摄入,
- 把它覆盖的键又写了一遍) ⇒ 与用户令「重算要按 <data/raw> 的最新变化」不符(既会重复计数, 又不是源集的函数)。
- 现在按**键集合**判: 全件并集去重 (键 = Name/Alarmcode/TimeOn), 同一键只留一行, 归属取"最窄的那个源件"
- (行数最少者优先, 同名再按文件名) —— 结果只由"盘上现在有哪些件"决定, 加件只会加键、删件只会去键。
- 2. **产物 = 当前 data/raw 的函数** (用户令 2026-09-20): 只保留**仍在盘**的源件带来的行; 已不在盘的行如实
- 丢掉并大声报出来(哪几个源件、多少行、涉及哪些月), 让人一眼分清"误删"还是"现场撤数"。
- 用法:
- python scripts/windscada_alarms_ingest.py --dry-run # 只报计划与证据
- python scripts/windscada_alarms_ingest.py # 写 <store>/alarms.parquet
- python scripts/windscada_alarms_ingest.py --farm rudong --raw D:\\别的现场
- """
- from __future__ import annotations
- import argparse
- import pathlib
- import re
- import sys
- ROOT = pathlib.Path(__file__).resolve().parents[1]
- sys.path.insert(0, str(ROOT))
- import pandas as pd # noqa: E402
- COLS = ["file", "turbine", "code", "text", "group", "t_on", "t_off", "dur_s", "dur_true_s"]
- KEY = ["turbine", "code", "t_on"]
- ATTR = re.compile(r'(\w+)="([^"]*)"')
- def parse_event_file(p: pathlib.Path) -> pd.DataFrame | None:
- """XML 报警导出 → DataFrame(COLS); 不是事件导出 (缺 TimeOn/Alarmcode) 返回 None。"""
- txt = p.read_text(encoding="utf-8-sig", errors="replace")
- rows = []
- for m in re.finditer(r"<row\b([^>]*)/?>", txt):
- a = dict(ATTR.findall(m.group(1)))
- if "TimeOn" not in a or "Alarmcode" not in a:
- continue
- t_on = pd.to_datetime(a.get("TimeOn") or None, errors="coerce")
- t_off = pd.to_datetime(a.get("TimeOff") or None, errors="coerce")
- dur_s = pd.to_numeric(a.get("DurationSec") or None, errors="coerce")
- rows.append(dict(file=p.name, turbine=a.get("Name"), code=a.get("Alarmcode"),
- text=a.get("AlarmText"), group=a.get("AlarmGroupText"),
- t_on=t_on, t_off=t_off, dur_s=dur_s,
- dur_true_s=(t_off - t_on).total_seconds() if pd.notna(t_off) and pd.notna(t_on) else float("nan")))
- if not rows:
- return None
- d = pd.DataFrame(rows, columns=COLS)
- d["code"] = d["code"].astype("string")
- for c in ("file", "turbine", "text", "group"):
- d[c] = d[c].astype("string")
- return d
- def main() -> int:
- ap = argparse.ArgumentParser()
- ap.add_argument("--farm", default=None)
- ap.add_argument("--raw", default=None, help="覆盖原始件根目录 (默认取场配置 data/raw)")
- ap.add_argument("--dry-run", action="store_true")
- a = ap.parse_args()
- from src.windscada.config import farm
- cfg = farm(a.farm)
- src = pathlib.Path(a.raw) if a.raw else pathlib.Path(str(cfg["src_alarm"]))
- store = pathlib.Path(cfg["store"])
- out = store / "alarms.parquet"
- print(f"源: {src}\n出: {out}\n")
- files = sorted(src.glob("*.xls"))
- if not files:
- raise SystemExit(f"源目录没有 .xls: {src}")
- parsed, skipped = {}, []
- for p in files:
- d = parse_event_file(p)
- if d is None:
- skipped.append((p.name, "XML 里没有 TimeOn/Alarmcode 行 —— 不是报警事件导出 (统计/报表件)"))
- continue
- parsed[p.name] = d
- # ★口径 (2026-09-20 用户令「重算要按 data/raw 最新变化」): **全件并集去重**, 结果只由盘上现有哪些件决定。
- # 同一键只留一行; 归属(file)取"最窄的源件"= 行数最少者优先, 同名再按文件名 —— 这样:
- # · 加件: 只可能新增键(或把某键的归属改到更窄的件上), 不会重复计数;
- # · 删件: 只可能去掉"只有它有"的键, 其余键的归属自动落到覆盖它的快照件上(并在下面报出来)。
- # 旧做法按"累加顺序 + 新键为零就跳过"判快照, 实测**删一个件反而多出 1.7 万行**(重复计数)。
- frames = []
- for n, d in parsed.items():
- dd = d.copy()
- dd["file"] = n
- frames.append(dd)
- allrows = pd.concat(frames, ignore_index=True) if frames else pd.DataFrame(columns=COLS)
- if len(allrows):
- allrows["_n"] = allrows["file"].map({n: len(d) for n, d in parsed.items()})
- allrows["_ord"] = list(range(len(allrows))) # 稳定序: 同件内保持原行序
- key = allrows[KEY].astype(str).agg("|".join, axis=1)
- allrows = (allrows.assign(_k=key)
- .sort_values(["_k", "_n", "file", "_ord"], kind="stable")
- .drop_duplicates("_k", keep="first")
- .drop(columns=["_n", "_ord", "_k"])
- .reset_index(drop=True))
- new = allrows[COLS] if len(allrows) else pd.DataFrame(columns=COLS)
- dup_rows = sum(len(d) for d in parsed.values()) - len(new)
- # 各源件的"贡献行数"与"被覆盖(并入更窄件)"情况 —— 让人看清快照件为什么不再是它的功劳
- contrib = new["file"].value_counts().to_dict() if len(new) else {}
- print("== 源件 (按贡献行数) ==")
- for n in sorted(parsed, key=lambda x: -contrib.get(x, 0)):
- tag = "摄入" if contrib.get(n) else "快照·全被覆盖"
- print(f" [{tag:10s}] {n:22s} {len(parsed[n]):6d} 行 → 贡献 {contrib.get(n, 0):6d} 行")
- for n, why in skipped:
- print(f" [非事件件·跳过] {n:22s} {why}")
- if dup_rows:
- print(f"\n== 去重证据: 逐件行数合计 {sum(len(d) for d in parsed.values())} ⇒ 键并集 {len(new)} "
- f"(去掉 {dup_rows} 行重复键; 键 = {'/'.join(KEY)}) ==")
- # 逐源件幂等: 只替换本次涉及的 file
- old = pd.read_parquet(out) if out.exists() else pd.DataFrame(columns=COLS)
- for c in ("file", "turbine", "code", "text", "group"):
- if c in old:
- old[c] = old[c].astype("string")
- # ★产物 = **当前源件的函数** (用户令 2026-09-20「重算要按 data/raw 的最新变化」):
- # 不再与旧产物合并(旧做法会把"源件已不在盘"的行一直留着), 旧产物只用来报"这回少了什么"。
- merged = new.sort_values(["t_on", "turbine", "code"], na_position="last").reset_index(drop=True)
- present = {p.name for p in files}
- gone_files = sorted(set(old["file"].dropna().unique()) - present) if len(old) else []
- old_keys = set(old[KEY].astype(str).agg("|".join, axis=1)) if len(old) else set()
- new_keys = set(merged[KEY].astype(str).agg("|".join, axis=1)) if len(merged) else set()
- lost = old_keys - new_keys
- gained = new_keys - old_keys
- print(f"\n== 结果 ==\n 旧 {len(old)} 行 → 新 {len(merged)} 行 "
- f"(键 {len(old_keys)} → {len(new_keys)}: 新增 {len(gained)} · 消失 {len(lost)})")
- if gone_files:
- print(f" [!] 有 {len(gone_files)} 个源件**已不在盘**: " + "、".join(gone_files[:6])
- + (" …" if len(gone_files) > 6 else "")
- + " —— 它们曾贡献的行照实不再产出(口径: 产物=当前 data/raw 的函数)")
- if lost:
- _lk = pd.DataFrame([k.split("|", 2) for k in lost], columns=KEY)
- _m = pd.to_datetime(_lk["t_on"], errors="coerce")
- print(f" 消失的键涉及月份: {_m.dt.strftime('%Y-%m').dropna().value_counts().sort_index().to_dict()}")
- print(" 若是**误删**: 把源件放回重跑本脚本即恢复; 若是现场撤数, 这就是正确结果。")
- print(f" 覆盖: {merged.t_on.min() if len(merged) else '—'} ~ {merged.t_on.max() if len(merged) else '—'}"
- f" 码 {merged.code.nunique() if len(merged) else 0} 个 台 {merged.turbine.nunique() if len(merged) else 0} 台")
- if a.dry_run:
- print("\n(dry-run, 未写盘)")
- return 0
- store.mkdir(parents=True, exist_ok=True)
- merged.to_parquet(out, index=False)
- print(f"\n已写 {out}")
- return 0
- if __name__ == "__main__":
- sys.exit(main())
|