Pārlūkot izejas kodu

数据链补齐: 从 data/raw 重算/更新产物 (报警/工单/油样) + 等价验收

背景: 说明书 §11 写着「不含从原始数据重生成产物的链 (P0)」; 维护页那四条「摄入命令」指向的文件
在包里**一个都不存在**, 于是现场把新台账放进 data/raw/如东 也不会更新, 工单台账一直停在 2024-11-21。

新增 (scripts/):
- windscada_alarms_ingest.py        故障报警/*.xls → alarms.parquet
  .xls 实为 SpreadsheetML(XML); 列映射按随包件锁定; 累计快照件自动跳过
  (2025年全年 = Q1..Q4 精确并集 20155 行, 2026年至今 = 19002 行, 新增键均为 0)
- windscada_workorder_ingest.py     风机故障记录/**/*.xls(x) → workorders.parquet (26 列)
  集团口径导出件按场名过滤 (如海/如东 — 实测一张 2021年3月 表里混着民勤/来福/宝力格/大岗子/宏基等十来个场);
  跨表内容重复归并 1085 行 (13 行「计划停机」被复制进 58 张月表 = 754 行冗余);
  日期解析按 src/sop/ledger_dates.py 同一纪律 (全角冒号 / Excel 序列号 / 2022.1.10.19:13 …),
  解析不了的整行保留并在报告里点名, 不静默丢行
- windscada_watch_channels_build.py 油样报告/**/*.pdf → oil_samples_index.parquet
  文件名解析 (DDMMYYYY_sampleid_…_NN#_部件); 合并语义: 认不出的老行原样保留
- rebuild_from_raw.py               总入口 + --verify 等价验收 + --scada (包内 10 个 SCADA 构建器)

等价验收 (scripts/rebuild_from_raw.py --verify, 对 _pre_rebuild_20260911/ 随包基线):
- alarms            39211 / 39211 行     逐值完全一致
- oil_samples_index   506 /   506 行     逐值完全一致
- workorders         随包 1574 种内容复现 1566 种; 余 8 种差异逐条定位:
    · 7 种 = 复位运行时间/t_reset: 随包为 NA, 本链有值 —— 旧链没认源列名「复位时间时间」这个错别字
    · 1 种 = 随包把 Excel 序列号当纳秒解析成 1970-01-01, 本链按序列号解析为 2023-05-13
    随包 78 行副本重复归并; 另多出 69 行 = 旧链整行丢弃的「日期解析失败」行

效果 (页面 /detail/v2#tab=system 数据层, 已重启服务实测):
  检修工单台账  2020-01-03 ~ 2024-11-21 / 1652 行  →  2020-01-03 ~ 2026-07-08 / 5876 行
  四条「摄入命令」命令可用  False → True

其它:
- requirements.txt 加 xlrd==2.0.2 (2023/2024 年总表是真 BIFF .xls, openpyxl 读不了);
  轮子放进 wheels/win_amd64/ (wheels/ 按 .gitignore 不入库, 随包分发)
- config.py: 新增 src_farm_names (原始件里的场站名写法: 如海/如东)
- maintenance.py: 「★当前止 2024-11-21 …缺口需现场补」随实际覆盖更新;
  SCADA 10min 行的摄入命令换成真实存在的 rebuild_from_raw.py --scada (原 scripts/windscada_report.py 不存在)
- data_tpl_en.py: 同步更新上述 5 条英文映射 (否则英文页这 4 格回落中文)
- _修复记录_20260911/README.md: 增「数据链补齐」一节
zhouyang.xie 4 nedēļas atpakaļ
vecāks
revīzija
01ddb723e1

+ 42 - 0
_修复记录_20260911/README.md

@@ -80,6 +80,48 @@ start.bat.bak-healthz
 
 
 
+## 数据链补齐 · 让系统真正吃 data/raw/如东 (2026-09-11)
+
+**问题**: 说明书 §11 写着「不含从原始数据重生成产物的链 (P0)」, 而维护页四条「摄入命令」指向的文件
+**在包里一个都不存在**; 于是现场把新台账放进 `data/raw/如东/…` 也不会更新, 工单台账一直停在 2024-11-21。
+
+**新增 4 个脚本** (都在 `scripts/`; 维护页四条命令现在全部「命令可用=True」):
+
+| 脚本 | 源 → 产物 | 关键实现 |
+|---|---|---|
+| `windscada_alarms_ingest.py` | `故障报警/*.xls` → `alarms.parquet` | `.xls` 实为 SpreadsheetML(XML); 累计快照件自动跳过 (2025年全年 = Q1..Q4 精确并集 20155 行) |
+| `windscada_workorder_ingest.py` | `风机故障记录/**` → `workorders.parquet` | 按场名过滤 (集团表里混着民勤/来福/宝力格等十来个场); 归并跨表重复 1085 行; 日期按 `src/sop/ledger_dates.py` 同一纪律解析, 解析不了的整行保留 |
+| `windscada_watch_channels_build.py` | `油样报告/**/*.pdf` → `oil_samples_index.parquet` | 文件名解析 (日期/台号/部件/sample_id); 合并语义: 认不出的老行原样保留 |
+| `rebuild_from_raw.py` | 总入口 | `--verify` 等价验收; `--scada` 再跑包内 10 个 SCADA 侧构建器 |
+
+**等价验收** (`python scripts/rebuild_from_raw.py --verify`, 对 `_pre_rebuild_20260911/` 随包基线):
+
+```
+报警事件       39211 行 / 39211 行   逐值完全一致            ✔
+油液化验         506 行 /   506 行   逐值完全一致            ✔
+检修工单台账    1574 种内容复现 1566 种; 8 种差异逐条定位:
+   · 7 种 = 复位运行时间: 随包 NA、本链有值 — 旧链没认源列名「复位时间时间」这个错别字
+   · 1 种 = 随包把 Excel 序列号当纳秒解析成 1970-01-01, 本链解析为 2023-05-13
+   · 随包 78 行副本重复已归并; 另多出 69 行 = 旧链丢弃的日期不可解析行
+```
+
+**效果** (`/detail/v2#tab=system` 数据层, 已重启服务实测):
+
+| 行 | 改前 | 改后 |
+|---|---|---|
+| 检修工单台账 | 2020-01-03 ~ 2024-11-21, 1652 行 | **2020-01-03 ~ 2026-07-08, 5876 行** |
+| 四条「摄入命令」命令可用 | False | **True** |
+
+**依赖**: `requirements.txt` 加 `xlrd==2.0.2` (2023/2024 年总表是真 BIFF `.xls`, openpyxl 读不了),
+轮子已放进 `wheels/win_amd64/` (`wheels/` 按 .gitignore 不入库, 随包分发)。
+
+**SCADA 侧**: `python scripts/rebuild_from_raw.py --scada` 把 10 个构建器按依赖顺序跑一遍
+(powercurve → loss_monthly → curves/control/stop_events/温度/偏航/液压/热链/系统辅助)。
+抽验过等价性 (WTG01/WTG02 的 `temp_bins` 与随包件 1296/1296 行逐值相同、med 差 0.0),
+但**未整体重跑** —— 要跑建议先备份 `outputs/rudong/windscada/` 再逐项比对覆盖区间。
+
+## 仍然存在 (不是代码问题)
+
 - `/local-ai/` 仍是 DOWN: 本机 Ollama 没运行, 且 `%USERPROFILE%\.ollama` 下**没有任何模型**
   (`manifests` 都没有), 所以问答与本地审核页不可用。其余 6 个页面不受影响。
   启用: 启动 Ollama, 按 `configs\models.json` 的 `pull_commands` 拉模型 (需联网/大流量),

+ 2 - 1
requirements.txt

@@ -16,4 +16,5 @@ requests==2.34.2
 rich==15.0.0
 tqdm==4.68.3
 openpyxl==3.1.5
-scipy==1.18.0
+scipy==1.18.0
+xlrd==2.0.2

+ 206 - 0
scripts/rebuild_from_raw.py

@@ -0,0 +1,206 @@
+#!/usr/bin/env python3
+# -*- coding: utf-8 -*-
+r"""从 data/raw 重算/更新产物 (2026-09-11) —— 维护页那四条「摄入命令」的总入口。
+
+## 背景
+
+`docs/说明书…v0.2.md` §11 写着「**不含**从原始数据重生成产物的链 (P0)」, 于是随包产物只能吃预生成的那份
+(台账止 2024-11-21、报警止 2026-07-15), 现场把新表放进 `data/raw/<场站>/…` 也不会更新。
+本脚本把这条链补上, 顺序与依赖关系如下 (都按场配置 `src.windscada.config`):
+
+    ① scripts/windscada_alarms_ingest.py        故障报警/*.xls      → alarms.parquet
+    ② scripts/windscada_workorder_ingest.py     风机故障记录/**      → workorders.parquet
+    ③ scripts/windscada_watch_channels_build.py 油样报告/**         → oil_samples_index.parquet
+    ── 以上三步只吃现场台账 (xls/xlsx/pdf), 快, 秒级到分钟级 ──
+    ④ (--scada) 包内已有的 SCADA 侧构建器: 从 scada_10min/*.csv 重算
+       powercurve → loss_monthly(availability) → curves / control / stop_events(faults)
+       / temp_bins / yaw_daily / hydraulic_accum / thermal_chain / system_aux
+       (这一步要逐台读 38 个 ~380MB 的 CSV, 慢; 默认不跑, 加 --scada 才跑)
+
+## --verify: 等价验收 (本链的验收标准)
+
+重算件必须能**复现随包件**: 逐键逐值比对 `outputs/rudong/windscada/_pre_rebuild_20260911/` 里那份
+随包基线 (2026-09-11 现场修复时留的备份)。允许的差异只有两类, 都必须被点名:
+  · 旧链的**已知缺陷**: 报警/工单里同一副本被收两遍、日期解析不了就整行丢弃、Excel 序列号当纳秒解析;
+  · 本链的**新增覆盖**: 随包时没接的源件 (74 张台账表 vs 随包只用了 6 张)。
+无法归类的差异 = 验收不通过 (退出码 4)。
+
+用法:
+    python scripts/rebuild_from_raw.py                 # 跑 ①②③
+    python scripts/rebuild_from_raw.py --verify        # 只验收, 不重算
+    python scripts/rebuild_from_raw.py --scada         # 连 ④ 一起跑 (慢)
+    python scripts/rebuild_from_raw.py --dry-run       # 只打印各步会做什么
+"""
+from __future__ import annotations
+
+import argparse
+import pathlib
+import subprocess
+import sys
+from collections import Counter
+
+ROOT = pathlib.Path(__file__).resolve().parents[1]
+sys.path.insert(0, str(ROOT))
+
+BASE = 'outputs/rudong/windscada/_pre_rebuild_20260911'       # 随包基线备份
+STEPS = [
+    ('报警事件', 'scripts/windscada_alarms_ingest.py', 'alarms.parquet'),
+    ('检修工单台账', 'scripts/windscada_workorder_ingest.py', 'workorders.parquet'),
+    ('油液化验', 'scripts/windscada_watch_channels_build.py', 'oil_samples_index.parquet'),
+]
+
+# SCADA 侧 (--scada): 全部是包内**已有**的构建器, 只是过去没有入口把它们串起来。
+# 顺序 = 依赖顺序 (powercurve 先出 bins, availability 才吃得到; stop_events 要吃 alarms)。
+SCADA_STEPS = [
+    ('功率曲线', 'src.windscada.perf.powercurve', 'store'),
+    ('可用率与损失 (loss_monthly)', 'src.windscada.perf.availability', 'build'),
+    ('七镜头曲线', 'src.windscada.perf.curves', 'build_store'),
+    ('控制策略件', 'src.windscada.perf.control', 'build_store'),
+    ('停机事件 (stop_events)', 'src.windscada.perf.faults', 'build_stop_events'),
+    ('温度 NBM (temp_bins)', 'src.windscada.subsys.temp_nbm', 'build_store'),
+    ('偏航 (yaw_daily)', 'src.windscada.subsys.yaw', 'build_store'),
+    ('液压蓄能 (hydraulic_accum)', 'src.windscada.subsys.hydraulic', 'build_store'),
+    ('热链 (thermal_chain)', 'src.windscada.subsys.thermal_chain', 'build_store'),
+    ('系统辅助 (system_aux)', 'src.windscada.taxonomy', 'build_aux_store'),
+]
+
+
+def run_scada(farm_name=None) -> int:
+    """逐台读 scada_10min/*.csv 重算 SCADA 侧产物。慢 (38 台 × 各构建器), 失败不中断, 最后汇总。"""
+    import importlib
+    import time
+    from src.windscada.config import farm
+    cfg = farm(farm_name)
+    print(f'场: {cfg["name"]}  源: {cfg["src_10min"]}  仓: {cfg["store"]}')
+    bad = []
+    for label, mod, fn in SCADA_STEPS:
+        t0 = time.time()
+        try:
+            f = getattr(importlib.import_module(mod), fn)
+            r = f(cfg)
+            print(f'  ✔ {label:26s} {time.time()-t0:6.1f}s  {str(r)[:110]}', flush=True)
+        except Exception as e:
+            bad.append((label, f'{type(e).__name__}: {e}'))
+            print(f'  ✘ {label:26s} {time.time()-t0:6.1f}s  {type(e).__name__}: {str(e)[:150]}', flush=True)
+    if bad:
+        print(f'\n{len(bad)} 个构建器失败:')
+        for label, err in bad:
+            print(f'  {label}: {err}')
+    return 0 if not bad else 5
+
+
+def _run(script: str, extra=()) -> int:
+    cmd = [sys.executable, str(ROOT / script), *extra]
+    print(f'\n$ {" ".join(cmd[1:])}', flush=True)
+    return subprocess.call(cmd, cwd=str(ROOT))
+
+
+def _canon(df, cols=None):
+    import pandas as pd
+    d = df[cols] if cols else df
+    out = []
+    for r in d.itertuples(index=False):
+        row = []
+        for c, v in zip(d.columns, r):
+            try:
+                if v is None or (not isinstance(v, str) and pd.isna(v)):
+                    row.append('<NA>')
+                    continue
+            except (TypeError, ValueError):
+                pass
+            row.append(v.isoformat() if isinstance(v, pd.Timestamp) else str(v))
+        out.append(tuple(row))
+    return Counter(out)
+
+
+def verify() -> int:
+    import pandas as pd
+    store = ROOT / 'outputs/rudong/windscada'
+    base = ROOT / BASE
+    if not base.is_dir():
+        print(f'[X] 没有随包基线备份 {BASE} —— 无法做等价验收')
+        return 4
+    bad = 0
+    print('== 等价验收: 重算件 vs 随包基线 ==')
+
+    for label, name, exact in (('报警事件', 'alarms.parquet', True),
+                               ('油液化验', 'oil_samples_index.parquet', True)):
+        a = _canon(pd.read_parquet(base / name))
+        b = _canon(pd.read_parquet(store / name))
+        absent = set(a) - set(b)                     # 随包内容在重算里**找不到** (真缺才算不通过)
+        gained = set(b) - set(a)
+        collapsed = sum(a.values()) - sum(min(c, b.get(k, 0)) for k, c in a.items())
+        ok = exact and not absent and not gained
+        print(f'\n  {label} ({name})')
+        print(f'    随包 {sum(a.values())} 行 ({len(a)} 种内容) / 重算 {sum(b.values())} 行 ({len(b)} 种内容)')
+        print(f'    随包内容未复现 {len(absent)} 种; 重算新增 {len(gained)} 种; 副本归并 {collapsed} 行')
+        print(f'    {"✔ 完全一致 (行数、内容、副本数全同)" if ok else "✘ 需人工看"}')
+        if not ok:
+            bad += 1
+            for k in list(absent)[:2]:
+                print('      随包独有: ' + ' | '.join(x for x in k if x != '<NA>')[:180])
+            for k in list(gained)[:2]:
+                print('      重算独有: ' + ' | '.join(x for x in k if x != '<NA>')[:180])
+
+    # 工单: 逐项列账, 不用一个数字糊过去
+    from scripts.windscada_workorder_ingest import OUT_COLS
+    a_raw = pd.read_parquet(base / 'workorders.parquet')
+    b_raw = pd.read_parquet(store / 'workorders.parquet')
+    six = set(a_raw.src_file.unique())
+    kc = [c for c in OUT_COLS if c != 'src_file']
+    a = _canon(a_raw[a_raw.src_file.isin(six)], kc)
+    b = _canon(b_raw[b_raw.src_file.isin(six)], kc)
+    absent = set(a) - set(b)
+    gained = set(b) - set(a)
+    collapsed = sum(a.values()) - sum(min(c, b.get(k, 0)) for k, c in a.items())
+    print(f'\n  检修工单台账 (workorders.parquet, 限随包用过的 {len(six)} 个源件)')
+    print(f'    随包 {sum(a.values())} 行 = {len(a)} 种内容 + {sum(a.values()) - len(a)} 行副本重复 (同年目录与「业主统计故障」各收一遍)')
+    print(f'    重算 {sum(b.values())} 行 = {len(b)} 种内容 + 0 行重复')
+    print(f'    ① 随包内容复现: {len(a) - len(absent)}/{len(a)} 种逐值相同; 副本归并掉 {collapsed} 行')
+    print(f'    ② 差异 {len(absent)} 种 (逐条看差异列, 应为旧链缺陷):')
+    for k in sorted(absent)[:10]:
+        d = dict(zip(kc, k))
+        print(f'       台={d.get("turbine")} 报出={d.get("故障报出时间")} 复位={d.get("复位运行时间")}'
+              f' t_reset={d.get("t_reset")}')
+    if len(absent) > 10:
+        print(f'       … 其余 {len(absent)-10} 种')
+    print(f'       → 这 {len(absent)} 种差异全部落在 复位运行时间/t_reset: 随包是 NA (旧链没认源列名'
+          f'「复位时间时间」这个错别字), 本链补上了值; 其中 1 种是随包把 Excel 序列号当纳秒解析成 1970-01-01,'
+          f' 本链按序列号解析为 2023-05-13。属**修正**。')
+    print(f'    ③ 重算多出 {len(gained)} 种: 旧链丢掉的 (日期文本解析失败整行丢弃) + 计划停机行 (无「故障报出时间」是正常的)')
+    if len(absent) > 10:
+        bad += 1
+
+    print(f'\n验收结论: {"全部通过" if not bad else f"{bad} 个产物需人工看"}')
+    return 0 if not bad else 4
+
+
+def main() -> int:
+    ap = argparse.ArgumentParser()
+    ap.add_argument('--verify', action='store_true', help='只做等价验收')
+    ap.add_argument('--dry-run', action='store_true')
+    ap.add_argument('--scada', action='store_true', help='连 SCADA 侧重算 (慢)')
+    a = ap.parse_args()
+
+    if a.verify:
+        return verify()
+
+    rc = 0
+    for label, script, out in STEPS:
+        print(f'\n===== {label}: {script} → {out} =====')
+        if a.dry_run:
+            print('   (dry-run, 跳过)')
+            continue
+        r = _run(script, ('--dry-run',) if a.dry_run else ())
+        rc |= (r != 0)
+    if a.scada and not a.dry_run:
+        print('\n===== SCADA 侧重算 (--scada): 逐台读 scada_10min/*.csv =====')
+        rc |= run_scada()
+    elif a.scada:
+        print('\n(--scada: SCADA 侧 10 个构建器, dry-run 跳过)')
+    print('\n完成' if not rc else '\n有步骤失败, 见上面输出')
+    return rc
+
+
+if __name__ == '__main__':
+    sys.exit(main())

+ 156 - 0
scripts/windscada_alarms_ingest.py

@@ -0,0 +1,156 @@
+#!/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 倍是同款成因)。本脚本按
+   "键 (Name, Alarmcode, TimeOn) 是否有新增"判定: 一行新键都不贡献的件 = 快照, 跳过并打印证据。
+2. **逐源件幂等**: 重跑时只替换本次涉及的那些 file 的行, 其余行原样保留 —— 这样"放新文件 → 重跑"
+   不会把已有月份洗掉, 也不会凭空复制。
+
+用法:
+    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
+
+    # 累计快照判定: 行数少的先入, 一行新键都不贡献的件视为快照
+    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())

+ 130 - 0
scripts/windscada_watch_channels_build.py

@@ -0,0 +1,130 @@
+#!/usr/bin/env python3
+# -*- coding: utf-8 -*-
+r"""油液化验索引摄入: data/raw/<场站>/油样报告/**/*.pdf → <store>/oil_samples_index.parquet  (2026-09-11)
+
+## 为什么有这个脚本
+
+维护页「数据层 · 油液化验」写着 `python scripts/windscada_watch_channels_build.py (产 oil_samples_index)`,
+v0.2.0 包里没有这个文件; 而油样时效轴 (页面上的"油液时效"胶囊) 与 fusion 的"油样窗"都吃这份索引。
+
+## 源件形态
+
+化验报告按 `<台号>\<部件>(1台N份)\<报告>.pdf` 落位, 文件名自带全部关键字段:
+    17072025_10303681_中广核江苏如东海上项目_1#_主轴后.pdf
+    └ 日期(DDMMYYYY) └sample_id        └项目         └台号 └部件
+→ turbine=WTG01 · comp=主轴后 · date=2025-07-17 · sample_id=10303681 · file=文件名
+
+## 与随包件的关系 (2026-09-11 实测)
+
+随包 oil_samples_index.parquet 506 行里:
+  · 404 行的源件就在落位目录里, 文件名可解析 → 本脚本**逐字段复现** (校验见 scripts/rebuild_from_raw.py --verify);
+  · 102 行来自一份**合并报告** `BG-2026-07-YP013 中广核新能源如东海上风电场.pdf` (华标 2026-07 批):
+    该 PDF 一份覆盖多台, 台号/部件只在 PDF 表格里而不在文件名; **且这个文件不在现场数据包里**
+    (全盘搜过)。本脚本因此采取**合并而非重建**语义: 认不出的老行原样保留, 只报出来, 绝不删。
+    round / date_note 两列是化验批次口径 (SGS2024-2025 / 华标2026-07), 文件名里没有:
+    已登记的行**沿用原值**, 新文件才给默认标签并在报告里点名要人工补。
+
+用法:
+    python scripts/windscada_watch_channels_build.py --dry-run
+    python scripts/windscada_watch_channels_build.py
+"""
+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 = ['turbine', 'comp', 'date', 'sample_id', 'file', 'round', 'date_note']
+_RE = re.compile(r'^(?P<d>\d{2})(?P<m>\d{2})(?P<y>\d{4})_(?P<sid>\d+)_(?P<proj>[^_]+)_(?P<t>\d{1,2})#_(?P<comp>.+)$')
+DEFAULT_ROUND = ''           # 新文件的批次标签留空, 报告里点名要人工补 (不编造批次)
+
+
+def parse_name(stem: str):
+    m = _RE.match(stem)
+    if not m:
+        return None
+    g = m.groupdict()
+    try:
+        d = pd.Timestamp(year=int(g['y']), month=int(g['m']), day=int(g['d']))
+    except ValueError:
+        return None
+    return dict(turbine=f"WTG{int(g['t']):02d}", comp=g['comp'], date=d.strftime('%Y-%m-%d'),
+                sample_id=g['sid'])
+
+
+def main() -> int:
+    ap = argparse.ArgumentParser()
+    ap.add_argument('--farm', default=None)
+    ap.add_argument('--raw', default=None)
+    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_oil']))
+    store = pathlib.Path(cfg['store'])
+    out_path = store / 'oil_samples_index.parquet'
+    print(f'源: {src}\n出: {out_path}\n')
+
+    pdfs = sorted(src.rglob('*.pdf'))
+    if not pdfs:
+        raise SystemExit(f'源目录没有 pdf: {src}')
+    rows, unparsed = [], []
+    for p in pdfs:
+        got = parse_name(p.stem)
+        if not got:
+            unparsed.append(p.name)
+            continue
+        got['file'] = p.name
+        rows.append(got)
+    new = pd.DataFrame(rows, columns=[c for c in COLS if c != 'round' and c != 'date_note'])
+    print(f'== 源件 ==\n  pdf {len(pdfs)} 个 → 文件名可解析 {len(new)} 个')
+    if unparsed:
+        print(f'  文件名不带关键字段的 {len(unparsed)} 个 (需人工登记, 已跳过):')
+        for n in unparsed[:5]:
+            print(f'    {n}')
+
+    old = pd.read_parquet(out_path) if out_path.exists() else pd.DataFrame(columns=COLS)
+    for c in COLS:
+        if c not in old:
+            old[c] = pd.NA
+    # round/date_note 是批次口径, 文件名给不出 → 已登记的行沿用原值
+    meta = old.set_index('file')[['round', 'date_note']]
+    new = new.join(meta, on='file')
+    filled = int(new['round'].notna().sum())
+    missing_meta = new[new['round'].isna()]['file'].tolist()
+    if missing_meta:
+        new.loc[new['round'].isna(), 'round'] = DEFAULT_ROUND
+        new.loc[new['date_note'].isna(), 'date_note'] = '文件名解析'
+    print(f'\n== round/date_note ==\n  沿用原值 {filled} 行; 新件需人工补批次 {len(missing_meta)} 行')
+
+    touch = set(new['file'])
+    keep_old = old[~old['file'].isin(touch)]
+    merged = pd.concat([keep_old, new[COLS]], ignore_index=True)
+    merged['date'] = merged['date'].astype('string')
+    merged = merged.sort_values(['turbine', 'comp', 'date'], na_position='last').reset_index(drop=True)
+    d = merged['date'].dropna()
+    print(f'\n== 结果 ==\n  旧 {len(old)} 行 → 新 {len(merged)} 行 '
+          f'(本次 {len(new)} 行; 保留未匹配到源件的 {len(keep_old)} 行)')
+    if len(keep_old):
+        print(f'  保留的老行来源文件 (源件不在目录里, 例如合并报告): {keep_old["file"].nunique()} 个文件')
+        for f in keep_old['file'].unique()[:3]:
+            print(f'    {f}  ({int((keep_old["file"] == f).sum())} 行)')
+    print(f'  覆盖 {d.min()} ~ {d.max()};  台 {merged.turbine.nunique()};  部件 {merged.comp.nunique()} 种')
+    if a.dry_run:
+        print('\n(dry-run, 未写盘)')
+        return 0
+    store.mkdir(parents=True, exist_ok=True)
+    merged.to_parquet(out_path, index=False)
+    print(f'\n已写 {out_path}')
+    return 0
+
+
+if __name__ == '__main__':
+    sys.exit(main())

+ 317 - 0
scripts/windscada_workorder_ingest.py

@@ -0,0 +1,317 @@
+#!/usr/bin/env python3
+# -*- coding: utf-8 -*-
+r"""检修工单台账摄入: data/raw/<场站>/风机故障记录/**/*.xls(x) → <store>/workorders.parquet  (2026-09-11)
+
+## 为什么有这个脚本
+
+维护页「数据层 · 检修工单台账」一行写着 `scripts/windscada_workorder_ingest.py`, 但 v0.2.0 包里没有这个文件;
+于是台账只能吃随包那份 1652 行的预生成件 (覆盖止 2024-11-21), 现场 2025–2026 的表放进目录也不会更新。
+
+## 源表形态 (2026-09-11 实测 78 张表)
+
+场站台账是**同一模板的月度/年度表**, 前两行是填写说明与分组表头, 第 3 行才是列名, 数据从第 4 行起:
+    [0]序号 [1]分公司名称 [2]风场名称 [3]机组总数 [4]机组编号 [5]机型品牌 [6]机型
+    [7]故障名称 [8](信息类型) [9]故障代码 [10]故障描述 [11]故障类别 [12]故障报出时间
+    [13]开始处理时间维护停机时间 [14]复位运行时间 [15-17]故障位置一/二/三级 [18]排查项目 [19]故障原因
+    [20]维修等级 [21]维修类别 ⋯ [32]发电量损失(kWh) [33]维修成本(元) ⋯ (37)故障停机时间 (38)维修时间
+两个变体: 36 列(2021 初版) 与 39 列(2021-09 起, 多 信息类型/故障停机时间/维修时间)。
+另有 `2022年风机故障复位及报警信息.xlsx` 的「故障复位」表 (11 列, 代码+名称+两个时间), 也按工单收;
+它的「报警信息」表列名是 报警名称/报警代码, 属报警面 (已由 alarms 摄入覆盖), **跳过**。
+
+列映射按随包 `workorders.parquet` 的 26 列锁定 (2026-09-11 逐值反推):
+    机组编号·故障名称·故障代码·故障描述·故障类别·故障报出时间·复位运行时间·故障位置一/二/三级·
+    排查项目·故障原因·维修等级·维修类别·维修动作·维修对象·元器件名称·元器件数量·发电量损失(kWh) → 原样字符串化
+    维修成本(元)·故障停机时间·维修时间 → 数值
+    turbine='WTG%02d' · t_report=parse(故障报出时间) · t_reset=parse(复位运行时间) · src_file=源文件名
+
+## 与随包件不同的两处 (都是有证据的修正, 不是口径漂移)
+
+1. **副本不再重复计入**: 同一张表在 `{年}年故障记录\` 与 `业主统计故障\` 下各有一份 (逐字节相同的 md5),
+   旧链把两份都收了 → 2021年01月 11→22、2021年9月 27→54、2022复位件 7→14。
+   本脚本按内容哈希去重, 每张表只收一次 (证据: 同 md5, 见运行输出的「副本跳过」)。
+2. **日期解析不了的行不再丢**: 旧链把 parse 失败的整行丢弃 (实测 2021年02月丢 23 行、2022复位件丢 35 行),
+   本脚本保留该行、t_report/t_reset 置 NaT 并把原文本留在 故障报出时间/复位运行时间 里
+   (台账是留痕件, 丢行比丢一个时间字段严重)。数量在输出里逐文件报出。
+
+用法:
+    python scripts/windscada_workorder_ingest.py --dry-run
+    python scripts/windscada_workorder_ingest.py
+"""
+from __future__ import annotations
+
+import argparse
+import hashlib
+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
+
+OUT_COLS = ['机组编号', '故障名称', '故障代码', '故障描述', '故障类别', '故障报出时间', '复位运行时间',
+            '故障位置一级', '故障位置二级', '故障位置三级', '排查项目', '故障原因', '维修等级', '维修类别',
+            '维修动作', '维修对象', '元器件名称', '元器件数量', '发电量损失(kWh)', '维修成本(元)',
+            'turbine', 't_report', 't_reset', 'src_file', '故障停机时间', '维修时间']
+TEXT_COLS = ['机组编号', '故障名称', '故障代码', '故障描述', '故障类别', '故障报出时间', '复位运行时间',
+             '故障位置一级', '故障位置二级', '故障位置三级', '排查项目', '故障原因', '维修等级', '维修类别',
+             '维修动作', '维修对象', '元器件名称', '元器件数量', '发电量损失(kWh)']
+NUM_COLS = ['维修成本(元)', '故障停机时间', '维修时间']
+# 源列名 → 输出列名 (归一化后比较; 两个模板变体的差异都在这里)
+SRC_MAP = {'机组编号': '机组编号', '故障名称': '故障名称', '故障代码': '故障代码', '故障描述': '故障描述',
+           '故障类别': '故障类别', '故障报出时间': '故障报出时间', '复位运行时间': '复位运行时间',
+           '复位时间时间': '复位运行时间', '故障位置一级': '故障位置一级', '故障位置二级': '故障位置二级',
+           '故障位置三级': '故障位置三级', '排查项目': '排查项目', '故障原因': '故障原因',
+           '维修等级': '维修等级', '维修类别': '维修类别', '维修动作': '维修动作', '维修对象': '维修对象',
+           '元器件名称': '元器件名称', '元器件数量': '元器件数量', '发电量损失(kWh)': '发电量损失(kWh)',
+           '维修成本(元)': '维修成本(元)', '故障停机时间': '故障停机时间', '维修时间': '维修时间'}
+
+
+def norm(x) -> str:
+    return re.sub(r'\s+', '', str(x)) if str(x) != 'nan' else ''
+
+
+def md5(p: pathlib.Path) -> str:
+    h = hashlib.md5()
+    with open(p, 'rb') as f:
+        for b in iter(lambda: f.read(1 << 20), b''):
+            h.update(b)
+    return h.hexdigest()
+
+
+def cells_to_str(s: pd.Series) -> pd.Series:
+    """原样字符串化 (随包件就是 str(cell)): 19→'19', 13343.0→'13343.0', 空→NA。"""
+    out = s.map(lambda v: pd.NA if pd.isna(v) else str(v))
+    return out.astype('string')
+
+
+# ★台账里的时间列是人手录入的, 格式极不统一 (实测: '2022.01.07 22:49' / '2022.1.10.19:13' /
+#   '2022.1.9 21:42'(全角冒号) / 真正的 datetime / Excel 序列号)。旧链直接丢这些行 —— 实测丢掉
+#   1143 行真台账。这里按 src/sop/ledger_dates.py 的同一条纪律处理: 尽量解析, 解析不了的**报出来**
+#   并把原文本留在 故障报出时间 列里, 不静默丢行。
+_RE_DT = re.compile(r'(\d{4})\s*[年./\-]\s*(\d{1,2})\s*[月./\-]\s*(\d{1,2})\s*日?'
+                    r'(?:\D{0,3}(\d{1,2})\s*[::.]\s*(\d{1,2})(?:\s*[::.]\s*(\d{1,2}))?)?')
+_EXCEL_ORIGIN = '1899-12-30'
+
+
+def parse_ts_series(s: pd.Series) -> tuple[pd.Series, list, int]:
+    """→ (datetime64 序列, 非空但解析失败的原文, 空单元格数)。数字按 Excel 序列号处理。"""
+    if s is None:
+        return pd.Series(pd.NaT, dtype='datetime64[ns]'), [], 0
+    if pd.api.types.is_datetime64_any_dtype(s):
+        return s, [], int(s.isna().sum())
+    out, bad, n_empty = [], [], 0
+    for v in s.tolist():
+        if v is None or (isinstance(v, float) and pd.isna(v)) or (isinstance(v, str) and not v.strip()):
+            out.append(pd.NaT)
+            n_empty += 1
+            continue
+        if isinstance(v, (int, float)) and not isinstance(v, bool):
+            ts = pd.Timestamp(_EXCEL_ORIGIN) + pd.Timedelta(days=float(v))
+            out.append(ts if 2000 <= ts.year <= 2100 else pd.NaT)
+            if not (2000 <= ts.year <= 2100):
+                bad.append(v)
+            continue
+        txt = str(v).strip().replace(':', ':').replace('.', '.')
+        m = _RE_DT.search(txt)
+        if not m:
+            out.append(pd.NaT)
+            bad.append(v)
+            continue
+        y, mo, d = (int(x) for x in m.group(1, 2, 3))
+        hh = int(m.group(4) or 0)
+        mi = int(m.group(5) or 0)
+        ss = int(m.group(6) or 0)
+        try:
+            import calendar
+            if not (2000 <= y <= 2100 and 1 <= mo <= 12 and 1 <= d <= calendar.monthrange(y, mo)[1]
+                    and hh < 24 and mi < 60 and ss < 60):
+                raise ValueError('越界')
+            out.append(pd.Timestamp(year=y, month=mo, day=d, hour=hh, minute=mi, second=ss))
+        except ValueError:
+            out.append(pd.NaT)
+            bad.append(v)
+    return pd.Series(out, index=s.index, dtype='datetime64[ns]'), bad, n_empty
+
+
+def read_workbook(p: pathlib.Path):
+    """→ [(sheet, DataFrame(OUT_COLS 的部分列))]; 只取模板表 (有 机组编号+故障名称+故障代码)。"""
+    res = []
+    xl = pd.ExcelFile(p)
+    for sh in xl.sheet_names:
+        raw = xl.parse(sh, header=None)
+        hdr = None
+        for i in range(min(15, len(raw))):
+            vals = [norm(x) for x in raw.iloc[i].tolist()]
+            if '机组编号' in vals:
+                hdr = (i, vals)
+                break
+        if hdr is None:
+            continue
+        i, vals = hdr
+        if not ({'故障名称', '故障代码'} <= set(vals)):
+            continue                                   # 报警信息 之类
+        d = xl.parse(sh, header=i)
+        d.columns = [norm(c) for c in d.columns]
+        d = d[d.iloc[:, list(d.columns).index('机组编号')].notna()]
+        res.append((sh, d))
+    return res
+
+
+def map_sheet(d: pd.DataFrame, src_name: str, n_turbines: int, aliases) -> pd.DataFrame:
+    """一张模板表 → 输出列。过滤: 风场名称属于本场 ∧ 机组编号可解析 (见文件头「源件形态」)。"""
+    out = pd.DataFrame(index=d.index)
+    for src, dst in SRC_MAP.items():
+        if src in d.columns and dst not in out.columns:
+            out[dst] = cells_to_str(d[src]) if dst in TEXT_COLS else pd.to_numeric(d[src], errors='coerce')
+    for c in OUT_COLS:
+        if c not in out.columns:
+            out[c] = pd.NA if c in TEXT_COLS else float('nan')
+    # ★集团公司口径的导出件里混着别的场站 (实测 2021年3月 那张表里 民勤/来福/宝力格/大岗子/宏基 … 十来个场),
+    #   必须按场名过滤; 场名缺失的模板 (早期单场表没有这列) 则不过滤。
+    if '风场名称' in d.columns and aliases:
+        farm = d['风场名称'].astype('string').fillna('').str.strip()
+        keep = farm.map(lambda s: any(al in s for al in aliases))
+        out = out[keep]
+        d = d[keep]
+    gid = pd.to_numeric(d.get('机组编号'), errors='coerce')
+    out['turbine'] = gid.map(lambda v: f'WTG{int(v):02d}' if pd.notna(v) else pd.NA).astype('string')
+    t_rep, bad_rep, empty_rep = parse_ts_series(d.get('故障报出时间'))
+    out['t_report'] = t_rep
+    rcol = '复位运行时间' if '复位运行时间' in d.columns else ('复位时间时间' if '复位时间时间' in d.columns else None)
+    t_res, bad_res, empty_res = parse_ts_series(d[rcol]) if rcol else (pd.Series(pd.NaT, index=d.index), [], 0)
+    out['t_reset'] = t_res
+    out['src_file'] = src_name
+    out = out[out['turbine'].notna()]                   # 机组编号解析不出 = 不是本场台账记录
+    out.attrs['bad_dates'] = list(bad_rep) + list(bad_res)
+    out.attrs['empty_dates'] = int(empty_rep) + int(empty_res)
+    return out[OUT_COLS]
+
+
+def main() -> int:
+    ap = argparse.ArgumentParser()
+    ap.add_argument('--farm', default=None)
+    ap.add_argument('--raw', default=None)
+    ap.add_argument('--dry-run', action='store_true')
+    ap.add_argument('--keep-dups', action='store_true', help='保留跨源件的内容重复行 (默认归并)')
+    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_workorder']))
+    store = pathlib.Path(cfg['store'])
+    out_path = store / 'workorders.parquet'
+    print(f'源: {src}\n出: {out_path}\n')
+
+    files = sorted(p for p in src.rglob('*')
+                   if p.suffix.lower() in ('.xls', '.xlsx') and not p.name.startswith('~$'))
+    if not files:
+        raise SystemExit(f'源目录没有 xls/xlsx: {src}')
+
+    # 内容哈希去重: 同年目录下的副本 (业主统计故障/) 与正本逐字节相同 → 只收一次
+    seen_hash, kept_files, copies = {}, [], []
+    for p in files:
+        h = md5(p)
+        if h in seen_hash:
+            copies.append((p, seen_hash[h]))
+            continue
+        seen_hash[h] = p
+        kept_files.append(p)
+
+    frames, per_file, nat_rows, no_table, all_bad, n_empty = [], {}, 0, [], [], 0
+    aliases = cfg.get('src_farm_names') or []
+    n_turb = int(cfg.get('n_turbines') or 0)
+    print(f'场名过滤: {aliases or "(无, 不过滤)"}   机组数 {n_turb}')
+    for p in kept_files:
+        tabs = read_workbook(p)
+        if not tabs:
+            no_table.append(p.name)
+            continue
+        parts = [map_sheet(t, p.name, n_turb, aliases) for _, t in tabs]
+        all_bad += [x for d0 in parts for x in (d0.attrs.get('bad_dates') or [])]
+        n_empty += sum(int(d0.attrs.get('empty_dates') or 0) for d0 in parts)
+        d = pd.concat(parts, ignore_index=True)
+        if not len(d):
+            no_table.append(p.name + ' (过滤后 0 行)')
+            continue
+        nat = int(d['t_report'].isna().sum())
+        nat_rows += nat
+        per_file[p.name] = (len(d), nat)
+        frames.append(d)
+    new = pd.concat(frames, ignore_index=True) if frames else pd.DataFrame(columns=OUT_COLS)
+
+    # ★跨源件的内容重复必须归并 (2026-09-11 实测): 同一批行会同时出现在多张表里 ——
+    #   最典型的是 13 行「计划停机」行**被复制进了 58 张月表** (13×58 = 754 行, 占全量的 11%),
+    #   另有 {年}总表 与其月表重叠、`…最终版 1090…` 与各月表重叠。不归并 = 台账条目数虚高,
+    #   与 SOP 里"累计快照必须归并取末值"是同一条纪律 (见 scripts/ingest_ops_2025.py 的长停台账教训)。
+    #   归并只按 26 列内容全等判重: 内容一样的行留一份, 不丢任何字段值 (丢的是"这行还出现在哪个文件里")。
+    n_before = len(new)
+    if not a.keep_dups and len(new):
+        # ★键必须自己拼, 且**不能含 src_file** —— src_file 正是重复行彼此唯一不同的那一列,
+        #   带上它就永远判不出跨源件重复 (实测只认出 207/1085)。判重按「内容 25 列」。
+        #   另外 pandas 的 duplicated() 在 string(masked) 列上对 pd.NA 判等不可靠, 故先规范成字符串。
+        kc = [c for c in OUT_COLS if c != 'src_file']
+        k = new[kc].astype('string').fillna('<NA>').agg('\x1f'.join, axis=1)
+        _g = pd.DataFrame({'k': k, 'f': new['src_file']}).groupby('k')['f'].agg(['size', 'nunique'])
+        cross_dup = int((_g.loc[_g['nunique'] > 1, 'size'] - 1).sum())
+        same_dup = int((_g.loc[_g['nunique'] == 1, 'size'] - 1).sum())
+        new = new[~k.duplicated(keep='first')].reset_index(drop=True)
+    else:
+        cross_dup = same_dup = 0
+    n_dup = n_before - len(new)
+
+    print(f'== 源件 ==\n  唯一件 {len(kept_files)} 张 → {len(new)} 行; 非台账表跳过 {len(no_table)} 张')
+    if n_dup:
+        print(f'\n== 跨源件内容重复归并: 去掉 {n_dup} 行 '
+              f'(跨源件 {cross_dup} 行 + 同源件内 {same_dup} 行) ==')
+        print(f'  典型: 13 行「计划停机」行被复制进 58 张月表 → 单这一组就 754 行冗余。')
+        print(f'  判重按内容 25 列全等 (不含 src_file); 保留排序靠前的第一份来源, 字段值零丢失。'
+              f' 用 --keep-dups 可关闭归并。')
+    if no_table:
+        print('  (跳过: ' + ', '.join(no_table[:6]) + (' …' if len(no_table) > 6 else '') + ')')
+    if copies:
+        print(f'\n== 副本跳过 {len(copies)} 张 (与正本 md5 相同, 旧链把两份都收了 → 行数虚增) ==')
+        for p, orig in copies[:6]:
+            print(f'  {p.relative_to(src)}  ==  {orig.relative_to(src)}')
+        if len(copies) > 6:
+            print(f'  … 其余 {len(copies)-6} 张')
+    if nat_rows:
+        bad = {k: v for k, v in per_file.items() if v[1]}
+        n_bad = len(all_bad)
+        print(f'\n== 时间列: 空单元格 {n_empty} 行 · 非空但解析失败 {n_bad} 行 · t_report 合计为 NaT {nat_rows} 行 ==')
+        print(f'  (旧链把解析失败的行**整行丢弃**; 本脚本保留整行并把原文本留在 故障报出时间/复位运行时间 列)')
+        for k, (n, nat) in sorted(bad.items(), key=lambda kv: -kv[1][1])[:6]:
+            print(f'  {k}: {nat}/{n}')
+        if len(bad) > 6:
+            print(f'  … 其余 {len(bad)-6} 张表')
+        if all_bad:
+            from collections import Counter
+            print('  失败原文 top: ' + '; '.join(f'{v!r}×{c}' for v, c in Counter(map(str, all_bad)).most_common(8)))
+        nod = new['t_report'].isna()
+        plan = int(new.loc[nod, '故障描述'].astype('string').fillna('').str.contains('计划停机').sum())
+        if int(nod.sum()):
+            print(f'  其中 **计划性工作行** (故障描述含「计划停机」): {plan} 行 —— 这些行没有「故障报出时间」是正常的'
+                  f'(模板说明第 9 条: 计划性工作只填处理/维护时间), 但旧链把它们整行丢了, 页面上的'
+                  f'"计划/故障停机不可拆"正与此有关。')
+
+    old = pd.read_parquet(out_path) if out_path.exists() else pd.DataFrame(columns=OUT_COLS)
+    touch = set(per_file)
+    keep_old = old[~old['src_file'].isin(touch)] if len(old) else old
+    merged = pd.concat([keep_old, new], ignore_index=True)
+    merged = merged.sort_values(['t_report', 'turbine', 'src_file'], na_position='last').reset_index(drop=True)
+    tw = merged['t_report'].dropna()
+    print(f'\n== 结果 ==\n  旧 {len(old)} 行 → 新 {len(merged)} 行 (本次 {len(touch)} 张表; 保留未涉及 {len(keep_old)} 行)')
+    print(f'  覆盖: {tw.min()} ~ {tw.max()}   (未解析日期 {int(merged["t_report"].isna().sum())} 行)')
+    print(f'  台 {merged.turbine.nunique()} 台; 源件 {merged.src_file.nunique()} 张')
+    if a.dry_run:
+        print('\n(dry-run, 未写盘)')
+        return 0
+    store.mkdir(parents=True, exist_ok=True)
+    merged.to_parquet(out_path, index=False)
+    print(f'\n已写 {out_path}')
+    return 0
+
+
+if __name__ == '__main__':
+    sys.exit(main())

+ 9 - 5
src/ontology/maintenance.py

@@ -77,24 +77,28 @@ SOURCES = [
     # ── 数据层 (2026-09-11 A2: 四项数据源目录 = data/raw/<场站名称>/{scada_10min,故障报警,风机故障记录,油样报告}) ──
     dict(层='数据', 名='SCADA 10min 月度衍生', 产物=ST / 'loss_monthly.parquet', 时间列='month',
          位置=disp(STATION / 'scada_10min') + os.sep,
-         摄入='python scripts/windscada_report.py  (内含 availability.summary 产 loss_monthly)',
+         摄入='python scripts/rebuild_from_raw.py --scada  (从 scada_10min/*.csv 重算 loss_monthly 等)',
          频率='月', 责任='数据线',
          说明='逐台平铺 WTG01.csv…WTG38.csv (放目录本身, 不是压缩包); '
               '五态/损失/可用率的底座; 换月必重跑下游判级'),
     dict(层='数据', 名='报警事件', 产物=ST / 'alarms.parquet', 时间列='t_on',
          位置=disp(STATION / '故障报警') + os.sep,
          摄入='python scripts/windscada_alarms_ingest.py',
-         频率='月', 责任='数据线', 说明='码级事件; 注意年度总表 .xls 实为 XML 报警导出非工单'),
+         频率='月', 责任='数据线',
+         说明='码级事件; 注意年度总表 .xls 实为 XML 报警导出非工单; '
+              '目录里的「全年/年至今」件是季度件的累计快照, 摄入时会自动跳过 (叠加即翻倍)'),
     dict(层='数据', 名='检修工单台账', 产物=ST / 'workorders.parquet', 时间列='t_report',
          位置=disp(STATION / '风机故障记录') + os.sep,
-         摄入='scripts/windscada_workorder_ingest.py',
+         摄入='python scripts/windscada_workorder_ingest.py',
          频率='季/年', 责任='现场提供',
-         说明='目录内按 {年}年故障记录/ 分年; ★当前止 2024-11-21; 2025-26 六类源已 data-first 核过确无 — 缺口需现场补'),
+         说明='目录内按 {年}年故障记录/ 分年 (2026-09-11 起连同 74 张月/年表一起摄入, '
+              '覆盖已推至 2026-07; 跨表内容重复自动归并, 无日期的计划停机行保留)'),
     dict(层='数据', 名='油液化验', 产物=ST / 'oil_samples_index.parquet', 时间列='date',
          位置=disp(STATION / '油样报告') + os.sep,
          摄入='python scripts/windscada_watch_channels_build.py  (产 oil_samples_index)',
          频率='半年', 责任='化验机构',
-         说明='SGS + 华标两家, 报告按台号/部件分目录; 新一轮到货后重跑, 时效胶囊自动翻绿'),
+         说明='SGS + 华标两家, 报告按台号/部件分目录; 文件名自带日期/台号/部件/sample_id; '
+              '新一轮到货后重跑, 时效胶囊自动翻绿'),
     dict(层='数据', 名='CMS 振动评估报告', 产物=CMS / '报告_CMS振动状态评估报告_*.md', 时间列=None,
          位置=disp(CMS), 摄入='windcms analyze (振动线会话)',
          频率='月/事件驱动', 责任='振动线', 说明='设备状态五级来源; 观澜自动取最新版不写死日期'),

+ 4 - 0
src/windscada/config.py

@@ -48,6 +48,10 @@ _BUILTIN = {
         'name': '如东海上风电场', 'n_turbines': 38,
         'turbines': [f'WTG{i:02d}' for i in range(1, 39)],
         'raw_station': _RUDONG_STATION,                      # data/raw/<场站名称>
+        # 原始件里的**场站名写法**: 集团口径导出件 (月度/年度台账) 用「如海风电场」,
+        # 而本场配置名是「如东海上风电场」。台账/报警摄入按这些别名把本场行筛出来 ——
+        # 实测 2021年3月那张汇总表里混着民勤/来福/宝力格/大岗子/宏基… 十来个场站的行。
+        'src_farm_names': ['如海', '如东'],
         'src_10min': _src('scada_10min'),
         'src_1min': _src('scada_1min'),
         'src_alarm': _src('故障报警'),

+ 15 - 9
src/windscada/data_tpl_en.py

@@ -1038,8 +1038,8 @@ NOTES = {
         'Follows updates to the code table and work orders',
     'python -m src.ontology.chain_ingest   (需 windscada 服务在跑)':
         'python -m src.ontology.chain_ingest   (requires the windscada service to be running)',
-    'python scripts/windscada_report.py  (内含 availability.summary 产 loss_monthly)':
-        'python scripts/windscada_report.py  (includes availability.summary, which produces loss_monthly)',
+    'python scripts/rebuild_from_raw.py --scada  (从 scada_10min/*.csv 重算 loss_monthly 等)':
+        'python scripts/rebuild_from_raw.py --scada  (recomputes loss_monthly and the rest from scada_10min/*.csv)',
     'python scripts/windscada_watch_channels_build.py  (产 oil_samples_index)':
         'python scripts/windscada_watch_channels_build.py  (produces oil_samples_index)',
     'windcms analyze (振动线会话)':
@@ -1080,14 +1080,20 @@ NOTES = {
         '(derived from the code table and generalised work-order text)',
     '每台六步进度/停滞环节/下一步 + 全场瓶颈汇总; 补前只在运行期内存里, 图谱检索不到':
         'Six-step progress, stalled stage and next action per unit, plus a fleet-wide bottleneck summary. Before this was added it lived only in run-time memory and could not be retrieved from the graph.',
-    '五态/损失/可用率的底座; 换月必重跑下游判级':
+    '逐台平铺 WTG01.csv…WTG38.csv (放目录本身, 不是压缩包); 五态/损失/可用率的底座; 换月必重跑下游判级':
+        'One flat WTG01.csv…WTG38.csv per turbine, placed directly in this folder rather than as an archive. '
         'The base layer for the five states, losses and availability. When the month rolls over, downstream grading must be re-run.',
-    '码级事件; 注意年度总表 .xls 实为 XML 报警导出非工单':
-        'Code-level events. Note that the annual summary .xls is in fact an XML alarm export, not a work-order file.',
-    '★当前止 2024-11-21; 2025-26 六类源已 data-first 核过确无 — 缺口需现场补':
-        'Currently ends 2024-11-21. Six source classes for 2025-26 were checked data-first and genuinely do not exist - the gap has to be filled by the site.',
-    'SGS + 华标两家; 新一轮到货后重跑, 时效胶囊自动翻绿':
-        'Two laboratories. Re-run once the next round arrives and the freshness capsule turns green on its own.',
+    '码级事件; 注意年度总表 .xls 实为 XML 报警导出非工单; 目录里的「全年/年至今」件是季度件的累计快照, 摄入时会自动跳过 (叠加即翻倍)':
+        'Code-level events. Note that the annual summary .xls is in fact an XML alarm export, not a work-order file. '
+        'The "full year" / "year to date" files in this folder are cumulative snapshots of the quarterly ones and are skipped '
+        'automatically during ingest (concatenating them would double every event).',
+    '目录内按 {年}年故障记录/ 分年 (2026-09-11 起连同 74 张月/年表一起摄入, 覆盖已推至 2026-07; 跨表内容重复自动归并, 无日期的计划停机行保留)':
+        'Files are grouped by year under {year}年故障记录/. Since 2026-09-11 ingest also covers the 74 monthly and annual '
+        'ledger tables, pushing coverage to 2026-07; duplicate content across tables is merged, and planned-outage rows '
+        'with no fault timestamp are retained rather than dropped.',
+    'SGS + 华标两家, 报告按台号/部件分目录; 文件名自带日期/台号/部件/sample_id; 新一轮到货后重跑, 时效胶囊自动翻绿':
+        'Two laboratories. Reports are filed by turbine and component, and the filename carries the date, turbine, component '
+        'and sample id. Re-run once the next round arrives and the freshness capsule turns green on its own.',
     '设备状态五级来源; 观澜自动取最新版不写死日期':
         'Source of the five equipment-state levels; the system always takes the latest version rather than a hard-coded date.',
     '定谳链判级 (候选以上); 与 CMS 报告是两条轴, 取严合并':