ingest_ops_2025.py 19 KB

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