| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227 |
- #!/usr/bin/env python3
- # -*- coding: utf-8 -*-
- r"""补齐"包内没有生成端"的产物 —— 只补缺件, 不动 raw 重算件; 并落一份**来源台账** (2026-09-12)。
- ## 为什么需要它
- `rebuild_from_raw.py` 能从 `data/raw` 重算出一部分产物 (L0 仓 16 件 + 本体 3 件); 但随包里另有
- 一批产物**全库只有读取方、0 处写入方** —— 没有生成端, 从零重算就是拿不到 (工作台主数据
- `/api/fleet` 一缺 `pitch/pitch_daily.parquet` / `temp_monthly.parquet` 就整个回 `err=no_products`,
- 页面看着像"没数据")。
- 用户令 (2026-09-12): **既要"从零重算", 又要"页面不缺"**。两者只能这样同时成立:
- · 能证明是"从 raw 重算出来的"→ 用重算件 (并且逐值对齐随包件才敢用);
- · 证明不了的 (没有生成端 / 规则复现不出) → 用随包件补齐, **但必须标明它不是重算件**。
- 本脚本就是后者的执行者, 并把每一件的来源写进 `outputs/<场>/_provenance.json`。
- ## 规则 (保守优先)
- 对随包产物目录里的每个文件:
- · 目标**已存在** → 跳过 (那是 raw 重算出来的, 不许被随包件覆盖);
- · 目标不存在 → 从随包件拷过去, 记为 `shipped`;
- 另: 只处理"产物目录"(outputs/<场>/ 下的仓), 不碰 release/ 与 data/。
- ## 用法
- python scripts/products_restore_missing.py --dry-run # 只报要补什么
- python scripts/products_restore_missing.py # 真补 + 写台账
- python scripts/products_restore_missing.py --stash <dir> # 指定随包件所在目录
- """
- from __future__ import annotations
- import argparse
- import json
- import pathlib
- import shutil
- import sys
- import time
- ROOT = pathlib.Path(__file__).resolve().parents[1]
- sys.path.insert(0, str(ROOT))
- from src import paths as P # noqa: E402
- # 由 data/raw 重算出来的 (重建器 + 验证依据) —— 这些**永不被随包件覆盖**
- RAW_DERIVED = {
- # 三门台账 (rebuild_from_raw.py; --verify 与随包基线逐值比对过)
- 'windscada/alarms.parquet': 'scripts/windscada_alarms_ingest.py',
- 'windscada/workorders.parquet': 'scripts/windscada_workorder_ingest.py',
- 'windscada/oil_samples_index.parquet': 'scripts/windscada_watch_channels_build.py',
- # SCADA 侧 10 个构建器 (rebuild_from_raw.py --scada)
- 'windscada/powercurve_dev.parquet': 'src.windscada.perf.powercurve',
- 'windscada/powercurve_bins.parquet': 'src.windscada.perf.powercurve',
- 'windscada/loss_monthly.parquet': 'src.windscada.perf.availability',
- 'windscada/curve_lenses.parquet': 'src.windscada.perf.curves',
- 'windscada/curve_liveness.parquet': 'src.windscada.perf.curves',
- 'windscada/control_profile.parquet': 'src.windscada.perf.control',
- 'windscada/control_schedule.parquet': 'src.windscada.perf.control',
- 'windscada/stop_events.parquet': 'src.windscada.perf.faults',
- 'windscada/temp_bins.parquet': 'src.windscada.subsys.temp_nbm',
- 'windscada/yaw_daily.parquet': 'src.windscada.subsys.yaw',
- 'windscada/hydraulic_accum.parquet': 'src.windscada.subsys.hydraulic',
- 'windscada/thermal_chain.parquet': 'src.windscada.subsys.thermal_chain',
- 'windscada/system_aux.parquet': 'src.windscada.taxonomy',
- # 月度派生件 (本轮新增; 逐值对齐随包件后才算数)
- 'windscada/temp_monthly.parquet': 'scripts/windscada_monthly_build.py (19494/19494 逐值一致)',
- # 本体层 (kb_ingest → populate → chain_ingest → trend_ingest; 全链跑完后 audit 0 问题)
- 'ontology/objects.json': 'python -m src.ontology.kb_ingest + populate + chain_ingest + trend_ingest (audit 0 问题)',
- 'ontology/retrieval_index.json': 'src.ontology.retrieval.build',
- 'ontology/turbine_params.parquet': 'src.ontology.maintenance.refresh_params',
- }
- def stash_dir(explicit=None) -> pathlib.Path:
- """随包产物原件在哪。
- 三个地方都可能是它(取决于用过哪些开关), 所以按"像不像随包整份产物"打分挑:
- · `_products_off/<场>/` = products_state --off 挪走的当前产物;
- · `_products_off_prev_<时间戳>/<场>/` = --off --archive-old 存档的上一代(= 随包原件);
- · `--stash` 显式指定。
- 判据用**只在随包里有的件**(windscada/temp_monthly.parquet 这类"包内无生成端"的产物):
- 含它的目录才是"标准答案"来源。原先只看第一个候选, 于是 --on 还原后暂存区空了 → 台账被写成 0/0。
- """
- if explicit:
- # 必须 resolve: 传相对路径时下面 `stash.relative_to(ROOT)` 会 ValueError (2026-09-12 实逮)
- return pathlib.Path(explicit).resolve()
- marker = 'windscada/temp_monthly.parquet'
- cands = []
- for pat in ('_products_off/*', '_products_off_prev_*/*'):
- cands += [p for p in sorted(ROOT.glob(pat)) if p.is_dir() and p.name != '_baseline_kept']
- scored = []
- for c in cands:
- n = sum(1 for _ in c.rglob('*') if _.is_file())
- scored.append(((1 if (c / marker).is_file() else 0), n, c))
- scored.sort(key=lambda x: (-x[0], -x[1]))
- if scored:
- return scored[0][2]
- raise StashMissing('找不到随包产物原件 (_products_off/** 或 _products_off_prev_**); 用 --stash 指定')
- class StashMissing(SystemExit):
- """暂存区不存在 —— **不是崩溃, 是"这一步没得做"**。
- 2026-09-16 用户令"系统运行中/未运行都能重算"时实逮: `rebuild_all.py` 的第⑤步 (补齐"包内没有生成端"
- 的产物) 在**没有暂存区**的检出上必抛 SystemExit → 重建链当场中断, 后面的 ⑦本体链 与 ⑧审计**全不执行**
- (脚本的中断语义是"后面的步骤依赖它")。而暂存区只有"清除产物"动作会产生 —— 首次全量重算、或
- 清完又还原过的机器上就都没有, 于是"从零重算"在交付包上跑不完。
- 处置: 本类带退出码 6, 由 rebuild_all.py 显式容忍并打印"跳过随包补齐" —— 缺的是"随包件"这一路,
- 不是重算本身; 如实报出来比假装通过好, 也比把整条链打断好。
- """
- exit_code = 6
- def main() -> int:
- ap = argparse.ArgumentParser()
- ap.add_argument('--dry-run', action='store_true')
- ap.add_argument('--stash', default=None)
- a = ap.parse_args()
- try:
- stash = stash_dir(a.stash)
- except StashMissing as e:
- # 没有暂存区 → 这一步没得做。**不许静默、也不许把整条重算链打断**(见 StashMissing 的说明)。
- print(f'[跳过] {e}')
- print(' 处置: ① 从**交付包 zip** 按需补齐: --stash <交付包.zip> (推荐; 2026-09-16 用户令'
- '"清除产物不留备份"之后, 交付包就是随包件的来源); '
- '② --stash <随包产物目录> 显式指定一个目录; ③ 只想让页面有数 → 产物就在位, 无需本步。')
- return StashMissing.exit_code
- # ★2026-09-16: 原写法 `P.STORE if hasattr(P,'STORE') else ROOT/'outputs'/'rudong'` —— P.STORE 根本
- # 不存在, 于是**恒**落到硬编码的 rudong, 多场部署下会把 rudong 的随包件补进别的场 (静默串场)。
- dest_root = P.out_root()
- print(f'随包件: {P.rel(stash)}')
- print(f'产物仓: {dest_root.relative_to(ROOT)}\n')
- prov = {}
- def take(rel: str, src_path=None, src_bytes=None):
- """补齐一件 (或在位则只记账)。src_path=目录来源, src_bytes=zip 来源。"""
- dst = dest_root / rel
- if dst.exists():
- # 已存在的件分两类: 我们自己重算的(raw-derived) 与 上轮已补齐的(shipped)。
- # ★不能因为"本轮没拷"就不记账 —— 早先写成"存在即不列", 于是重复跑一次台账就被清成 0/0,
- # 页面上的来源统计跟着一起错 (2026-09-12 实逮)。
- if rel in RAW_DERIVED:
- prov[rel] = dict(source='raw-derived', builder=RAW_DERIVED[rel])
- else:
- prov[rel] = dict(source='shipped', why='包内无生成端 / 规则未复现 → 随包件补齐 (本轮已在位)')
- return
- prov[rel] = dict(source='shipped', why='包内无生成端 / 规则未复现 → 用随包件补齐')
- if a.dry_run:
- return
- dst.parent.mkdir(parents=True, exist_ok=True)
- if src_bytes is not None:
- dst.write_bytes(src_bytes)
- else:
- shutil.copy2(src_path, dst)
- if stash.is_file():
- # 交付包 zip 形态 (2026-09-16): 按需从包里取 outputs/<场>/** 补齐 ——
- # 这是"清除不留备份"之后**唯一**的随包件来源, 所以把它做成一等输入。
- import zipfile
- with zipfile.ZipFile(stash) as zf:
- pref = f'outputs/{P.farm()}/'
- members = [n for n in zf.namelist() if n.startswith(pref) and not n.endswith('/')]
- print(f' 从交付包读取 {len(members)} 件 (前缀 {pref})')
- for name in sorted(members):
- if a.dry_run:
- take(name[len(pref):])
- else:
- take(name[len(pref):], src_bytes=zf.read(name))
- else:
- for src in sorted(stash.rglob('*')):
- if not src.is_file():
- continue
- take(src.relative_to(stash).as_posix(), src_path=src)
- # 重算件里有些**不在随包件里**(如我们新造的 turbine_params.parquet), 也要记进台账
- for rel, builder in RAW_DERIVED.items():
- if rel not in prov and (dest_root / rel).exists():
- prov[rel] = dict(source='raw-derived', builder=builder)
- # 构建脚本的**自登记** (src/derived_manifest.py): 振动侧的产物名随窗/分片变, 写不进上面的精确表,
- # 按名字通配又会误伤同名旧件 → 由"谁算的谁登记", 这里只认登记。缺失不影响其余记账。
- try:
- from src.derived_manifest import load as _load_manifest, prune as _prune_manifest
- _gone = _prune_manifest(dest_root) # 先清掉"登记了但盘上已删"的条目 (如被删的 _reimport 窗)
- if _gone:
- print(f' 自登记清理: {_gone} 条指向已删文件的条目')
- _man = (_load_manifest(dest_root).get('files') or {})
- except Exception:
- _man = {}
- for rel, info in _man.items():
- if (dest_root / rel).exists():
- prov[rel] = dict(source='raw-derived',
- builder=f"{info.get('builder', '?')} (自登记 by {info.get('by', '?')})")
- else:
- prov[rel] = dict(source='raw-derived',
- builder=f"{info.get('builder', '?')} (自登记, 但该件当前不在盘上)")
- by_store = {}
- for rel, m in prov.items():
- store = rel.split('/')[0]
- d = by_store.setdefault(store, {'raw-derived': 0, 'shipped': 0})
- d[m['source']] += 1
- print(f'{"仓":26s} {"raw 重算":>9s} {"随包补齐":>9s}')
- for store, d in sorted(by_store.items()):
- print(f' {store:24s} {d["raw-derived"]:9d} {d["shipped"]:9d}')
- tot = {k: sum(d[k] for d in by_store.values()) for k in ('raw-derived', 'shipped')}
- print(f' {"合计":24s} {tot["raw-derived"]:9d} {tot["shipped"]:9d}')
- if a.dry_run:
- print('\n(dry-run, 未写盘)')
- return 0
- out = dest_root / '_provenance.json'
- out.write_text(json.dumps(dict(
- at=time.strftime('%Y-%m-%d %H:%M:%S'),
- note='逐件来源台账: raw-derived = 由 data/raw 重算(含验证依据); shipped = 包内无生成端, 用随包件补齐',
- counts=tot, files=prov), ensure_ascii=False, indent=1), encoding='utf-8')
- print(f'\n已写来源台账 {P.rel(out)}')
- return 0
- if __name__ == '__main__':
- # 控制台可能是 GBK: print 里的非 GBK 字形编不出来会抛异常, 干成了却退出码 1
- for _s in (sys.stdout, sys.stderr):
- try: _s.reconfigure(errors='replace')
- except Exception: pass
- sys.exit(main())
|