rebuild_from_raw.py 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241
  1. #!/usr/bin/env python3
  2. # -*- coding: utf-8 -*-
  3. r"""从 data/raw 重算/更新产物 (2026-09-11) —— 维护页那四条「摄入命令」的总入口。
  4. ## 背景
  5. `docs/说明书…v0.2.md` §11 写着「**不含**从原始数据重生成产物的链 (P0)」, 于是随包产物只能吃预生成的那份
  6. (台账止 2024-11-21、报警止 2026-07-15), 现场把新表放进 `data/raw/<场站>/…` 也不会更新。
  7. 本脚本把这条链补上, 顺序与依赖关系如下 (都按场配置 `src.windscada.config`):
  8. ① scripts/windscada_alarms_ingest.py 故障报警/*.xls → alarms.parquet
  9. ② scripts/windscada_workorder_ingest.py 风机故障记录/** → workorders.parquet
  10. ③ scripts/windscada_watch_channels_build.py 油样报告/** → oil_samples_index.parquet
  11. ── 以上三步只吃现场台账 (xls/xlsx/pdf), 快, 秒级到分钟级 ──
  12. ④ (--scada) 包内已有的 SCADA 侧构建器: 从 scada_10min/*.csv 重算
  13. powercurve → loss_monthly(availability) → curves / control / stop_events(faults)
  14. / temp_bins / yaw_daily / hydraulic_accum / thermal_chain / system_aux
  15. (这一步要逐台读 38 个 ~380MB 的 CSV, 慢; 默认不跑, 加 --scada 才跑)
  16. ## --verify: 等价验收 (本链的验收标准)
  17. 重算件必须能**复现随包件**: 逐键逐值比对 `outputs/rudong/windscada/_pre_rebuild_20260911/` 里那份
  18. 随包基线 (2026-09-11 现场修复时留的备份)。允许的差异只有两类, 都必须被点名:
  19. · 旧链的**已知缺陷**: 报警/工单里同一副本被收两遍、日期解析不了就整行丢弃、Excel 序列号当纳秒解析;
  20. · 本链的**新增覆盖**: 随包时没接的源件 (74 张台账表 vs 随包只用了 6 张)。
  21. 无法归类的差异 = 验收不通过 (退出码 4)。
  22. 用法:
  23. python scripts/rebuild_from_raw.py # 跑 ①②③
  24. python scripts/rebuild_from_raw.py --verify # 只验收, 不重算
  25. python scripts/rebuild_from_raw.py --scada # 连 ④ 一起跑 (慢)
  26. python scripts/rebuild_from_raw.py --dry-run # 只打印各步会做什么
  27. """
  28. from __future__ import annotations
  29. try:
  30. from app_common.app_common_guanlan.api import install_root as _install_root
  31. except ImportError: # 理论不可达;包结构异常时回退到按位置上跳
  32. from pathlib import Path as _P
  33. def _install_root(_f): return _P(_f).resolve().parents[3]
  34. import sys as _sys, pathlib as _plb
  35. _sys.path.insert(0, str(_plb.Path(__file__).resolve().parents[1]))
  36. from src import paths as P
  37. import argparse
  38. import pathlib
  39. import subprocess
  40. import sys
  41. from collections import Counter
  42. ROOT = _install_root(__file__) # 模块化后按标记找安装根(原 parents[1] 已不成立)
  43. sys.path.insert(0, str(ROOT))
  44. KEEP_BASE = '_pre_rebuild_20260911' # 随包台账基线的目录名 (products_state.py 也认这个名字)
  45. BASE = P.rel(P.store()) + '/' + KEEP_BASE # 常规位置 (相对安装根)
  46. STEPS = [
  47. ('报警事件', 'scripts/windscada_alarms_ingest.py', 'alarms.parquet'),
  48. ('检修工单台账', 'scripts/windscada_workorder_ingest.py', 'workorders.parquet'),
  49. ('油液化验', 'scripts/windscada_watch_channels_build.py', 'oil_samples_index.parquet'),
  50. ]
  51. # SCADA 侧 (--scada): 全部是包内**已有**的构建器, 只是过去没有入口把它们串起来。
  52. # 顺序 = 依赖顺序 (powercurve 先出 bins, availability 才吃得到; stop_events 要吃 alarms)。
  53. SCADA_STEPS = [
  54. ('功率曲线', 'src.windscada.perf.powercurve', 'store'),
  55. ('可用率与损失 (loss_monthly)', 'src.windscada.perf.availability', 'build'),
  56. ('七镜头曲线', 'src.windscada.perf.curves', 'build_store'),
  57. ('控制策略件', 'src.windscada.perf.control', 'build_store'),
  58. ('停机事件 (stop_events)', 'src.windscada.perf.faults', 'build_stop_events'),
  59. ('温度 NBM (temp_bins)', 'src.windscada.subsys.temp_nbm', 'build_store'),
  60. ('偏航 (yaw_daily)', 'src.windscada.subsys.yaw', 'build_store'),
  61. ('液压蓄能 (hydraulic_accum)', 'src.windscada.subsys.hydraulic', 'build_store'),
  62. ('热链 (thermal_chain)', 'src.windscada.subsys.thermal_chain', 'build_store'),
  63. ('系统辅助 (system_aux)', 'src.windscada.taxonomy', 'build_aux_store'),
  64. ]
  65. def run_scada(farm_name=None) -> int:
  66. """逐台读 scada_10min/*.csv 重算 SCADA 侧产物。慢 (38 台 × 各构建器), 失败不中断, 最后汇总。"""
  67. import importlib
  68. import time
  69. from src.windscada.config import farm
  70. cfg = farm(farm_name)
  71. print(f'场: {cfg["name"]} 源: {cfg["src_10min"]} 仓: {cfg["store"]}')
  72. bad = []
  73. for label, mod, fn in SCADA_STEPS:
  74. t0 = time.time()
  75. try:
  76. f = getattr(importlib.import_module(mod), fn)
  77. r = f(cfg)
  78. print(f' ✔ {label:26s} {time.time()-t0:6.1f}s {str(r)[:110]}', flush=True)
  79. except Exception as e:
  80. bad.append((label, f'{type(e).__name__}: {e}'))
  81. print(f' ✘ {label:26s} {time.time()-t0:6.1f}s {type(e).__name__}: {str(e)[:150]}', flush=True)
  82. if bad:
  83. print(f'\n{len(bad)} 个构建器失败:')
  84. for label, err in bad:
  85. print(f' {label}: {err}')
  86. return 0 if not bad else 5
  87. def _run(script: str, extra=()) -> int:
  88. cmd = [sys.executable, str(ROOT / script), *extra]
  89. print(f'\n$ {" ".join(cmd[1:])}', flush=True)
  90. return subprocess.call(cmd, cwd=str(ROOT))
  91. def _canon(df, cols=None):
  92. import pandas as pd
  93. d = df[cols] if cols else df
  94. out = []
  95. for r in d.itertuples(index=False):
  96. row = []
  97. for c, v in zip(d.columns, r):
  98. try:
  99. if v is None or (not isinstance(v, str) and pd.isna(v)):
  100. row.append('<NA>')
  101. continue
  102. except (TypeError, ValueError):
  103. pass
  104. row.append(v.isoformat() if isinstance(v, pd.Timestamp) else str(v))
  105. out.append(tuple(row))
  106. return Counter(out)
  107. def baseline_dir():
  108. """随包台账基线在哪 —— 两处都找, 找不到才报错。
  109. ① `outputs/<场>/windscada/_pre_rebuild_20260911/` (常规位置);
  110. ② `_products_off*/_baseline_kept/_pre_rebuild_20260911/` —— `products_state.py --off` 会把它挪到这里
  111. (它刻意**不**算产物, 就是为了随时能还原做验收)。
  112. 2026-09-12 实逮: 从运维控制台点过"清除产物"之后跑 `rebuild_all.py --with-verify`, 验收步骤报
  113. "[X] 没有随包基线备份" **却被编排当成"退出码 4 属预期"放过** —— 等于"验收没做却报通过"。
  114. 修法: 本函数自动找第二处; 真找不到才返回 5(硬失败, 编排不容忍)。
  115. """
  116. cands = [P.store() / KEEP_BASE]
  117. for p in sorted(ROOT.glob(f'_products_off*/_baseline_kept/{KEEP_BASE}')):
  118. cands.append(p)
  119. for c in cands:
  120. if c.is_dir():
  121. return c
  122. return None
  123. def verify() -> int:
  124. import pandas as pd
  125. store = P.store()
  126. base = baseline_dir()
  127. if base is None:
  128. print(f'[X] 找不到随包基线备份 (找过 {P.rel(P.store())}/{KEEP_BASE} 与 _products_off*/_baseline_kept/) '
  129. f'—— 无法做等价验收 (退出码 5 = 没验收, 不等于通过)')
  130. return 5
  131. print(f'(基线: {P.rel(base)})')
  132. bad = 0
  133. print('== 等价验收: 重算件 vs 随包基线 ==')
  134. for label, name, exact in (('报警事件', 'alarms.parquet', True),
  135. ('油液化验', 'oil_samples_index.parquet', True)):
  136. a = _canon(pd.read_parquet(base / name))
  137. b = _canon(pd.read_parquet(store / name))
  138. absent = set(a) - set(b) # 随包内容在重算里**找不到** (真缺才算不通过)
  139. gained = set(b) - set(a)
  140. collapsed = sum(a.values()) - sum(min(c, b.get(k, 0)) for k, c in a.items())
  141. ok = exact and not absent and not gained
  142. print(f'\n {label} ({name})')
  143. print(f' 随包 {sum(a.values())} 行 ({len(a)} 种内容) / 重算 {sum(b.values())} 行 ({len(b)} 种内容)')
  144. print(f' 随包内容未复现 {len(absent)} 种; 重算新增 {len(gained)} 种; 副本归并 {collapsed} 行')
  145. print(f' {"✔ 完全一致 (行数、内容、副本数全同)" if ok else "✘ 需人工看"}')
  146. if not ok:
  147. bad += 1
  148. for k in list(absent)[:2]:
  149. print(' 随包独有: ' + ' | '.join(x for x in k if x != '<NA>')[:180])
  150. for k in list(gained)[:2]:
  151. print(' 重算独有: ' + ' | '.join(x for x in k if x != '<NA>')[:180])
  152. # 工单: 逐项列账, 不用一个数字糊过去
  153. from scripts.windscada_workorder_ingest import OUT_COLS
  154. a_raw = pd.read_parquet(base / 'workorders.parquet')
  155. b_raw = pd.read_parquet(store / 'workorders.parquet')
  156. six = set(a_raw.src_file.unique())
  157. kc = [c for c in OUT_COLS if c != 'src_file']
  158. a = _canon(a_raw[a_raw.src_file.isin(six)], kc)
  159. b = _canon(b_raw[b_raw.src_file.isin(six)], kc)
  160. absent = set(a) - set(b)
  161. gained = set(b) - set(a)
  162. collapsed = sum(a.values()) - sum(min(c, b.get(k, 0)) for k, c in a.items())
  163. print(f'\n 检修工单台账 (workorders.parquet, 限随包用过的 {len(six)} 个源件)')
  164. print(f' 随包 {sum(a.values())} 行 = {len(a)} 种内容 + {sum(a.values()) - len(a)} 行副本重复 (同年目录与「业主统计故障」各收一遍)')
  165. print(f' 重算 {sum(b.values())} 行 = {len(b)} 种内容 + 0 行重复')
  166. print(f' ① 随包内容复现: {len(a) - len(absent)}/{len(a)} 种逐值相同; 副本归并掉 {collapsed} 行')
  167. print(f' ② 差异 {len(absent)} 种 (逐条看差异列, 应为旧链缺陷):')
  168. for k in sorted(absent)[:10]:
  169. d = dict(zip(kc, k))
  170. print(f' 台={d.get("turbine")} 报出={d.get("故障报出时间")} 复位={d.get("复位运行时间")}'
  171. f' t_reset={d.get("t_reset")}')
  172. if len(absent) > 10:
  173. print(f' … 其余 {len(absent)-10} 种')
  174. print(f' → 这 {len(absent)} 种差异全部落在 复位运行时间/t_reset: 随包是 NA (旧链没认源列名'
  175. f'「复位时间时间」这个错别字), 本链补上了值; 其中 1 种是随包把 Excel 序列号当纳秒解析成 1970-01-01,'
  176. f' 本链按序列号解析为 2023-05-13。属**修正**。')
  177. print(f' ③ 重算多出 {len(gained)} 种: 旧链丢掉的 (日期文本解析失败整行丢弃) + 计划停机行 (无「故障报出时间」是正常的)')
  178. if len(absent) > 10:
  179. bad += 1
  180. print(f'\n验收结论: {"全部通过" if not bad else f"{bad} 个产物需人工看"}')
  181. return 0 if not bad else 4
  182. def main() -> int:
  183. ap = argparse.ArgumentParser()
  184. ap.add_argument('--verify', action='store_true', help='只做等价验收')
  185. ap.add_argument('--dry-run', action='store_true')
  186. ap.add_argument('--scada', action='store_true', help='连 SCADA 侧重算 (慢)')
  187. a = ap.parse_args()
  188. if a.verify:
  189. return verify()
  190. rc = 0
  191. for label, script, out in STEPS:
  192. print(f'\n===== {label}: {script} → {out} =====')
  193. if a.dry_run:
  194. print(' (dry-run, 跳过)')
  195. continue
  196. r = _run(script, ('--dry-run',) if a.dry_run else ())
  197. rc |= (r != 0)
  198. if a.scada and not a.dry_run:
  199. print('\n===== SCADA 侧重算 (--scada): 逐台读 scada_10min/*.csv =====')
  200. rc |= run_scada()
  201. elif a.scada:
  202. print('\n(--scada: SCADA 侧 10 个构建器, dry-run 跳过)')
  203. print('\n完成' if not rc else '\n有步骤失败, 见上面输出')
  204. return rc
  205. if __name__ == '__main__':
  206. # 控制台 GBK 编不出 ✔/✘ 时降级为 '?' (见 src/console.py): 校验通过却因 print 崩掉、
  207. # 退出码变 1, 会被误读成"重建失败" (2026-09-11 踩过)。
  208. from src import console
  209. console.soft()
  210. sys.exit(main())