#!/usr/bin/env python3 # -*- coding: utf-8 -*- r"""检修工单台账摄入: data/raw/<场站>/风机故障记录/**/*.xls(x) → /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('').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) # ★产物 = **当前源件的函数** (用户令 2026-09-20「重算要按 data/raw 的最新变化」): 原先 `keep_old` 会把 # "源件已不在盘"的行一直留着(页面照旧按已不存在的台账显示检修记录)。现在只保留仍在盘的源件的行, # 消失的行如实丢掉并报出来(哪几件、多少行、涉及哪些月), 让人分清误删/撤数。 present = {p.name for p in files} gone = sorted(set(old['src_file'].dropna().unique()) - present) if len(old) else [] drop_gone = old[old['src_file'].isin(gone)] if gone else old.iloc[0:0] keep_old = old[~old['src_file'].isin(touch | set(gone))] 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)} 行)') if len(drop_gone): print(f' [!] 有 {len(drop_gone)} 行来自**已不在盘**的台账件 ({len(gone)} 个): ' + '、'.join(f'{n}({int((drop_gone["src_file"] == n).sum())} 行)' for n in gone[:6]) + ' —— 按"产物=当前 data/raw 的函数"口径不再产出; 误删就把文件放回重跑') _m = pd.to_datetime(drop_gone['t_report'], errors='coerce') if _m.notna().any(): print(f' 涉及月份: {_m.dt.strftime("%Y-%m").dropna().value_counts().sort_index().to_dict()}') 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())