ingest_ops_2025.py 19 KB

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