windscada_workorder_ingest.py 19 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330
  1. #!/usr/bin/env python3
  2. # -*- coding: utf-8 -*-
  3. r"""检修工单台账摄入: data/raw/<场站>/风机故障记录/**/*.xls(x) → <store>/workorders.parquet (2026-09-11)
  4. ## 为什么有这个脚本
  5. 维护页「数据层 · 检修工单台账」一行写着 `scripts/windscada_workorder_ingest.py`, 但 v0.2.0 包里没有这个文件;
  6. 于是台账只能吃随包那份 1652 行的预生成件 (覆盖止 2024-11-21), 现场 2025–2026 的表放进目录也不会更新。
  7. ## 源表形态 (2026-09-11 实测 78 张表)
  8. 场站台账是**同一模板的月度/年度表**, 前两行是填写说明与分组表头, 第 3 行才是列名, 数据从第 4 行起:
  9. [0]序号 [1]分公司名称 [2]风场名称 [3]机组总数 [4]机组编号 [5]机型品牌 [6]机型
  10. [7]故障名称 [8](信息类型) [9]故障代码 [10]故障描述 [11]故障类别 [12]故障报出时间
  11. [13]开始处理时间维护停机时间 [14]复位运行时间 [15-17]故障位置一/二/三级 [18]排查项目 [19]故障原因
  12. [20]维修等级 [21]维修类别 ⋯ [32]发电量损失(kWh) [33]维修成本(元) ⋯ (37)故障停机时间 (38)维修时间
  13. 两个变体: 36 列(2021 初版) 与 39 列(2021-09 起, 多 信息类型/故障停机时间/维修时间)。
  14. 另有 `2022年风机故障复位及报警信息.xlsx` 的「故障复位」表 (11 列, 代码+名称+两个时间), 也按工单收;
  15. 它的「报警信息」表列名是 报警名称/报警代码, 属报警面 (已由 alarms 摄入覆盖), **跳过**。
  16. 列映射按随包 `workorders.parquet` 的 26 列锁定 (2026-09-11 逐值反推):
  17. 机组编号·故障名称·故障代码·故障描述·故障类别·故障报出时间·复位运行时间·故障位置一/二/三级·
  18. 排查项目·故障原因·维修等级·维修类别·维修动作·维修对象·元器件名称·元器件数量·发电量损失(kWh) → 原样字符串化
  19. 维修成本(元)·故障停机时间·维修时间 → 数值
  20. turbine='WTG%02d' · t_report=parse(故障报出时间) · t_reset=parse(复位运行时间) · src_file=源文件名
  21. ## 与随包件不同的两处 (都是有证据的修正, 不是口径漂移)
  22. 1. **副本不再重复计入**: 同一张表在 `{年}年故障记录\` 与 `业主统计故障\` 下各有一份 (逐字节相同的 md5),
  23. 旧链把两份都收了 → 2021年01月 11→22、2021年9月 27→54、2022复位件 7→14。
  24. 本脚本按内容哈希去重, 每张表只收一次 (证据: 同 md5, 见运行输出的「副本跳过」)。
  25. 2. **日期解析不了的行不再丢**: 旧链把 parse 失败的整行丢弃 (实测 2021年02月丢 23 行、2022复位件丢 35 行),
  26. 本脚本保留该行、t_report/t_reset 置 NaT 并把原文本留在 故障报出时间/复位运行时间 里
  27. (台账是留痕件, 丢行比丢一个时间字段严重)。数量在输出里逐文件报出。
  28. 用法:
  29. python scripts/windscada_workorder_ingest.py --dry-run
  30. python scripts/windscada_workorder_ingest.py
  31. """
  32. from __future__ import annotations
  33. import argparse
  34. import hashlib
  35. import pathlib
  36. import re
  37. import sys
  38. ROOT = pathlib.Path(__file__).resolve().parents[1]
  39. sys.path.insert(0, str(ROOT))
  40. import pandas as pd # noqa: E402
  41. OUT_COLS = ['机组编号', '故障名称', '故障代码', '故障描述', '故障类别', '故障报出时间', '复位运行时间',
  42. '故障位置一级', '故障位置二级', '故障位置三级', '排查项目', '故障原因', '维修等级', '维修类别',
  43. '维修动作', '维修对象', '元器件名称', '元器件数量', '发电量损失(kWh)', '维修成本(元)',
  44. 'turbine', 't_report', 't_reset', 'src_file', '故障停机时间', '维修时间']
  45. TEXT_COLS = ['机组编号', '故障名称', '故障代码', '故障描述', '故障类别', '故障报出时间', '复位运行时间',
  46. '故障位置一级', '故障位置二级', '故障位置三级', '排查项目', '故障原因', '维修等级', '维修类别',
  47. '维修动作', '维修对象', '元器件名称', '元器件数量', '发电量损失(kWh)']
  48. NUM_COLS = ['维修成本(元)', '故障停机时间', '维修时间']
  49. # 源列名 → 输出列名 (归一化后比较; 两个模板变体的差异都在这里)
  50. SRC_MAP = {'机组编号': '机组编号', '故障名称': '故障名称', '故障代码': '故障代码', '故障描述': '故障描述',
  51. '故障类别': '故障类别', '故障报出时间': '故障报出时间', '复位运行时间': '复位运行时间',
  52. '复位时间时间': '复位运行时间', '故障位置一级': '故障位置一级', '故障位置二级': '故障位置二级',
  53. '故障位置三级': '故障位置三级', '排查项目': '排查项目', '故障原因': '故障原因',
  54. '维修等级': '维修等级', '维修类别': '维修类别', '维修动作': '维修动作', '维修对象': '维修对象',
  55. '元器件名称': '元器件名称', '元器件数量': '元器件数量', '发电量损失(kWh)': '发电量损失(kWh)',
  56. '维修成本(元)': '维修成本(元)', '故障停机时间': '故障停机时间', '维修时间': '维修时间'}
  57. def norm(x) -> str:
  58. return re.sub(r'\s+', '', str(x)) if str(x) != 'nan' else ''
  59. def md5(p: pathlib.Path) -> str:
  60. h = hashlib.md5()
  61. with open(p, 'rb') as f:
  62. for b in iter(lambda: f.read(1 << 20), b''):
  63. h.update(b)
  64. return h.hexdigest()
  65. def cells_to_str(s: pd.Series) -> pd.Series:
  66. """原样字符串化 (随包件就是 str(cell)): 19→'19', 13343.0→'13343.0', 空→NA。"""
  67. out = s.map(lambda v: pd.NA if pd.isna(v) else str(v))
  68. return out.astype('string')
  69. # ★台账里的时间列是人手录入的, 格式极不统一 (实测: '2022.01.07 22:49' / '2022.1.10.19:13' /
  70. # '2022.1.9 21:42'(全角冒号) / 真正的 datetime / Excel 序列号)。旧链直接丢这些行 —— 实测丢掉
  71. # 1143 行真台账。这里按 src/sop/ledger_dates.py 的同一条纪律处理: 尽量解析, 解析不了的**报出来**
  72. # 并把原文本留在 故障报出时间 列里, 不静默丢行。
  73. _RE_DT = re.compile(r'(\d{4})\s*[年./\-]\s*(\d{1,2})\s*[月./\-]\s*(\d{1,2})\s*日?'
  74. r'(?:\D{0,3}(\d{1,2})\s*[::.]\s*(\d{1,2})(?:\s*[::.]\s*(\d{1,2}))?)?')
  75. _EXCEL_ORIGIN = '1899-12-30'
  76. def parse_ts_series(s: pd.Series) -> tuple[pd.Series, list, int]:
  77. """→ (datetime64 序列, 非空但解析失败的原文, 空单元格数)。数字按 Excel 序列号处理。"""
  78. if s is None:
  79. return pd.Series(pd.NaT, dtype='datetime64[ns]'), [], 0
  80. if pd.api.types.is_datetime64_any_dtype(s):
  81. return s, [], int(s.isna().sum())
  82. out, bad, n_empty = [], [], 0
  83. for v in s.tolist():
  84. if v is None or (isinstance(v, float) and pd.isna(v)) or (isinstance(v, str) and not v.strip()):
  85. out.append(pd.NaT)
  86. n_empty += 1
  87. continue
  88. if isinstance(v, (int, float)) and not isinstance(v, bool):
  89. ts = pd.Timestamp(_EXCEL_ORIGIN) + pd.Timedelta(days=float(v))
  90. out.append(ts if 2000 <= ts.year <= 2100 else pd.NaT)
  91. if not (2000 <= ts.year <= 2100):
  92. bad.append(v)
  93. continue
  94. txt = str(v).strip().replace(':', ':').replace('.', '.')
  95. m = _RE_DT.search(txt)
  96. if not m:
  97. out.append(pd.NaT)
  98. bad.append(v)
  99. continue
  100. y, mo, d = (int(x) for x in m.group(1, 2, 3))
  101. hh = int(m.group(4) or 0)
  102. mi = int(m.group(5) or 0)
  103. ss = int(m.group(6) or 0)
  104. try:
  105. import calendar
  106. if not (2000 <= y <= 2100 and 1 <= mo <= 12 and 1 <= d <= calendar.monthrange(y, mo)[1]
  107. and hh < 24 and mi < 60 and ss < 60):
  108. raise ValueError('越界')
  109. out.append(pd.Timestamp(year=y, month=mo, day=d, hour=hh, minute=mi, second=ss))
  110. except ValueError:
  111. out.append(pd.NaT)
  112. bad.append(v)
  113. return pd.Series(out, index=s.index, dtype='datetime64[ns]'), bad, n_empty
  114. def read_workbook(p: pathlib.Path):
  115. """→ [(sheet, DataFrame(OUT_COLS 的部分列))]; 只取模板表 (有 机组编号+故障名称+故障代码)。"""
  116. res = []
  117. xl = pd.ExcelFile(p)
  118. for sh in xl.sheet_names:
  119. raw = xl.parse(sh, header=None)
  120. hdr = None
  121. for i in range(min(15, len(raw))):
  122. vals = [norm(x) for x in raw.iloc[i].tolist()]
  123. if '机组编号' in vals:
  124. hdr = (i, vals)
  125. break
  126. if hdr is None:
  127. continue
  128. i, vals = hdr
  129. if not ({'故障名称', '故障代码'} <= set(vals)):
  130. continue # 报警信息 之类
  131. d = xl.parse(sh, header=i)
  132. d.columns = [norm(c) for c in d.columns]
  133. d = d[d.iloc[:, list(d.columns).index('机组编号')].notna()]
  134. res.append((sh, d))
  135. return res
  136. def map_sheet(d: pd.DataFrame, src_name: str, n_turbines: int, aliases) -> pd.DataFrame:
  137. """一张模板表 → 输出列。过滤: 风场名称属于本场 ∧ 机组编号可解析 (见文件头「源件形态」)。"""
  138. out = pd.DataFrame(index=d.index)
  139. for src, dst in SRC_MAP.items():
  140. if src in d.columns and dst not in out.columns:
  141. out[dst] = cells_to_str(d[src]) if dst in TEXT_COLS else pd.to_numeric(d[src], errors='coerce')
  142. for c in OUT_COLS:
  143. if c not in out.columns:
  144. out[c] = pd.NA if c in TEXT_COLS else float('nan')
  145. # ★集团公司口径的导出件里混着别的场站 (实测 2021年3月 那张表里 民勤/来福/宝力格/大岗子/宏基 … 十来个场),
  146. # 必须按场名过滤; 场名缺失的模板 (早期单场表没有这列) 则不过滤。
  147. if '风场名称' in d.columns and aliases:
  148. farm = d['风场名称'].astype('string').fillna('').str.strip()
  149. keep = farm.map(lambda s: any(al in s for al in aliases))
  150. out = out[keep]
  151. d = d[keep]
  152. gid = pd.to_numeric(d.get('机组编号'), errors='coerce')
  153. out['turbine'] = gid.map(lambda v: f'WTG{int(v):02d}' if pd.notna(v) else pd.NA).astype('string')
  154. t_rep, bad_rep, empty_rep = parse_ts_series(d.get('故障报出时间'))
  155. out['t_report'] = t_rep
  156. rcol = '复位运行时间' if '复位运行时间' in d.columns else ('复位时间时间' if '复位时间时间' in d.columns else None)
  157. t_res, bad_res, empty_res = parse_ts_series(d[rcol]) if rcol else (pd.Series(pd.NaT, index=d.index), [], 0)
  158. out['t_reset'] = t_res
  159. out['src_file'] = src_name
  160. out = out[out['turbine'].notna()] # 机组编号解析不出 = 不是本场台账记录
  161. out.attrs['bad_dates'] = list(bad_rep) + list(bad_res)
  162. out.attrs['empty_dates'] = int(empty_rep) + int(empty_res)
  163. return out[OUT_COLS]
  164. def main() -> int:
  165. ap = argparse.ArgumentParser()
  166. ap.add_argument('--farm', default=None)
  167. ap.add_argument('--raw', default=None)
  168. ap.add_argument('--dry-run', action='store_true')
  169. ap.add_argument('--keep-dups', action='store_true', help='保留跨源件的内容重复行 (默认归并)')
  170. a = ap.parse_args()
  171. from src.windscada.config import farm
  172. cfg = farm(a.farm)
  173. src = pathlib.Path(a.raw) if a.raw else pathlib.Path(str(cfg['src_workorder']))
  174. store = pathlib.Path(cfg['store'])
  175. out_path = store / 'workorders.parquet'
  176. print(f'源: {src}\n出: {out_path}\n')
  177. files = sorted(p for p in src.rglob('*')
  178. if p.suffix.lower() in ('.xls', '.xlsx') and not p.name.startswith('~$'))
  179. if not files:
  180. raise SystemExit(f'源目录没有 xls/xlsx: {src}')
  181. # 内容哈希去重: 同年目录下的副本 (业主统计故障/) 与正本逐字节相同 → 只收一次
  182. seen_hash, kept_files, copies = {}, [], []
  183. for p in files:
  184. h = md5(p)
  185. if h in seen_hash:
  186. copies.append((p, seen_hash[h]))
  187. continue
  188. seen_hash[h] = p
  189. kept_files.append(p)
  190. frames, per_file, nat_rows, no_table, all_bad, n_empty = [], {}, 0, [], [], 0
  191. aliases = cfg.get('src_farm_names') or []
  192. n_turb = int(cfg.get('n_turbines') or 0)
  193. print(f'场名过滤: {aliases or "(无, 不过滤)"} 机组数 {n_turb}')
  194. for p in kept_files:
  195. tabs = read_workbook(p)
  196. if not tabs:
  197. no_table.append(p.name)
  198. continue
  199. parts = [map_sheet(t, p.name, n_turb, aliases) for _, t in tabs]
  200. all_bad += [x for d0 in parts for x in (d0.attrs.get('bad_dates') or [])]
  201. n_empty += sum(int(d0.attrs.get('empty_dates') or 0) for d0 in parts)
  202. d = pd.concat(parts, ignore_index=True)
  203. if not len(d):
  204. no_table.append(p.name + ' (过滤后 0 行)')
  205. continue
  206. nat = int(d['t_report'].isna().sum())
  207. nat_rows += nat
  208. per_file[p.name] = (len(d), nat)
  209. frames.append(d)
  210. new = pd.concat(frames, ignore_index=True) if frames else pd.DataFrame(columns=OUT_COLS)
  211. # ★跨源件的内容重复必须归并 (2026-09-11 实测): 同一批行会同时出现在多张表里 ——
  212. # 最典型的是 13 行「计划停机」行**被复制进了 58 张月表** (13×58 = 754 行, 占全量的 11%),
  213. # 另有 {年}总表 与其月表重叠、`…最终版 1090…` 与各月表重叠。不归并 = 台账条目数虚高,
  214. # 与 SOP 里"累计快照必须归并取末值"是同一条纪律 (见 scripts/ingest_ops_2025.py 的长停台账教训)。
  215. # 归并只按 26 列内容全等判重: 内容一样的行留一份, 不丢任何字段值 (丢的是"这行还出现在哪个文件里")。
  216. n_before = len(new)
  217. if not a.keep_dups and len(new):
  218. # ★键必须自己拼, 且**不能含 src_file** —— src_file 正是重复行彼此唯一不同的那一列,
  219. # 带上它就永远判不出跨源件重复 (实测只认出 207/1085)。判重按「内容 25 列」。
  220. # 另外 pandas 的 duplicated() 在 string(masked) 列上对 pd.NA 判等不可靠, 故先规范成字符串。
  221. kc = [c for c in OUT_COLS if c != 'src_file']
  222. k = new[kc].astype('string').fillna('<NA>').agg('\x1f'.join, axis=1)
  223. _g = pd.DataFrame({'k': k, 'f': new['src_file']}).groupby('k')['f'].agg(['size', 'nunique'])
  224. cross_dup = int((_g.loc[_g['nunique'] > 1, 'size'] - 1).sum())
  225. same_dup = int((_g.loc[_g['nunique'] == 1, 'size'] - 1).sum())
  226. new = new[~k.duplicated(keep='first')].reset_index(drop=True)
  227. else:
  228. cross_dup = same_dup = 0
  229. n_dup = n_before - len(new)
  230. print(f'== 源件 ==\n 唯一件 {len(kept_files)} 张 → {len(new)} 行; 非台账表跳过 {len(no_table)} 张')
  231. if n_dup:
  232. print(f'\n== 跨源件内容重复归并: 去掉 {n_dup} 行 '
  233. f'(跨源件 {cross_dup} 行 + 同源件内 {same_dup} 行) ==')
  234. print(f' 典型: 13 行「计划停机」行被复制进 58 张月表 → 单这一组就 754 行冗余。')
  235. print(f' 判重按内容 25 列全等 (不含 src_file); 保留排序靠前的第一份来源, 字段值零丢失。'
  236. f' 用 --keep-dups 可关闭归并。')
  237. if no_table:
  238. print(' (跳过: ' + ', '.join(no_table[:6]) + (' …' if len(no_table) > 6 else '') + ')')
  239. if copies:
  240. print(f'\n== 副本跳过 {len(copies)} 张 (与正本 md5 相同, 旧链把两份都收了 → 行数虚增) ==')
  241. for p, orig in copies[:6]:
  242. print(f' {p.relative_to(src)} == {orig.relative_to(src)}')
  243. if len(copies) > 6:
  244. print(f' … 其余 {len(copies)-6} 张')
  245. if nat_rows:
  246. bad = {k: v for k, v in per_file.items() if v[1]}
  247. n_bad = len(all_bad)
  248. print(f'\n== 时间列: 空单元格 {n_empty} 行 · 非空但解析失败 {n_bad} 行 · t_report 合计为 NaT {nat_rows} 行 ==')
  249. print(f' (旧链把解析失败的行**整行丢弃**; 本脚本保留整行并把原文本留在 故障报出时间/复位运行时间 列)')
  250. for k, (n, nat) in sorted(bad.items(), key=lambda kv: -kv[1][1])[:6]:
  251. print(f' {k}: {nat}/{n}')
  252. if len(bad) > 6:
  253. print(f' … 其余 {len(bad)-6} 张表')
  254. if all_bad:
  255. from collections import Counter
  256. print(' 失败原文 top: ' + '; '.join(f'{v!r}×{c}' for v, c in Counter(map(str, all_bad)).most_common(8)))
  257. nod = new['t_report'].isna()
  258. plan = int(new.loc[nod, '故障描述'].astype('string').fillna('').str.contains('计划停机').sum())
  259. if int(nod.sum()):
  260. print(f' 其中 **计划性工作行** (故障描述含「计划停机」): {plan} 行 —— 这些行没有「故障报出时间」是正常的'
  261. f'(模板说明第 9 条: 计划性工作只填处理/维护时间), 但旧链把它们整行丢了, 页面上的'
  262. f'"计划/故障停机不可拆"正与此有关。')
  263. old = pd.read_parquet(out_path) if out_path.exists() else pd.DataFrame(columns=OUT_COLS)
  264. touch = set(per_file)
  265. # ★产物 = **当前源件的函数** (用户令 2026-09-20「重算要按 data/raw 的最新变化」): 原先 `keep_old` 会把
  266. # "源件已不在盘"的行一直留着(页面照旧按已不存在的台账显示检修记录)。现在只保留仍在盘的源件的行,
  267. # 消失的行如实丢掉并报出来(哪几件、多少行、涉及哪些月), 让人分清误删/撤数。
  268. present = {p.name for p in files}
  269. gone = sorted(set(old['src_file'].dropna().unique()) - present) if len(old) else []
  270. drop_gone = old[old['src_file'].isin(gone)] if gone else old.iloc[0:0]
  271. keep_old = old[~old['src_file'].isin(touch | set(gone))] if len(old) else old
  272. merged = pd.concat([keep_old, new], ignore_index=True)
  273. merged = merged.sort_values(['t_report', 'turbine', 'src_file'], na_position='last').reset_index(drop=True)
  274. tw = merged['t_report'].dropna()
  275. print(f'\n== 结果 ==\n 旧 {len(old)} 行 → 新 {len(merged)} 行 (本次 {len(touch)} 张表; 保留未涉及 {len(keep_old)} 行)')
  276. if len(drop_gone):
  277. print(f' [!] 有 {len(drop_gone)} 行来自**已不在盘**的台账件 ({len(gone)} 个): '
  278. + '、'.join(f'{n}({int((drop_gone["src_file"] == n).sum())} 行)' for n in gone[:6])
  279. + ' —— 按"产物=当前 data/raw 的函数"口径不再产出; 误删就把文件放回重跑')
  280. _m = pd.to_datetime(drop_gone['t_report'], errors='coerce')
  281. if _m.notna().any():
  282. print(f' 涉及月份: {_m.dt.strftime("%Y-%m").dropna().value_counts().sort_index().to_dict()}')
  283. print(f' 覆盖: {tw.min()} ~ {tw.max()} (未解析日期 {int(merged["t_report"].isna().sum())} 行)')
  284. print(f' 台 {merged.turbine.nunique()} 台; 源件 {merged.src_file.nunique()} 张')
  285. if a.dry_run:
  286. print('\n(dry-run, 未写盘)')
  287. return 0
  288. store.mkdir(parents=True, exist_ok=True)
  289. merged.to_parquet(out_path, index=False)
  290. print(f'\n已写 {out_path}')
  291. return 0
  292. if __name__ == '__main__':
  293. sys.exit(main())