#!/usr/bin/env python3 # -*- coding: utf-8 -*- r"""报警事件摄入: data/raw/<场站>/故障报警/*.xls → /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), 每行一个报警事件: 9 个季度/月度导出件 (2025Q1…2026.7) 的 `` 数与随包 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**(原本被跳过的大快照件因为少了一个键而被整件摄入, 把它覆盖的键又写了一遍) ⇒ 与用户令「重算要按 的最新变化」不符(既会重复计数, 又不是源集的函数)。 现在按**键集合**判: 全件并集去重 (键 = Name/Alarmcode/TimeOn), 同一键只留一行, 归属取"最窄的那个源件" (行数最少者优先, 同名再按文件名) —— 结果只由"盘上现在有哪些件"决定, 加件只会加键、删件只会去键。 2. **产物 = 当前 data/raw 的函数** (用户令 2026-09-20): 只保留**仍在盘**的源件带来的行; 已不在盘的行如实 丢掉并大声报出来(哪几个源件、多少行、涉及哪些月), 让人一眼分清"误删"还是"现场撤数"。 用法: python scripts/windscada_alarms_ingest.py --dry-run # 只报计划与证据 python scripts/windscada_alarms_ingest.py # 写 /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"]*)/?>", 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())