ingest_ops_2025.py 21 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319
  1. #!/usr/bin/env python3
  2. # -*- coding: utf-8 -*-
  3. """接入 2025-2026 检修工单与运行台账 (2026-08-31)。
  4. ## 为什么这批必须接
  5. 本体 WorkOrder 1651 条止于 **2024-11-21**; 工作库里 maintenance_work_order.parquet
  6. 覆盖 **2025-01-03 → 2026-07-08** 共 818 条, 与本体 (机组,日期) 键**零重叠** ——
  7. 不是重复数据, 是缺失的两年。页面上那句"台账止 2024-11, 本期检修执行与闭环无法自动判定"
  8. 正由此而来: 数据一直在, 只是没接线。
  9. ## 它带来的不只是行数
  10. 新表有旧表没有的四类字段, 每类都改变某条判据能否成立:
  11. · work_type 计划性停机 736 / 故障停机 57 → 直接回答双模审"混淆计划性与故障性更换"
  12. · part_brand/spec 供应商与件号 → 成熟度"共性识别"维原判不可评估的理由消失
  13. · check_item 检查项 647 条 → 闭环验证有了记录载体
  14. · downtime_h_calc / repair_h_calc 工时 → 停机与修复时长可分开算
  15. ★但**字段存在 ≠ 覆盖足够**: part_brand 仅 19 条、part_spec 13 条。
  16. 重算成熟度时必须按实际非空率判, 不能因为"字段有了"就宣布该维可评估 —
  17. 那会把 19/818 的覆盖说成能力达标。
  18. """
  19. import pathlib as _pl, sys as _sys
  20. _sys.path.insert(0, str(_pl.Path(__file__).resolve().parents[1])) # 2026-09-16: 与其它摄入脚本一致
  21. from src import paths as _P
  22. import argparse, json, os, pathlib, sys, warnings
  23. warnings.filterwarnings('ignore')
  24. import pandas as pd
  25. def _ledger_range(raw):
  26. """台账维修时间单元格 → (start_iso, end_iso); 解析不出返回 (None, None) 并**响亮打印**。
  27. 不 raise 是因为 ingest 要能跑完整批; 但也绝不静默 —— 解析失败逐条打出原文,
  28. 否则又变回 `维修时间_dt` 那种 "20/20 全 NaT 却无人察觉" 的静默失效。
  29. """
  30. if not raw:
  31. return None, None
  32. try:
  33. from src.sop.ledger_dates import parse_range
  34. a, b = parse_range(raw)
  35. return (a.strftime('%Y-%m-%d') if a is not None else None,
  36. b.strftime('%Y-%m-%d') if b is not None else None)
  37. except Exception as e:
  38. print('[大部件] 日期解析失败, 原文 %r — %s' % (raw, e), file=sys.stderr)
  39. return None, None
  40. OBJ = _P.objects_json()
  41. # ★2026-09-16 修: 原为 `pathlib.Path('outputs/rudong/ontology/objects.json')` —— **cwd 相对** + 写死场名。
  42. # 本文件在模块级用它, 于是: ① 从别处调用/换工作目录就会写到别处; ② 它还是全仓唯一没有
  43. # `sys.path.insert(0, str(ROOT))` 的摄入脚本, 干净环境下连 `src.paths` 都导不进来。
  44. # 注意: 该脚本此前因为缺 `import os` (上面那行用 os.environ) 在模块级就 NameError, 从未真正跑过 ——
  45. # 所以"补 import os"必须和"改掉 cwd 相对路径 + 走 Store 落盘"同一次做完, 否则只是把一把静默的凶器上膛。
  46. OPS = pathlib.Path(os.environ.get('WINDSCADA_OPS_LIB') or (_P.RAW_ROOT / '工作库' / '20_structured' / 'ops'))
  47. # 原为 mac 盘路径 (本机不存在): 改为 env 可配 + 安装目录相对默认位置
  48. def _s(v):
  49. if v is None or (isinstance(v, float) and pd.isna(v)):
  50. return None
  51. s = str(v).strip()
  52. return s if s and s.lower() not in ('nan', 'nat', 'none') else None
  53. def main():
  54. ap = argparse.ArgumentParser()
  55. ap.add_argument('--dry-run', action='store_true')
  56. a = ap.parse_args()
  57. db = json.loads(OBJ.read_text(encoding='utf-8'))
  58. n0 = len(db)
  59. added = {}
  60. # ── ① 2025-2026 检修工单 ──
  61. d = pd.read_parquet(OPS / 'maintenance_work_order.parquet')
  62. n_wo = 0
  63. for i, r in d.iterrows():
  64. tid = r.get('tid')
  65. if pd.isna(tid):
  66. continue
  67. t = f'WTG{int(tid):02d}'
  68. date = pd.to_datetime(r.get('time_start'), errors='coerce')
  69. if pd.isna(date):
  70. continue
  71. ds = date.strftime('%Y-%m-%d')
  72. sq = r.get('seq')
  73. sq = int(sq) if pd.notna(sq) else int(i)
  74. oid = f'workorder/{t}/{ds}/{sq}'
  75. # ★保留下游脚本挂上的派生字段 (2026-08-31 实逮回归)。
  76. # 本脚本原来 `db[oid] = dict(...)` 是**整体替换**: 重跑一次就把
  77. # ingest_workorder_en.py 挂的 ge_system 全抹掉 —— 实测 2022/2025/2026 三年
  78. # 共 766 条工单的 ge_system 归零, 而依赖它的 tier_trend finding **还留在本体里**,
  79. # 于是出现"结论在、支撑它的数据没了"。
  80. # ⇒ ingest 类脚本必须**合并**而非替换: 自己产的字段照写, 不认识的字段原样保留。
  81. _keep = {k: v for k, v in ((db.get(oid) or {}).get('props') or {}).items()
  82. if k in ('ge_system', 'ge_tier', 'name_en', 'name_en_source')}
  83. props = dict(turbine=t, date=ds,
  84. name=_s(r.get('fault_name')) or _s(r.get('fault_desc')) or '(未命名)',
  85. loc2=_s(r.get('component_std')), part=_s(r.get('part_name')),
  86. work_type=_s(r.get('work_type')),
  87. fault_desc=_s(r.get('fault_desc')),
  88. check_item=_s(r.get('check_item')),
  89. part_spec=_s(r.get('part_spec')), part_qty=_s(r.get('part_qty')),
  90. part_brand=_s(r.get('part_brand')),
  91. downtime_h=None if pd.isna(r.get('downtime_h_calc')) else float(r['downtime_h_calc']),
  92. repair_h=None if pd.isna(r.get('repair_h_calc')) else float(r['repair_h_calc']),
  93. time_start=ds, time_end=_s(r.get('time_end')),
  94. source='工作库/20_structured/ops/maintenance_work_order.parquet')
  95. db[oid] = dict(id=oid, type='WorkOrder',
  96. props=dict(_keep, **{k: v for k, v in props.items() if v is not None}),
  97. links=dict(about=[f'turbine/{t}']))
  98. n_wo += 1
  99. added['WorkOrder(2025-2026)'] = n_wo
  100. # ── ② 大部件更换 ──
  101. d = pd.read_parquet(OPS / 'major_component_repair.parquet')
  102. n_mc = 0
  103. for i, r in d.iterrows():
  104. tid = r.get('tid')
  105. t = f'WTG{int(tid):02d}' if pd.notna(tid) else None
  106. oid = f'majorrepair/{t or "unknown"}/{i}'
  107. # ★日期必须走 src/sop/ledger_dates.parse_range, 不能裸 _s() 取。
  108. # 台账 20 行**全部有日期**, 但人工录入格式极不统一 ('2024年5月13日至2024年5月19日' /
  109. # '2024/9/19至2024-10-13' / '2025/8/23 -2025-8-28' / '2025年12月4日-2025年12月70')。
  110. # 裸取的结果: 7/20 落成 null, 而"大部件更换"是留存期 life 的记录 —— 没有日期
  111. # 意味着十年后回答不了"这台什么时候换的主轴"。
  112. # 解析器 2026-08-26 就写好了并带三条硬校验, 只是**没接上生产路径** (落库≠接上)。
  113. raw_dt = _s(r.get('大部件维修时间')) or _s(r.get('维修时间')) or _s(r.get('time_on'))
  114. rng = _ledger_range(raw_dt)
  115. db[oid] = dict(id=oid, type='WorkOrder', props=dict(
  116. turbine=t, date=rng[0], date_end=rng[1], date_raw=raw_dt,
  117. name=_s(r.get('大部件更换项目')) or '大部件更换',
  118. loc2=_s(r.get('所属系统')), part=_s(r.get('component_std')),
  119. action='更换', major_component=True,
  120. model=_s(r.get('型号')), reason=_s(r.get('维修/更换原因')),
  121. damage_cause=_s(r.get('损坏原因')),
  122. failure_report=_s(r.get('失效分析报告')),
  123. source='工作库/20_structured/ops/major_component_repair.parquet'),
  124. links=dict(about=[f'turbine/{t}']) if t else {})
  125. n_mc += 1
  126. added['WorkOrder(大部件更换)'] = n_mc
  127. # ── ③ 限电区间 ──
  128. d = pd.read_parquet(OPS / 'curtailment_interval.parquet')
  129. n_cu = 0
  130. for i, r in d.iterrows():
  131. tid = r.get('tid')
  132. t = f'WTG{int(tid):02d}' if pd.notna(tid) else 'fleet'
  133. oid = f'curtailment/{t}/{i}'
  134. db[oid] = dict(id=oid, type='Evidence', props=dict(
  135. name=f'{t} 限电区间', facet='外部记录',
  136. claim_window=f"{_s(r.get('time_on'))} ~ {_s(r.get('time_off'))}",
  137. source='工作库/20_structured/ops/curtailment_interval.parquet',
  138. turbine=t, time_on=_s(r.get('time_on')), time_off=_s(r.get('time_off')),
  139. # ★dur_h 在 time_off 缺失时会算出负值 (实测 33 条, 最小 -1108697h = -126 年)。
  140. # 负时长物理不可能 ⇒ 置 None 并标因, 不让它流进任何求和。
  141. dur_h=(lambda v: None if (pd.isna(v) or float(v) < 0) else float(v))(r.get('dur_h')),
  142. dur_h_invalid=(None if pd.isna(r.get('dur_h')) else
  143. ('negative: time_off missing' if float(r['dur_h']) < 0 else None)),
  144. p_limit_mw=None if pd.isna(r.get('p_limit_mw')) else float(r['p_limit_mw']),
  145. reason=_s(r.get('reason_class')) or _s(r.get('reason_raw')),
  146. loss_kwh=None if pd.isna(r.get('loss_kwh')) else float(r['loss_kwh'])),
  147. links=dict(about=[f'turbine/{t}']) if t != 'fleet' else {})
  148. n_cu += 1
  149. added['Evidence(限电区间)'] = n_cu
  150. # ── ④ 长停登记 (只取本场; 原表含射阳 43 条) ──
  151. #
  152. # ★★ 累计快照必须先归并 (2026-08-31 实测)
  153. # 原表如海 2277 行, 但按 (设备编号, 停机时间, 恢复时间) 判重只有 **59 个真实事件** ——
  154. # 同一条长停出现在 43 个不同的工作记事文件里 (2025年1月/3月/7月/9月…),
  155. # 因为每份工作记事都附一份"当前未结束长停"的全量台账。
  156. # 直接建对象会把 59 个事件记成 2277 条, **停机天数虚高 38 倍**:
  157. # 不归并 15383 天 = 全场时长的 17% (与可用率 94.8% 直接矛盾)
  158. # 归并后 437 天 = 0.5% (与可用率相容)
  159. # 这与 SOP §4.11 EL-6 记的"工单损失列=累计逐日快照, 直接 sum 虚高 3.4×"是同一个陷阱。
  160. # **判别法**: 同一业务键重复出现且各来自不同 src_file ⇒ 快照类, 必须归并取末值。
  161. d = pd.read_parquet(OPS / 'longstop_register.parquet')
  162. if '场站' in d.columns:
  163. d = d[d['场站'].astype(str).str.contains('如海', na=False)]
  164. n_raw = len(d)
  165. keycols = [c for c in ('设备编号', '停机时间', '恢复时间/预计投运时间') if c in d.columns]
  166. if keycols:
  167. if 'src_file' in d.columns:
  168. d = d.sort_values('src_file') # 取末值 = 最新一份快照里的状态
  169. d = d.drop_duplicates(subset=keycols, keep='last')
  170. print(f' 长停登记归并: {n_raw} 行 → {len(d)} 个事件 (累计快照去重)')
  171. # ★★ 产出 finding —— **数字必须在这里算**, 不在别处手写 (2026-08-31 独立审逮)。
  172. # 原来这条结论是会话内联建的对象, 仓库里没有生成端 ⇒ 按 SOP §0.2,
  173. # 【定论】的出具条件含"§9.1.5 可复现脚本", 没有脚本就够不上定论。
  174. # 审者据此判硬伤, 是对的。修法不是降级, 是把生成端补在**算它的地方**。
  175. _raw_all = pd.read_parquet(OPS / 'longstop_register.parquet')
  176. if '场站' in _raw_all.columns:
  177. _raw_all = _raw_all[_raw_all['场站'].astype(str).str.contains('如海', na=False)]
  178. _dc = '停机时长(天)'
  179. _days_raw = float(_raw_all[_dc].dropna().sum()) if _dc in _raw_all.columns else None
  180. _days_ded = float(d[_dc].dropna().sum()) if _dc in d.columns else None
  181. _nfile = int(_raw_all['src_file'].nunique()) if 'src_file' in _raw_all.columns else None
  182. # 分母 = 38 台 × 台账窗天数 (从数据本身取窗, 不写死)
  183. _ts = pd.to_datetime(_raw_all['停机时间'], errors='coerce')
  184. _span = (pd.Timestamp.today() - _ts.min()).days if _ts.notna().any() else None
  185. _fleet_days = 38 * _span if _span else None
  186. _pct_raw = round(100 * _days_raw / _fleet_days, 1) if (_days_raw and _fleet_days) else None
  187. _pct_ded = round(100 * _days_ded / _fleet_days, 2) if (_days_ded and _fleet_days) else None
  188. # ★两个不同的比, 必须分开命名 (2026-08-31 独立审后 data-first 复核逮到的自己的错):
  189. # _infl = **天数**虚高 = 15383/437 = 35.2× ← 这才是"停机天数虚高"
  190. # _rowr = **行数**比 = 2277/59 = 38.6× ← 首版把这个数写进了"天数虚高 38 倍"的标题
  191. # 两个数都对, 错在**标题的名词与数字不是同一个量** (memory headline-number-mislabels-own-decomposition)
  192. _infl = round(_days_raw / _days_ded, 1) if (_days_raw and _days_ded) else None
  193. _rowr = round(n_raw / len(d), 1) if len(d) else None
  194. # 场龄口径的占比: 首版分母用的是场龄 6.53 年, 而台账窗从 2022-12 才开始 ⇒ 分子分母窗不对齐,
  195. # 占比被系统性低估 (17% vs 台账窗口径 29.9%)。两个都报, 让读者看到口径差在哪。
  196. _age_days = 38 * int(6.53 * 365.25)
  197. _pct_raw_age = round(100 * _days_raw / _age_days, 1) if _days_raw else None
  198. # 复现次数最多的那一条 (支撑"同一事件出现在 N 份文件里")
  199. _top = None
  200. if keycols and 'src_file' in _raw_all.columns:
  201. _g = _raw_all.groupby(keycols, dropna=False)['src_file'].nunique().sort_values()
  202. if len(_g):
  203. _k, _n = _g.index[-1], int(_g.iloc[-1])
  204. _top = dict(key=[str(x)[:19] for x in (_k if isinstance(_k, tuple) else (_k,))], n_files=_n)
  205. _fid = 'finding/data_quality/cumulative_snapshot_inflation'
  206. _old = (db.get(_fid) or {}).get('props', {})
  207. db[_fid] = dict(id=_fid, type='Finding', props=dict(
  208. _old,
  209. name=f'长停台账是累计快照, 停机天数虚高 {_infl} 倍',
  210. name_zh=f'长停台账为累计快照: {len(d)} 条归并后记录散在 {_nfile} 份工作记事里共 {n_raw} 行 '
  211. f'(行数比 {_rowr}×), 不归并则**停机天数**虚高 {_infl} 倍 ({_pct_ded}% → {_pct_raw}%)',
  212. name_en=f'The long-stop register is a cumulative snapshot: {len(d)} de-duplicated '
  213. f'records appear as {n_raw} rows across {_nfile} work-log files (a row ratio of '
  214. f'{_rowr}x); left un-collapsed, **downtime days** are inflated {_infl}-fold '
  215. f'({_pct_ded}% to {_pct_raw}%)',
  216. criterion='R4', grade='established', count=len(d), count_basis='events',
  217. count_basis_note=f'{len(d)} = 归并后的唯一长停事件数',
  218. claim_window=f"{_s(_ts.min())[:10]} ~ {_s(_ts.max())[:10]} 长停台账",
  219. evidence_zh=[
  220. f'原表本场 {n_raw} 行, 按 (设备编号, 停机时间, 恢复时间) 判重仅 {len(d)} 个唯一事件',
  221. (f'复现次数最多的一条出现 {_top["n_files"]} 次, 分别来自 {_top["n_files"]} 个不同的'
  222. f'工作记事文件 —— 每份工作记事都附一份当前未结束长停的全量台账') if _top else '',
  223. f'不归并: {_days_raw:.0f} 天 = 台账窗内全场时长的 {_pct_raw}%, 与可用率 94.8% 直接矛盾',
  224. f'归并后: {_days_ded:.0f} 天 = {_pct_ded}%, 与可用率相容',
  225. f'★口径: 分母 = 38 台 × 台账窗 {_span} 天 (首条停机至今)。若改用场龄 6.53 年作分母则为 '
  226. f'{_pct_raw_age}% —— 但台账窗始于 2022-12, 用场龄会让分子分母窗不对齐、占比被低估',
  227. f'★两个比不同: 行数比 {_rowr}× (2277/59) vs 停机天数虚高 {_infl}× '
  228. f'({_days_raw:.0f}/{_days_ded:.0f})。首版标题误用行数比来描述天数'],
  229. evidence_en=[
  230. f'{n_raw} rows for this site collapse to {len(d)} unique events when de-duplicated on '
  231. f'(unit, stop time, resume time)',
  232. (f'The most-repeated entry appears {_top["n_files"]} times, once in each of '
  233. f'{_top["n_files"]} different work-log files, because every work log carries a full '
  234. f'copy of the currently-open long-stop register') if _top else '',
  235. f'Without de-duplication: {_days_raw:.0f} days = {_pct_raw}% of fleet time, which '
  236. f'directly contradicts the 94.8% availability figure',
  237. f'After de-duplication: {_days_ded:.0f} days = {_pct_ded}%, consistent with availability',
  238. f'Basis: the denominator is 38 units x {_span} register days (first stop to date). '
  239. f'Using the {6.53}-year age of the site instead gives {_pct_raw_age}%, but the register '
  240. f'only starts in December 2022, so that denominator does not match the numerator window',
  241. f'Two different ratios: rows {_rowr}x (2277/59) versus downtime days {_infl}x '
  242. f'({_days_raw:.0f}/{_days_ded:.0f}). The first version used the row ratio in a headline '
  243. f'that named downtime days'],
  244. discriminator_zh='**同一业务键重复出现且各来自不同 src_file ⇒ 快照类, 必须归并取末值**',
  245. discriminator_en='If the same business key repeats with each occurrence coming from a '
  246. 'different source file, the table is a cumulative snapshot and must be '
  247. 'collapsed to the latest state per key',
  248. related='与 SOP §4.11 EL-6 记录的"工单损失列=累计逐日快照, 直接 sum 虚高 3.4×"同一陷阱 '
  249. '(两个倍数各属不同表, 不是同一个数的两种说法)',
  250. caveat=(f'{len(d)} 是**归并后的台账行数**, 不等于"本场五年只发生过 {len(d)} 次长停" —— '
  251. f'若某条长停的起止时间在不同快照里被改写过, 会被判成多条; 反之跨快照周期的'
  252. f'同一次长停若键值一致则合为一条。真实长停次数可能与此不同。'
  253. f'另: 归并取末值, 对在最后一份快照前已结束的事件, 其结束时间取自最后一次出现的记录, '
  254. f'可能早于真实复役时间。'),
  255. caveat_en=(f'{len(d)} is the number of rows after de-duplication, not a claim that only '
  256. f'{len(d)} long stops occurred at this site in five years: an event whose start '
  257. f'or end time was edited between snapshots would be counted more than once, '
  258. f'while one that kept identical key values across snapshot periods collapses to '
  259. f'a single row. The true number of long stops may differ. Collapsing also keeps '
  260. f'the state from the latest snapshot in which each event appears, so for an '
  261. f'event that ended earlier the recorded resume time may precede the true one.'),
  262. provenance='scripts/ingest_ops_2025.py (数字实时算, 可复现)',
  263. ), links=dict(about=['farm/rudong']))
  264. print(f' → finding cumulative_snapshot_inflation: {n_raw}行/{_nfile}文件 → {len(d)}条, '
  265. f'行数比 {_rowr}× / 天数虚高 {_infl}× ({_pct_ded}% → {_pct_raw}%; 场龄口径 {_pct_raw_age}%)')
  266. n_ls = 0
  267. for i, r in d.iterrows():
  268. tid = r.get('tid')
  269. t = f'WTG{int(tid):02d}' if pd.notna(tid) else None
  270. oid = f'longstop/{t or "unknown"}/{i}'
  271. db[oid] = dict(id=oid, type='Evidence', props=dict(
  272. name=f'{t or "?"} 长停登记', facet='外部记录',
  273. claim_window=f"{_s(r.get('停机时间'))} ~ {_s(r.get('恢复时间/预计投运时间'))}",
  274. source='工作库/20_structured/ops/longstop_register.parquet (仅如海)',
  275. turbine=t, stop_time=_s(r.get('停机时间')),
  276. resume_time=_s(r.get('恢复时间/预计投运时间')),
  277. stop_class=_s(r.get('停机分类')),
  278. stop_days=None if pd.isna(r.get('停机时长(天)')) else float(r['停机时长(天)'])),
  279. links=dict(about=[f'turbine/{t}']) if t else {})
  280. n_ls += 1
  281. added['Evidence(长停登记)'] = n_ls
  282. for k, v in added.items():
  283. print(f' {k:26s} {v:5d}')
  284. print(f' 本体 {n0} → {len(db)}')
  285. if a.dry_run:
  286. print(' (dry-run, 未写盘)')
  287. return 0
  288. # 走本体 store 落盘 (2026-09-16): 原来直接 `OBJ.write_text(...)` —— 绕过了 store 的
  289. # ① 全库 validate 闸 ② F12 原子写 (tmp + os.replace)。对象库里混进一个坏对象时, 直写会
  290. # 让**整库**在下次读取时才炸, 而 store.save() 会当场报出是哪个对象不合法。
  291. from src.ontology.store import Store
  292. st = Store(OBJ)
  293. st.objects = db
  294. st.save()
  295. print(f' 已写入 {_P.rel(OBJ)} (经 Store 闸)')
  296. return 0
  297. if __name__ == '__main__':
  298. sys.exit(main())