| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319 |
- #!/usr/bin/env python3
- # -*- coding: utf-8 -*-
- """接入 2025-2026 检修工单与运行台账 (2026-08-31)。
- ## 为什么这批必须接
- 本体 WorkOrder 1651 条止于 **2024-11-21**; 工作库里 maintenance_work_order.parquet
- 覆盖 **2025-01-03 → 2026-07-08** 共 818 条, 与本体 (机组,日期) 键**零重叠** ——
- 不是重复数据, 是缺失的两年。页面上那句"台账止 2024-11, 本期检修执行与闭环无法自动判定"
- 正由此而来: 数据一直在, 只是没接线。
- ## 它带来的不只是行数
- 新表有旧表没有的四类字段, 每类都改变某条判据能否成立:
- · work_type 计划性停机 736 / 故障停机 57 → 直接回答双模审"混淆计划性与故障性更换"
- · part_brand/spec 供应商与件号 → 成熟度"共性识别"维原判不可评估的理由消失
- · check_item 检查项 647 条 → 闭环验证有了记录载体
- · downtime_h_calc / repair_h_calc 工时 → 停机与修复时长可分开算
- ★但**字段存在 ≠ 覆盖足够**: part_brand 仅 19 条、part_spec 13 条。
- 重算成熟度时必须按实际非空率判, 不能因为"字段有了"就宣布该维可评估 —
- 那会把 19/818 的覆盖说成能力达标。
- """
- import pathlib as _pl, sys as _sys
- _sys.path.insert(0, str(_pl.Path(__file__).resolve().parents[1])) # 2026-09-16: 与其它摄入脚本一致
- from src import paths as _P
- import argparse, json, os, pathlib, sys, warnings
- warnings.filterwarnings('ignore')
- import pandas as pd
- def _ledger_range(raw):
- """台账维修时间单元格 → (start_iso, end_iso); 解析不出返回 (None, None) 并**响亮打印**。
- 不 raise 是因为 ingest 要能跑完整批; 但也绝不静默 —— 解析失败逐条打出原文,
- 否则又变回 `维修时间_dt` 那种 "20/20 全 NaT 却无人察觉" 的静默失效。
- """
- if not raw:
- return None, None
- try:
- from src.sop.ledger_dates import parse_range
- a, b = parse_range(raw)
- return (a.strftime('%Y-%m-%d') if a is not None else None,
- b.strftime('%Y-%m-%d') if b is not None else None)
- except Exception as e:
- print('[大部件] 日期解析失败, 原文 %r — %s' % (raw, e), file=sys.stderr)
- return None, None
- OBJ = _P.objects_json()
- # ★2026-09-16 修: 原为 `pathlib.Path('outputs/rudong/ontology/objects.json')` —— **cwd 相对** + 写死场名。
- # 本文件在模块级用它, 于是: ① 从别处调用/换工作目录就会写到别处; ② 它还是全仓唯一没有
- # `sys.path.insert(0, str(ROOT))` 的摄入脚本, 干净环境下连 `src.paths` 都导不进来。
- # 注意: 该脚本此前因为缺 `import os` (上面那行用 os.environ) 在模块级就 NameError, 从未真正跑过 ——
- # 所以"补 import os"必须和"改掉 cwd 相对路径 + 走 Store 落盘"同一次做完, 否则只是把一把静默的凶器上膛。
- OPS = pathlib.Path(os.environ.get('WINDSCADA_OPS_LIB') or (_P.RAW_ROOT / '工作库' / '20_structured' / 'ops'))
- # 原为 mac 盘路径 (本机不存在): 改为 env 可配 + 安装目录相对默认位置
- def _s(v):
- if v is None or (isinstance(v, float) and pd.isna(v)):
- return None
- s = str(v).strip()
- return s if s and s.lower() not in ('nan', 'nat', 'none') else None
- def main():
- ap = argparse.ArgumentParser()
- ap.add_argument('--dry-run', action='store_true')
- a = ap.parse_args()
- db = json.loads(OBJ.read_text(encoding='utf-8'))
- n0 = len(db)
- added = {}
- # ── ① 2025-2026 检修工单 ──
- d = pd.read_parquet(OPS / 'maintenance_work_order.parquet')
- n_wo = 0
- for i, r in d.iterrows():
- tid = r.get('tid')
- if pd.isna(tid):
- continue
- t = f'WTG{int(tid):02d}'
- date = pd.to_datetime(r.get('time_start'), errors='coerce')
- if pd.isna(date):
- continue
- ds = date.strftime('%Y-%m-%d')
- sq = r.get('seq')
- sq = int(sq) if pd.notna(sq) else int(i)
- oid = f'workorder/{t}/{ds}/{sq}'
- # ★保留下游脚本挂上的派生字段 (2026-08-31 实逮回归)。
- # 本脚本原来 `db[oid] = dict(...)` 是**整体替换**: 重跑一次就把
- # ingest_workorder_en.py 挂的 ge_system 全抹掉 —— 实测 2022/2025/2026 三年
- # 共 766 条工单的 ge_system 归零, 而依赖它的 tier_trend finding **还留在本体里**,
- # 于是出现"结论在、支撑它的数据没了"。
- # ⇒ ingest 类脚本必须**合并**而非替换: 自己产的字段照写, 不认识的字段原样保留。
- _keep = {k: v for k, v in ((db.get(oid) or {}).get('props') or {}).items()
- if k in ('ge_system', 'ge_tier', 'name_en', 'name_en_source')}
- props = dict(turbine=t, date=ds,
- name=_s(r.get('fault_name')) or _s(r.get('fault_desc')) or '(未命名)',
- loc2=_s(r.get('component_std')), part=_s(r.get('part_name')),
- work_type=_s(r.get('work_type')),
- fault_desc=_s(r.get('fault_desc')),
- check_item=_s(r.get('check_item')),
- part_spec=_s(r.get('part_spec')), part_qty=_s(r.get('part_qty')),
- part_brand=_s(r.get('part_brand')),
- downtime_h=None if pd.isna(r.get('downtime_h_calc')) else float(r['downtime_h_calc']),
- repair_h=None if pd.isna(r.get('repair_h_calc')) else float(r['repair_h_calc']),
- time_start=ds, time_end=_s(r.get('time_end')),
- source='工作库/20_structured/ops/maintenance_work_order.parquet')
- db[oid] = dict(id=oid, type='WorkOrder',
- props=dict(_keep, **{k: v for k, v in props.items() if v is not None}),
- links=dict(about=[f'turbine/{t}']))
- n_wo += 1
- added['WorkOrder(2025-2026)'] = n_wo
- # ── ② 大部件更换 ──
- d = pd.read_parquet(OPS / 'major_component_repair.parquet')
- n_mc = 0
- for i, r in d.iterrows():
- tid = r.get('tid')
- t = f'WTG{int(tid):02d}' if pd.notna(tid) else None
- oid = f'majorrepair/{t or "unknown"}/{i}'
- # ★日期必须走 src/sop/ledger_dates.parse_range, 不能裸 _s() 取。
- # 台账 20 行**全部有日期**, 但人工录入格式极不统一 ('2024年5月13日至2024年5月19日' /
- # '2024/9/19至2024-10-13' / '2025/8/23 -2025-8-28' / '2025年12月4日-2025年12月70')。
- # 裸取的结果: 7/20 落成 null, 而"大部件更换"是留存期 life 的记录 —— 没有日期
- # 意味着十年后回答不了"这台什么时候换的主轴"。
- # 解析器 2026-08-26 就写好了并带三条硬校验, 只是**没接上生产路径** (落库≠接上)。
- raw_dt = _s(r.get('大部件维修时间')) or _s(r.get('维修时间')) or _s(r.get('time_on'))
- rng = _ledger_range(raw_dt)
- db[oid] = dict(id=oid, type='WorkOrder', props=dict(
- turbine=t, date=rng[0], date_end=rng[1], date_raw=raw_dt,
- name=_s(r.get('大部件更换项目')) or '大部件更换',
- loc2=_s(r.get('所属系统')), part=_s(r.get('component_std')),
- action='更换', major_component=True,
- model=_s(r.get('型号')), reason=_s(r.get('维修/更换原因')),
- damage_cause=_s(r.get('损坏原因')),
- failure_report=_s(r.get('失效分析报告')),
- source='工作库/20_structured/ops/major_component_repair.parquet'),
- links=dict(about=[f'turbine/{t}']) if t else {})
- n_mc += 1
- added['WorkOrder(大部件更换)'] = n_mc
- # ── ③ 限电区间 ──
- d = pd.read_parquet(OPS / 'curtailment_interval.parquet')
- n_cu = 0
- for i, r in d.iterrows():
- tid = r.get('tid')
- t = f'WTG{int(tid):02d}' if pd.notna(tid) else 'fleet'
- oid = f'curtailment/{t}/{i}'
- db[oid] = dict(id=oid, type='Evidence', props=dict(
- name=f'{t} 限电区间', facet='外部记录',
- claim_window=f"{_s(r.get('time_on'))} ~ {_s(r.get('time_off'))}",
- source='工作库/20_structured/ops/curtailment_interval.parquet',
- turbine=t, time_on=_s(r.get('time_on')), time_off=_s(r.get('time_off')),
- # ★dur_h 在 time_off 缺失时会算出负值 (实测 33 条, 最小 -1108697h = -126 年)。
- # 负时长物理不可能 ⇒ 置 None 并标因, 不让它流进任何求和。
- dur_h=(lambda v: None if (pd.isna(v) or float(v) < 0) else float(v))(r.get('dur_h')),
- dur_h_invalid=(None if pd.isna(r.get('dur_h')) else
- ('negative: time_off missing' if float(r['dur_h']) < 0 else None)),
- p_limit_mw=None if pd.isna(r.get('p_limit_mw')) else float(r['p_limit_mw']),
- reason=_s(r.get('reason_class')) or _s(r.get('reason_raw')),
- loss_kwh=None if pd.isna(r.get('loss_kwh')) else float(r['loss_kwh'])),
- links=dict(about=[f'turbine/{t}']) if t != 'fleet' else {})
- n_cu += 1
- added['Evidence(限电区间)'] = n_cu
- # ── ④ 长停登记 (只取本场; 原表含射阳 43 条) ──
- #
- # ★★ 累计快照必须先归并 (2026-08-31 实测)
- # 原表如海 2277 行, 但按 (设备编号, 停机时间, 恢复时间) 判重只有 **59 个真实事件** ——
- # 同一条长停出现在 43 个不同的工作记事文件里 (2025年1月/3月/7月/9月…),
- # 因为每份工作记事都附一份"当前未结束长停"的全量台账。
- # 直接建对象会把 59 个事件记成 2277 条, **停机天数虚高 38 倍**:
- # 不归并 15383 天 = 全场时长的 17% (与可用率 94.8% 直接矛盾)
- # 归并后 437 天 = 0.5% (与可用率相容)
- # 这与 SOP §4.11 EL-6 记的"工单损失列=累计逐日快照, 直接 sum 虚高 3.4×"是同一个陷阱。
- # **判别法**: 同一业务键重复出现且各来自不同 src_file ⇒ 快照类, 必须归并取末值。
- d = pd.read_parquet(OPS / 'longstop_register.parquet')
- if '场站' in d.columns:
- d = d[d['场站'].astype(str).str.contains('如海', na=False)]
- n_raw = len(d)
- keycols = [c for c in ('设备编号', '停机时间', '恢复时间/预计投运时间') if c in d.columns]
- if keycols:
- if 'src_file' in d.columns:
- d = d.sort_values('src_file') # 取末值 = 最新一份快照里的状态
- d = d.drop_duplicates(subset=keycols, keep='last')
- print(f' 长停登记归并: {n_raw} 行 → {len(d)} 个事件 (累计快照去重)')
- # ★★ 产出 finding —— **数字必须在这里算**, 不在别处手写 (2026-08-31 独立审逮)。
- # 原来这条结论是会话内联建的对象, 仓库里没有生成端 ⇒ 按 SOP §0.2,
- # 【定论】的出具条件含"§9.1.5 可复现脚本", 没有脚本就够不上定论。
- # 审者据此判硬伤, 是对的。修法不是降级, 是把生成端补在**算它的地方**。
- _raw_all = pd.read_parquet(OPS / 'longstop_register.parquet')
- if '场站' in _raw_all.columns:
- _raw_all = _raw_all[_raw_all['场站'].astype(str).str.contains('如海', na=False)]
- _dc = '停机时长(天)'
- _days_raw = float(_raw_all[_dc].dropna().sum()) if _dc in _raw_all.columns else None
- _days_ded = float(d[_dc].dropna().sum()) if _dc in d.columns else None
- _nfile = int(_raw_all['src_file'].nunique()) if 'src_file' in _raw_all.columns else None
- # 分母 = 38 台 × 台账窗天数 (从数据本身取窗, 不写死)
- _ts = pd.to_datetime(_raw_all['停机时间'], errors='coerce')
- _span = (pd.Timestamp.today() - _ts.min()).days if _ts.notna().any() else None
- _fleet_days = 38 * _span if _span else None
- _pct_raw = round(100 * _days_raw / _fleet_days, 1) if (_days_raw and _fleet_days) else None
- _pct_ded = round(100 * _days_ded / _fleet_days, 2) if (_days_ded and _fleet_days) else None
- # ★两个不同的比, 必须分开命名 (2026-08-31 独立审后 data-first 复核逮到的自己的错):
- # _infl = **天数**虚高 = 15383/437 = 35.2× ← 这才是"停机天数虚高"
- # _rowr = **行数**比 = 2277/59 = 38.6× ← 首版把这个数写进了"天数虚高 38 倍"的标题
- # 两个数都对, 错在**标题的名词与数字不是同一个量** (memory headline-number-mislabels-own-decomposition)
- _infl = round(_days_raw / _days_ded, 1) if (_days_raw and _days_ded) else None
- _rowr = round(n_raw / len(d), 1) if len(d) else None
- # 场龄口径的占比: 首版分母用的是场龄 6.53 年, 而台账窗从 2022-12 才开始 ⇒ 分子分母窗不对齐,
- # 占比被系统性低估 (17% vs 台账窗口径 29.9%)。两个都报, 让读者看到口径差在哪。
- _age_days = 38 * int(6.53 * 365.25)
- _pct_raw_age = round(100 * _days_raw / _age_days, 1) if _days_raw else None
- # 复现次数最多的那一条 (支撑"同一事件出现在 N 份文件里")
- _top = None
- if keycols and 'src_file' in _raw_all.columns:
- _g = _raw_all.groupby(keycols, dropna=False)['src_file'].nunique().sort_values()
- if len(_g):
- _k, _n = _g.index[-1], int(_g.iloc[-1])
- _top = dict(key=[str(x)[:19] for x in (_k if isinstance(_k, tuple) else (_k,))], n_files=_n)
- _fid = 'finding/data_quality/cumulative_snapshot_inflation'
- _old = (db.get(_fid) or {}).get('props', {})
- db[_fid] = dict(id=_fid, type='Finding', props=dict(
- _old,
- name=f'长停台账是累计快照, 停机天数虚高 {_infl} 倍',
- name_zh=f'长停台账为累计快照: {len(d)} 条归并后记录散在 {_nfile} 份工作记事里共 {n_raw} 行 '
- f'(行数比 {_rowr}×), 不归并则**停机天数**虚高 {_infl} 倍 ({_pct_ded}% → {_pct_raw}%)',
- name_en=f'The long-stop register is a cumulative snapshot: {len(d)} de-duplicated '
- f'records appear as {n_raw} rows across {_nfile} work-log files (a row ratio of '
- f'{_rowr}x); left un-collapsed, **downtime days** are inflated {_infl}-fold '
- f'({_pct_ded}% to {_pct_raw}%)',
- criterion='R4', grade='established', count=len(d), count_basis='events',
- count_basis_note=f'{len(d)} = 归并后的唯一长停事件数',
- claim_window=f"{_s(_ts.min())[:10]} ~ {_s(_ts.max())[:10]} 长停台账",
- evidence_zh=[
- f'原表本场 {n_raw} 行, 按 (设备编号, 停机时间, 恢复时间) 判重仅 {len(d)} 个唯一事件',
- (f'复现次数最多的一条出现 {_top["n_files"]} 次, 分别来自 {_top["n_files"]} 个不同的'
- f'工作记事文件 —— 每份工作记事都附一份当前未结束长停的全量台账') if _top else '',
- f'不归并: {_days_raw:.0f} 天 = 台账窗内全场时长的 {_pct_raw}%, 与可用率 94.8% 直接矛盾',
- f'归并后: {_days_ded:.0f} 天 = {_pct_ded}%, 与可用率相容',
- f'★口径: 分母 = 38 台 × 台账窗 {_span} 天 (首条停机至今)。若改用场龄 6.53 年作分母则为 '
- f'{_pct_raw_age}% —— 但台账窗始于 2022-12, 用场龄会让分子分母窗不对齐、占比被低估',
- f'★两个比不同: 行数比 {_rowr}× (2277/59) vs 停机天数虚高 {_infl}× '
- f'({_days_raw:.0f}/{_days_ded:.0f})。首版标题误用行数比来描述天数'],
- evidence_en=[
- f'{n_raw} rows for this site collapse to {len(d)} unique events when de-duplicated on '
- f'(unit, stop time, resume time)',
- (f'The most-repeated entry appears {_top["n_files"]} times, once in each of '
- f'{_top["n_files"]} different work-log files, because every work log carries a full '
- f'copy of the currently-open long-stop register') if _top else '',
- f'Without de-duplication: {_days_raw:.0f} days = {_pct_raw}% of fleet time, which '
- f'directly contradicts the 94.8% availability figure',
- f'After de-duplication: {_days_ded:.0f} days = {_pct_ded}%, consistent with availability',
- f'Basis: the denominator is 38 units x {_span} register days (first stop to date). '
- f'Using the {6.53}-year age of the site instead gives {_pct_raw_age}%, but the register '
- f'only starts in December 2022, so that denominator does not match the numerator window',
- f'Two different ratios: rows {_rowr}x (2277/59) versus downtime days {_infl}x '
- f'({_days_raw:.0f}/{_days_ded:.0f}). The first version used the row ratio in a headline '
- f'that named downtime days'],
- discriminator_zh='**同一业务键重复出现且各来自不同 src_file ⇒ 快照类, 必须归并取末值**',
- discriminator_en='If the same business key repeats with each occurrence coming from a '
- 'different source file, the table is a cumulative snapshot and must be '
- 'collapsed to the latest state per key',
- related='与 SOP §4.11 EL-6 记录的"工单损失列=累计逐日快照, 直接 sum 虚高 3.4×"同一陷阱 '
- '(两个倍数各属不同表, 不是同一个数的两种说法)',
- caveat=(f'{len(d)} 是**归并后的台账行数**, 不等于"本场五年只发生过 {len(d)} 次长停" —— '
- f'若某条长停的起止时间在不同快照里被改写过, 会被判成多条; 反之跨快照周期的'
- f'同一次长停若键值一致则合为一条。真实长停次数可能与此不同。'
- f'另: 归并取末值, 对在最后一份快照前已结束的事件, 其结束时间取自最后一次出现的记录, '
- f'可能早于真实复役时间。'),
- caveat_en=(f'{len(d)} is the number of rows after de-duplication, not a claim that only '
- f'{len(d)} long stops occurred at this site in five years: an event whose start '
- f'or end time was edited between snapshots would be counted more than once, '
- f'while one that kept identical key values across snapshot periods collapses to '
- f'a single row. The true number of long stops may differ. Collapsing also keeps '
- f'the state from the latest snapshot in which each event appears, so for an '
- f'event that ended earlier the recorded resume time may precede the true one.'),
- provenance='scripts/ingest_ops_2025.py (数字实时算, 可复现)',
- ), links=dict(about=['farm/rudong']))
- print(f' → finding cumulative_snapshot_inflation: {n_raw}行/{_nfile}文件 → {len(d)}条, '
- f'行数比 {_rowr}× / 天数虚高 {_infl}× ({_pct_ded}% → {_pct_raw}%; 场龄口径 {_pct_raw_age}%)')
- n_ls = 0
- for i, r in d.iterrows():
- tid = r.get('tid')
- t = f'WTG{int(tid):02d}' if pd.notna(tid) else None
- oid = f'longstop/{t or "unknown"}/{i}'
- db[oid] = dict(id=oid, type='Evidence', props=dict(
- name=f'{t or "?"} 长停登记', facet='外部记录',
- claim_window=f"{_s(r.get('停机时间'))} ~ {_s(r.get('恢复时间/预计投运时间'))}",
- source='工作库/20_structured/ops/longstop_register.parquet (仅如海)',
- turbine=t, stop_time=_s(r.get('停机时间')),
- resume_time=_s(r.get('恢复时间/预计投运时间')),
- stop_class=_s(r.get('停机分类')),
- stop_days=None if pd.isna(r.get('停机时长(天)')) else float(r['停机时长(天)'])),
- links=dict(about=[f'turbine/{t}']) if t else {})
- n_ls += 1
- added['Evidence(长停登记)'] = n_ls
- for k, v in added.items():
- print(f' {k:26s} {v:5d}')
- print(f' 本体 {n0} → {len(db)}')
- if a.dry_run:
- print(' (dry-run, 未写盘)')
- return 0
- # 走本体 store 落盘 (2026-09-16): 原来直接 `OBJ.write_text(...)` —— 绕过了 store 的
- # ① 全库 validate 闸 ② F12 原子写 (tmp + os.replace)。对象库里混进一个坏对象时, 直写会
- # 让**整库**在下次读取时才炸, 而 store.save() 会当场报出是哪个对象不合法。
- from src.ontology.store import Store
- st = Store(OBJ)
- st.objects = db
- st.save()
- print(f' 已写入 {_P.rel(OBJ)} (经 Store 闸)')
- return 0
- if __name__ == '__main__':
- sys.exit(main())
|