products_restore_missing.py 6.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135
  1. #!/usr/bin/env python3
  2. # -*- coding: utf-8 -*-
  3. r"""补齐"包内没有生成端"的产物 —— 只补缺件, 不动 raw 重算件; 并落一份**来源台账** (2026-09-12)。
  4. ## 为什么需要它
  5. `rebuild_from_raw.py` 能从 `data/raw` 重算出一部分产物 (L0 仓 16 件 + 本体 3 件); 但随包里另有
  6. 一批产物**全库只有读取方、0 处写入方** —— 没有生成端, 从零重算就是拿不到 (工作台主数据
  7. `/api/fleet` 一缺 `pitch/pitch_daily.parquet` / `temp_monthly.parquet` 就整个回 `err=no_products`,
  8. 页面看着像"没数据")。
  9. 用户令 (2026-09-12): **既要"从零重算", 又要"页面不缺"**。两者只能这样同时成立:
  10. · 能证明是"从 raw 重算出来的"→ 用重算件 (并且逐值对齐随包件才敢用);
  11. · 证明不了的 (没有生成端 / 规则复现不出) → 用随包件补齐, **但必须标明它不是重算件**。
  12. 本脚本就是后者的执行者, 并把每一件的来源写进 `outputs/<场>/_provenance.json`。
  13. ## 规则 (保守优先)
  14. 对随包产物目录里的每个文件:
  15. · 目标**已存在** → 跳过 (那是 raw 重算出来的, 不许被随包件覆盖);
  16. · 目标不存在 → 从随包件拷过去, 记为 `shipped`;
  17. 另: 只处理"产物目录"(outputs/<场>/ 下的仓), 不碰 release/ 与 data/。
  18. ## 用法
  19. python scripts/products_restore_missing.py --dry-run # 只报要补什么
  20. python scripts/products_restore_missing.py # 真补 + 写台账
  21. python scripts/products_restore_missing.py --stash <dir> # 指定随包件所在目录
  22. """
  23. from __future__ import annotations
  24. import argparse
  25. import json
  26. import pathlib
  27. import shutil
  28. import sys
  29. import time
  30. ROOT = pathlib.Path(__file__).resolve().parents[1]
  31. sys.path.insert(0, str(ROOT))
  32. from src import paths as P # noqa: E402
  33. # 由 data/raw 重算出来的 (重建器 + 验证依据) —— 这些**永不被随包件覆盖**
  34. RAW_DERIVED = {
  35. # 三门台账 (rebuild_from_raw.py; --verify 与随包基线逐值比对过)
  36. 'windscada/alarms.parquet': 'scripts/windscada_alarms_ingest.py',
  37. 'windscada/workorders.parquet': 'scripts/windscada_workorder_ingest.py',
  38. 'windscada/oil_samples_index.parquet': 'scripts/windscada_watch_channels_build.py',
  39. # SCADA 侧 10 个构建器 (rebuild_from_raw.py --scada)
  40. 'windscada/powercurve_dev.parquet': 'src.windscada.perf.powercurve',
  41. 'windscada/powercurve_bins.parquet': 'src.windscada.perf.powercurve',
  42. 'windscada/loss_monthly.parquet': 'src.windscada.perf.availability',
  43. 'windscada/curve_lenses.parquet': 'src.windscada.perf.curves',
  44. 'windscada/curve_liveness.parquet': 'src.windscada.perf.curves',
  45. 'windscada/control_profile.parquet': 'src.windscada.perf.control',
  46. 'windscada/control_schedule.parquet': 'src.windscada.perf.control',
  47. 'windscada/stop_events.parquet': 'src.windscada.perf.faults',
  48. 'windscada/temp_bins.parquet': 'src.windscada.subsys.temp_nbm',
  49. 'windscada/yaw_daily.parquet': 'src.windscada.subsys.yaw',
  50. 'windscada/hydraulic_accum.parquet': 'src.windscada.subsys.hydraulic',
  51. 'windscada/thermal_chain.parquet': 'src.windscada.subsys.thermal_chain',
  52. 'windscada/system_aux.parquet': 'src.windscada.taxonomy',
  53. # 月度派生件 (本轮新增; 逐值对齐随包件后才算数)
  54. 'windscada/temp_monthly.parquet': 'scripts/windscada_monthly_build.py (19494/19494 逐值一致)',
  55. # 本体层 (kb_ingest → populate → chain_ingest → trend_ingest; 全链跑完后 audit 0 问题)
  56. 'ontology/objects.json': 'python -m src.ontology.kb_ingest + populate + chain_ingest + trend_ingest (audit 0 问题)',
  57. 'ontology/retrieval_index.json': 'src.ontology.retrieval.build',
  58. 'ontology/turbine_params.parquet': 'src.ontology.maintenance.refresh_params',
  59. }
  60. def stash_dir(explicit=None) -> pathlib.Path:
  61. """随包产物暂存目录 (products_state --off 挪走后的地方)。"""
  62. if explicit:
  63. return pathlib.Path(explicit)
  64. for p in sorted((ROOT / '_products_off').rglob('rudong')):
  65. if p.is_dir():
  66. return p
  67. raise SystemExit('找不到随包产物暂存目录 (_products_off/**/rudong); 用 --stash 指定')
  68. def main() -> int:
  69. ap = argparse.ArgumentParser()
  70. ap.add_argument('--dry-run', action='store_true')
  71. ap.add_argument('--stash', default=None)
  72. a = ap.parse_args()
  73. stash = stash_dir(a.stash)
  74. dest_root = P.STORE if hasattr(P, 'STORE') else ROOT / 'outputs' / 'rudong'
  75. print(f'随包件: {stash.relative_to(ROOT)}')
  76. print(f'产物仓: {dest_root.relative_to(ROOT)}\n')
  77. prov = {}
  78. for src in sorted(stash.rglob('*')):
  79. if not src.is_file():
  80. continue
  81. rel = src.relative_to(stash).as_posix()
  82. dst = dest_root / rel
  83. if dst.exists():
  84. prov[rel] = dict(source='raw-derived', builder=RAW_DERIVED.get(rel, '(早期重算, 未登记)'))
  85. continue
  86. prov[rel] = dict(source='shipped', why='包内无生成端 / 规则未复现 → 用随包件补齐')
  87. if not a.dry_run:
  88. dst.parent.mkdir(parents=True, exist_ok=True)
  89. shutil.copy2(src, dst)
  90. by_store = {}
  91. for rel, m in prov.items():
  92. store = rel.split('/')[0]
  93. d = by_store.setdefault(store, {'raw-derived': 0, 'shipped': 0})
  94. d[m['source']] += 1
  95. print(f'{"仓":26s} {"raw 重算":>9s} {"随包补齐":>9s}')
  96. for store, d in sorted(by_store.items()):
  97. print(f' {store:24s} {d["raw-derived"]:9d} {d["shipped"]:9d}')
  98. tot = {k: sum(d[k] for d in by_store.values()) for k in ('raw-derived', 'shipped')}
  99. print(f' {"合计":24s} {tot["raw-derived"]:9d} {tot["shipped"]:9d}')
  100. if a.dry_run:
  101. print('\n(dry-run, 未写盘)')
  102. return 0
  103. out = dest_root / '_provenance.json'
  104. out.write_text(json.dumps(dict(
  105. at=time.strftime('%Y-%m-%d %H:%M:%S'),
  106. note='逐件来源台账: raw-derived = 由 data/raw 重算(含验证依据); shipped = 包内无生成端, 用随包件补齐',
  107. counts=tot, files=prov), ensure_ascii=False, indent=1), encoding='utf-8')
  108. print(f'\n已写来源台账 {P.rel(out)}')
  109. return 0
  110. if __name__ == '__main__':
  111. # 控制台可能是 GBK: print 里的非 GBK 字形编不出来会抛异常, 干成了却退出码 1
  112. for _s in (sys.stdout, sys.stderr):
  113. try: _s.reconfigure(errors='replace')
  114. except Exception: pass
  115. sys.exit(main())