#!/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 倍是同款成因)。本脚本按 "键 (Name, Alarmcode, TimeOn) 是否有新增"判定: 一行新键都不贡献的件 = 快照, 跳过并打印证据。 2. **逐源件幂等**: 重跑时只替换本次涉及的那些 file 的行, 其余行原样保留 —— 这样"放新文件 → 重跑" 不会把已有月份洗掉, 也不会凭空复制。 用法: 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 # 累计快照判定: 行数少的先入, 一行新键都不贡献的件视为快照 order = sorted(parsed, key=lambda n: len(parsed[n])) kept, seen, snapshots = {}, set(), [] for n in order: d = parsed[n] keys = set(map(tuple, d[KEY].astype(str).to_numpy())) new = keys - seen if not new: snapshots.append((n, len(d), len(keys))) continue kept[n] = d seen |= keys new = pd.concat(kept.values(), ignore_index=True) if kept else pd.DataFrame(columns=COLS) print("== 源件 ==") for n in sorted(parsed, key=lambda x: -len(parsed[x])): tag = "快照·跳过" if any(n == s[0] for s in snapshots) else "摄入" print(f" [{tag}] {n:22s} {len(parsed[n]):6d} 行") for n, why in skipped: print(f" [非事件件·跳过] {n:22s} {why}") if snapshots: print("\n== 累计快照证据 (精确等于既有件的并集, 叠加即翻倍) ==") for n, rows, keys in snapshots: print(f" {n}: {rows} 行 / {keys} 个唯一键, 新增键 0") # 逐源件幂等: 只替换本次涉及的 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") touch = set(kept) keep_old = old[~old["file"].isin(touch)] if len(old) else old merged = pd.concat([keep_old, new], ignore_index=True) if len(keep_old) or len(new) else old merged = merged.sort_values(["t_on", "turbine", "code"], na_position="last").reset_index(drop=True) print(f"\n== 结果 ==\n 旧 {len(old)} 行 → 新 {len(merged)} 行 " f"(本次涉及 {len(touch)} 个源件: {len(new)} 行; 保留未涉及的 {len(keep_old)} 行)") print(f" 覆盖: {merged.t_on.min()} ~ {merged.t_on.max()} 码 {merged.code.nunique()} 个 台 {merged.turbine.nunique()} 台") 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())