|
|
@@ -1,327 +1,432 @@
|
|
|
-#!/usr/bin/env python3
|
|
|
-# -*- coding: utf-8 -*-
|
|
|
-"""振动侧一键摄入: data/raw/<场站>/{windcms,m5_cms_tcm} 的原始件 → 窗索引/谱库 → (可选)CMS 报告与知识库.
|
|
|
-
|
|
|
-## 为什么需要它 (2026-09-12 用户令)
|
|
|
-
|
|
|
-用户令: 「数据层里的「CMS 振动评估报告」应遵循 `<安装目录>\\data\\raw\\如东\\windcms`、
|
|
|
-「振动线 handoff」应遵循 `<安装目录>\\data\\raw\\如东\\m5_cms_tcm` 存放; 修改系统支持振动数据参与
|
|
|
-运行、重算」。此前振动侧只有**产物**(`outputs/<场>/m5_cms_tcm` 90 件 + `outputs/<场>/windcms` 53 件),
|
|
|
-源件不在 data/raw, 生成端脚本 (振动线分支的 rudong_tcm_*.py 七件) 也没随包 ⇒ 从零重算时振动侧是断的。
|
|
|
-
|
|
|
-两个事实把这件事限定得很清楚:
|
|
|
- · `src/windcms/pipeline.py` 认的输入就是"含 `*_decode.json` 的目录"(Brande TCM Enterprise 导出),
|
|
|
- 这一步**必须能跑**, 否则数据进了 data/raw 也只是躺着;
|
|
|
- · 同分支的六层链 (oem_scan / energy_share / model_run / fusion) 脚本仍未随包 —— 本脚本**不假装**
|
|
|
- 能把那几步跑出来: 它只做"索引 + 谱 + (可选)报告/知识库", 其余在报告里如实写"缺失"。
|
|
|
-
|
|
|
-## 这一步跑完, 振动数据就真的参与了运行
|
|
|
-
|
|
|
- data/raw/<场站>/windcms/.../*_decode.json
|
|
|
- → outputs/<场>/m5_cms_tcm/windows/<窗>/index.parquet (54 列, 与包内 tcm_index.parquet 同构)
|
|
|
- → outputs/<场>/m5_cms_tcm/windows/<窗>/spectra/*.npz + spectra_meta.parquet
|
|
|
- 消费者 (无需改一行代码, 窗是自动发现的):
|
|
|
- src/windcms/data.py::windows() → 新窗进分析集
|
|
|
- src/windcms/data.py::load_scalars() → 标量进 CMS 报告/页面
|
|
|
- src/windcms/data.py::spectrum() → 谱图能取到 (此前 spectra 库整个缺失, 谱图是死的)
|
|
|
- scripts/windcms.py report / kb → 报告与知识库重生成
|
|
|
-
|
|
|
-## ★ 报告默认**跑**了 (2026-09-18 用户令后口径翻转)
|
|
|
-
|
|
|
-2026-09-12 的旧口径是"默认不跑": 那时拿**随包快照**当基准, 重生成会让 `windcms/report.md`
|
|
|
-从 16,857 B 掉到 467 B (−97%), 因为 `data.load_model()` 要读的 `m5/model_run_l6.parquet` 与
|
|
|
-`m5/fusion_38.csv` 属六层链 `model_run`/`fusion` 两步, 那两步脚本没随包。
|
|
|
-
|
|
|
-用户令 2026-09-18「`/detail` 各页面不能用旧版产出补, 必须基于输入数据重算」把这条基准**换掉了**:
|
|
|
-交付包里根本没有旧产物可比, 缺的那个"完整快照"也不该再出现在页面上。于是:
|
|
|
-
|
|
|
- · 报告里少的只有**融合级**那一列(如实写 `—`, 因为 model_run/fusion 两步确实没跑),
|
|
|
- 其余判级(状态等级/证据状态/行动等级/逐台页/总览)全部由 `data/raw/<场站>/windcms` 的
|
|
|
- `*_decode.json` 现算 —— 实测 2026-09-18: 38 台判级出来(报警 3 / 优秀 35), `/api/vibcms` 立即可用;
|
|
|
- · 故本脚本**默认**带报告与知识库 (`--no-report` 可关), `rebuild_all.py` 的 ④b 同步。
|
|
|
-
|
|
|
-## 用法
|
|
|
-
|
|
|
- python scripts/vib_raw_build.py --farm rudong # 摄入(索引+谱) + 报告/知识库
|
|
|
- python scripts/vib_raw_build.py --window w0316 # 指定窗名
|
|
|
- python scripts/vib_raw_build.py --jobs 12 --meas FFT_ # 并行/测量过滤
|
|
|
- python scripts/vib_raw_build.py --no-report # 只摄入, 不动报告/知识库
|
|
|
- python scripts/vib_raw_build.py --dry-run # 只报会做什么
|
|
|
-"""
|
|
|
-from __future__ import annotations
|
|
|
-
|
|
|
-import argparse
|
|
|
-import json
|
|
|
-import os
|
|
|
-import pathlib
|
|
|
-import shutil
|
|
|
-import subprocess
|
|
|
-import sys
|
|
|
-import time
|
|
|
-
|
|
|
-ROOT = pathlib.Path(__file__).resolve().parents[1]
|
|
|
-sys.path.insert(0, str(ROOT))
|
|
|
-from src.console import soft # noqa: E402
|
|
|
-
|
|
|
-soft()
|
|
|
-PY = sys.executable
|
|
|
-
|
|
|
-
|
|
|
-def human(n):
|
|
|
- for u in ('B', 'KB', 'MB', 'GB'):
|
|
|
- if n < 1024 or u == 'GB':
|
|
|
- return f'{n:.1f} {u}'
|
|
|
- n /= 1024
|
|
|
-
|
|
|
-
|
|
|
-def find_roots(station: pathlib.Path):
|
|
|
- """→ 振动侧原始件根列表 (含 *_decode.json 的目录或导出 zip)。"""
|
|
|
- roots = []
|
|
|
- for name in ('windcms', 'm5_cms_tcm'):
|
|
|
- d = station / name
|
|
|
- if not d.is_dir():
|
|
|
- continue
|
|
|
- for z in sorted(d.rglob('*.zip')):
|
|
|
- if z.name.upper().startswith('CMS') or 'decode' in z.name.lower():
|
|
|
- roots.append(z)
|
|
|
- for sub in sorted({p.parent for p in d.rglob('*_decode.json')}):
|
|
|
- # 取**最上层**那个含 decode 文件的目录 (measurement/), 避免每台机组一个 root
|
|
|
- top = sub
|
|
|
- while top.parent != d and any(top.parent.rglob('*_decode.json')):
|
|
|
- top = top.parent
|
|
|
- if top not in roots:
|
|
|
- roots.append(top)
|
|
|
- return roots
|
|
|
-
|
|
|
-
|
|
|
-def run(step, cmd, log):
|
|
|
- t0 = time.time()
|
|
|
- print(f'[{step}] {" ".join(str(c) for c in cmd)}', flush=True)
|
|
|
- p = subprocess.run([str(c) for c in cmd], cwd=str(ROOT), capture_output=True, text=True,
|
|
|
- encoding='utf-8', errors='replace')
|
|
|
- tail = (p.stdout or '') + (p.stderr or '')
|
|
|
- rec = dict(step=step, cmd=' '.join(str(c) for c in cmd), rc=p.returncode,
|
|
|
- seconds=round(time.time() - t0, 1), tail=tail[-1200:])
|
|
|
- log.append(rec)
|
|
|
- print(f' rc={p.returncode} {rec["seconds"]}s', flush=True)
|
|
|
- if p.returncode != 0:
|
|
|
- print(tail[-2000:], flush=True)
|
|
|
- raise SystemExit(f'步骤 {step} 失败 (rc={p.returncode})')
|
|
|
- return rec
|
|
|
-
|
|
|
-
|
|
|
-def _rels_of_window(win: pathlib.Path, m5: pathlib.Path) -> dict:
|
|
|
- """某窗的产物 → {相对产物仓的路径: 构建器说明}; 供正常摄入与 --register-only 共用 (单一实现)。"""
|
|
|
- rels = {f'm5_cms_tcm/windows/{win.name}/index.parquet':
|
|
|
- 'scripts/rudong_tcm_index.py (54 列, 与包内 tcm_index.parquet 同列名列序)',
|
|
|
- f'm5_cms_tcm/windows/{win.name}/spectra_meta.parquet':
|
|
|
- 'scripts/rudong_tcm_spectra.py',
|
|
|
- 'm5_cms_tcm/vib_raw_manifest.json': 'scripts/vib_raw_build.py'}
|
|
|
- sp = win / 'spectra'
|
|
|
- if sp.is_dir():
|
|
|
- for p in sp.rglob('*.npz'):
|
|
|
- rels[f'm5_cms_tcm/windows/{win.name}/spectra/{p.relative_to(sp).as_posix()}'] = \
|
|
|
- 'scripts/rudong_tcm_spectra.py (npz 分片: DataSets.DataSet.Values 空格串 → float 数组)'
|
|
|
- # 谱库目录里那份 meta 是同一内容的两条读取路径之一, 也登记
|
|
|
- rels[f'm5_cms_tcm/windows/{win.name}/spectra/spectra_meta.parquet'] = \
|
|
|
- 'scripts/rudong_tcm_spectra.py (与窗根同名件同内容: 兼顾两种读取约定)'
|
|
|
- return rels
|
|
|
-
|
|
|
-
|
|
|
-def main() -> int:
|
|
|
- ap = argparse.ArgumentParser()
|
|
|
- ap.add_argument('--farm', default='rudong')
|
|
|
- ap.add_argument('--window', default=None, help='窗名 (默认按数据起始日推 wMMDD)')
|
|
|
- ap.add_argument('--jobs', type=int, default=min(8, (os.cpu_count() or 4)))
|
|
|
- ap.add_argument('--meas', default='FFT_', help="谱摄入的测量过滤 (默认 FFT_; ALL=全转)")
|
|
|
- ap.add_argument('--skip-spectra', action='store_true', help='只做索引 (谱很占盘)')
|
|
|
- ap.add_argument('--no-report', dest='with_report', action='store_false',
|
|
|
- help='只摄入(索引+谱), 不跑 windcms.py report/kb')
|
|
|
- ap.set_defaults(with_report=True) # ★2026-09-18 口径翻转: 见文件头"报告默认跑了"
|
|
|
- ap.add_argument('--manifest-only', action='store_true',
|
|
|
- help='只补写 vib_raw_manifest.json(依已存在的窗; 不重跑摄入)')
|
|
|
- ap.add_argument('--register-only', action='store_true',
|
|
|
- help='不摄入, 只把**已存在**的窗体件补进来源自登记 (幂等修复: 例如先前的摄入跑在自登记'
|
|
|
- '功能之前, 台账就漏记了那些件)')
|
|
|
- ap.add_argument('--limit', type=int, default=0, help='冒烟: 只摄入前 N 个文件')
|
|
|
- ap.add_argument('--dry-run', action='store_true')
|
|
|
- a = ap.parse_args()
|
|
|
-
|
|
|
- from src.windscada.config import farm, raw_station_dir
|
|
|
- from src import paths as P
|
|
|
- cfg = farm(a.farm)
|
|
|
- station = pathlib.Path(raw_station_dir(a.farm))
|
|
|
- m5 = P.m5(a.farm)
|
|
|
- roots = find_roots(station)
|
|
|
- print(f'场站原始件目录: {station}')
|
|
|
- print(f'振动侧产物目录: {m5}')
|
|
|
-
|
|
|
- if a.register_only:
|
|
|
- # 只补登记: 找已存在的窗 (--window 指定, 否则取名字最大的那个), 逐件登记。
|
|
|
- # 放在"源件存在性检查"之前 —— 补登记不需要重新读原始件。
|
|
|
- cands = sorted(p for p in (m5 / 'windows').glob('w[0-9][0-9][0-9][0-9]')
|
|
|
- if (p / 'index.parquet').exists())
|
|
|
- if a.window:
|
|
|
- cands = [p for p in cands if p.name == a.window]
|
|
|
- if not cands:
|
|
|
- print('没有已存在的窗可登记 (先正常跑一次摄入)')
|
|
|
- return 0
|
|
|
- win = cands[-1]
|
|
|
- rels = _rels_of_window(win, m5)
|
|
|
- from src.derived_manifest import record as _rec
|
|
|
- p = _rec(P.out_root(a.farm), rels, by='scripts/vib_raw_build.py')
|
|
|
- print(f'已登记 {len(rels)} 件 (窗 {win.name}) → {p}')
|
|
|
- return 0
|
|
|
-
|
|
|
- if a.manifest_only:
|
|
|
- # 2026-09-19: 摄入已跑完、清单却没写成的情形(实测: ④b 在 report 步崩掉 → 本脚本在写清单**之前**
|
|
|
- # 就退出 ⇒ `m5_cms_tcm/vib_raw_manifest.json` 缺失, 页面归口审计按"溯源缺失"报 rc=6、整条链报失败)。
|
|
|
- # 本模式**只补写清单**, 不重跑摄入(重跑会再花 ~22 分钟并另立一个 _reimport 窗)。
|
|
|
- cands = sorted(p for p in (m5 / 'windows').glob('w[0-9][0-9][0-9][0-9]')
|
|
|
- if (p / 'index.parquet').exists())
|
|
|
- if a.window:
|
|
|
- cands = [p for p in cands if p.name == a.window]
|
|
|
- if not cands:
|
|
|
- print('没有已存在的窗可写清单 (先正常跑一次摄入)')
|
|
|
- return 0
|
|
|
- win = cands[-1]
|
|
|
- import pandas as _pd
|
|
|
- ix = _pd.read_parquet(win / 'index.parquet',
|
|
|
- columns=[c for c in ('turbine', 'sensor_name', 'meas_name', 'trigger_time')
|
|
|
- if c in _pd.read_parquet(win / 'index.parquet').columns])
|
|
|
- sm = win / 'spectra_meta.parquet'
|
|
|
- n_spec = int(len(_pd.read_parquet(sm, columns=['shard']))) if sm.is_file() else 0
|
|
|
- ts = _pd.to_datetime(ix['trigger_time'], errors='coerce') if 'trigger_time' in ix else None
|
|
|
- man = dict(farm=a.farm, window=win.name, window_dir=str(win.relative_to(ROOT)),
|
|
|
- sources=[str(r.relative_to(ROOT)) if str(r).startswith(str(ROOT)) else str(r) for r in roots],
|
|
|
- sources_note='路径相对<安装目录>; ★本清单由 --manifest-only 依**已存在的窗**补写 '
|
|
|
- '(2026-09-19), 不是摄入当次的原始记录 —— steps 留空, 时间取窗首末',
|
|
|
- rows=int(len(ix)), spectra=n_spec,
|
|
|
- time_min=(str(ts.min()) if ts is not None and len(ts) else ''), time_max=(str(ts.max()) if ts is not None and len(ts) else ''),
|
|
|
- turbines=int(ix.turbine.nunique()) if 'turbine' in ix else 0,
|
|
|
- sensors=int(ix.sensor_name.nunique()) if 'sensor_name' in ix else 0,
|
|
|
- meas_names=int(ix.meas_name.nunique()) if 'meas_name' in ix else 0,
|
|
|
- steps=[], finished=time.strftime('%Y-%m-%d %H:%M:%S'), seconds=None,
|
|
|
- missing_chain=[], missing_note='')
|
|
|
- (m5 / 'vib_raw_manifest.json').write_text(json.dumps(man, ensure_ascii=False, indent=1), encoding='utf-8')
|
|
|
- print(f'已补写 {P.rel(m5 / "vib_raw_manifest.json")}: 窗 {win.name} · {len(ix)} 行 · {n_spec} 谱')
|
|
|
- try:
|
|
|
- from src.derived_manifest import record as _rec
|
|
|
- rels = _rels_of_window(win, m5)
|
|
|
- _rec(P.out_root(a.farm), rels, by='scripts/vib_raw_build.py (manifest-only)')
|
|
|
- print(f' 已自登记 {len(rels)} 件')
|
|
|
- except Exception as _e:
|
|
|
- print(f' [i] 自登记跳过: {type(_e).__name__}: {_e}')
|
|
|
- return 0
|
|
|
-
|
|
|
- if not roots:
|
|
|
- print('没找到振动原始件 (*_decode.json 目录或 CMS 导出 zip)。'
|
|
|
- '先用 scripts/place_raw_data.py --scope vib 从现场包落位。')
|
|
|
- return 0
|
|
|
- for r in roots:
|
|
|
- n = len(list(r.rglob('*_decode.json'))) if r.is_dir() else '(zip)'
|
|
|
- sz = sum(p.stat().st_size for p in r.rglob('*_decode.json')) if r.is_dir() else r.stat().st_size
|
|
|
- print(f' 源件根: {r} {n} 个 decode 文件, {human(sz)}')
|
|
|
- if a.dry_run:
|
|
|
- print('(dry-run, 未摄入)')
|
|
|
- return 0
|
|
|
-
|
|
|
- windows_dir = m5 / 'windows'
|
|
|
- staging = windows_dir / '_staging_ingest' # 'test'/'_reimport' 之外的名字, 会被 data.windows() 看见 → 立即改名
|
|
|
- if staging.exists():
|
|
|
- shutil.rmtree(staging)
|
|
|
- staging.mkdir(parents=True, exist_ok=True)
|
|
|
-
|
|
|
- log = []
|
|
|
- t0 = time.time()
|
|
|
- root = roots[0] if len(roots) == 1 else station / 'windcms' # 多个根时交给各自目录 (index 支持 rglob)
|
|
|
- idx = staging / 'index.parquet'
|
|
|
- run('index', [PY, str(ROOT / 'scripts/rudong_tcm_index.py'), '--root', str(root), '--out', str(idx),
|
|
|
- '--jobs', str(a.jobs)] + (['--limit', str(a.limit)] if a.limit else []), log)
|
|
|
-
|
|
|
- n_spectra = 0
|
|
|
- if not a.skip_spectra:
|
|
|
- run('spectra', [PY, str(ROOT / 'scripts/rudong_tcm_spectra.py'), '--root', str(root),
|
|
|
- '--out', str(staging / 'spectra'), '--meas', a.meas, '--jobs', str(a.jobs)]
|
|
|
- + (['--limit', str(a.limit)] if a.limit else []), log)
|
|
|
- sm = staging / 'spectra' / 'spectra_meta.parquet'
|
|
|
- if sm.exists():
|
|
|
- import pandas as pd
|
|
|
- n_spectra = len(pd.read_parquet(sm))
|
|
|
-
|
|
|
- # 窗名: 用数据自身的起始日 (wMMDD), 与既有 w0127/w0707/w0811 同口径
|
|
|
- import pandas as pd
|
|
|
- d = pd.read_parquet(idx)
|
|
|
- tt = pd.to_datetime(d.trigger_time, errors='coerce')
|
|
|
- tmin, tmax = tt.min(), tt.max()
|
|
|
- win = a.window or f'w{tmin:%m%d}'
|
|
|
- final = windows_dir / win
|
|
|
- if final.exists():
|
|
|
- final = windows_dir / f'{win}_reimport_{time.strftime("%m%d%H%M")}'
|
|
|
- print(f'⚠ 窗 {win} 已存在 → 本次摄入落到 {final.name} (data.EXCLUDE_DEFAULT 会排除 _reimport 窗, 不污染生产集)')
|
|
|
- staging.rename(final)
|
|
|
- print(f'\n窗: {final.name} 数据窗 {tmin} → {tmax} 行 {len(d)} 谱 {n_spectra}')
|
|
|
-
|
|
|
- if a.with_report:
|
|
|
- # ★2026-09-19 补链序: 六层链的业务顺序是 oem_scan → energy_share → **model_run → fusion** → report,
|
|
|
- # 报告要读 model_run_l6.parquet 与 fusion_38.csv 才有"融合级"那一列。原先这里直接跑 report,
|
|
|
- # 于是清过产物的机器上报告永远缺融合级(整列 '—')。现在按序补齐(两步都是**按口径重建**, 无标准答案
|
|
|
- # 对拍 —— 口径与来路写在各自产物与台账里, 见 docs/振动六层链_接口规格与缺口_v0.1.md §3)。
|
|
|
- run('model_run', [PY, str(ROOT / 'scripts/rudong_model_run.py'), '--window', final.name], log)
|
|
|
- # 六层链的 energy_share 步(2026-09-19 补上: 此前该脚本未随包 ⇒ 链上少一环):
|
|
|
- # 线能量占比 = ±2bin 带内能量 ÷ 全谱能量; 只报占比不定级。
|
|
|
- run('energy_share', [PY, str(ROOT / 'scripts/rudong_line_energy_share.py'), '--window', final.name], log)
|
|
|
- run('fusion', [PY, str(ROOT / 'scripts/rudong_fusion_run.py'), '--window', final.name], log)
|
|
|
- # CMS 掩码阈值件(观澜自算: 中位+3/6×MAD 自适应基线; 现场导出件里没有厂商掩码定义)
|
|
|
- run('mask_thresholds', [PY, str(ROOT / 'scripts/tcm_mask_thresholds_build.py'), '--window', final.name], log)
|
|
|
- run('report', [PY, str(ROOT / 'scripts/windcms.py'), 'report', '--farm', a.farm], log)
|
|
|
- # ★2026-09-19 用户令「所有的计算均要形成观澜的源代码,确保别的电脑装完系统运行正常」:
|
|
|
- # 融合面 handoff 此前是**振动线人工交证件**、全库 0 处写入方 ⇒ 新机器上 `/detail/v2` 的
|
|
|
- # 「需要关注/全场状态」永远空着。现在由观澜自算生成同结构件(脚本在包内; 盘上若有正本则不覆盖)。
|
|
|
- run('fusion_handoff', [PY, str(ROOT / 'scripts/rudong_fusion_handoff.py')], log)
|
|
|
- run('kb', [PY, str(ROOT / 'scripts/windcms.py'), 'kb', '--farm', a.farm], log)
|
|
|
- # 六层链 fusion 步里**唯一已被标准答案验证过**的那一件 (fleet_scalar_z): 窗索引 → 同工况族中位
|
|
|
- # → 稳健 z。上面那次 fusion 已经一并出它 (同一条命令: --window final.name)。
|
|
|
- else:
|
|
|
- print('(按 --no-report 跳过 CMS 报告/知识库重生成 —— 页面上"振动评估"一栏会显示为无产物)')
|
|
|
-
|
|
|
- man = dict(farm=a.farm, window=final.name, window_dir=str(final.relative_to(ROOT)),
|
|
|
- # 源件路径写**相对安装根**的形态: 这文件会随包分发, 绝对路径换机后就是死链
|
|
|
- # (check_transferable.py 会把 outputs 下的绝对路径算作"机器相关路径")
|
|
|
- sources=[str(r.relative_to(ROOT)) if str(r).startswith(str(ROOT)) else str(r) for r in roots],
|
|
|
- sources_note='路径相对<安装目录>', rows=int(len(d)), spectra=int(n_spectra),
|
|
|
- time_min=str(tmin), time_max=str(tmax),
|
|
|
- turbines=int(d.turbine.nunique()), sensors=int(d.sensor_name.nunique()),
|
|
|
- meas_names=int(d.meas_name.nunique()),
|
|
|
- steps=[{k: v for k, v in r.items() if k != 'tail'} for r in log],
|
|
|
- finished=time.strftime('%Y-%m-%d %H:%M:%S'), seconds=round(time.time() - t0, 1),
|
|
|
- missing_chain=['rudong_tcm_oem_scan', 'rudong_line_energy_share', 'rudong_model_run',
|
|
|
- 'rudong_fusion_run'],
|
|
|
- missing_note='六层链的四步脚本未随包 (振动线分支 claude/vibration-data-diagnosis-32b69e); '
|
|
|
- '本脚本只做索引/谱/报告/知识库, 不冒充跑过那四步')
|
|
|
- (m5 / 'vib_raw_manifest.json').write_text(json.dumps(man, ensure_ascii=False, indent=1), encoding='utf-8')
|
|
|
-
|
|
|
- # 来源自登记: 让 _provenance.json 把这些件记成 raw-derived (见 src/derived_manifest.py 的说明)
|
|
|
- try:
|
|
|
- from src.derived_manifest import record as _rec
|
|
|
- rels = _rels_of_window(final, m5)
|
|
|
- if a.with_report:
|
|
|
- for p in (P.cms(a.farm)).glob('报告_CMS振动状态评估报告_*.md'):
|
|
|
- # 只登记"今天生成的"这一件 (随包/厂家转录的那些不冒充自算)
|
|
|
- if time.strftime('%Y-%m-%d') in p.name:
|
|
|
- rels[f'windcms/{p.name}'] = 'scripts/windcms.py report (数据窗: ' + final.name + ')'
|
|
|
- _rec(P.out_root(a.farm), rels, by='scripts/vib_raw_build.py')
|
|
|
- print(f' 自登记 {len(rels)} 件产物 → {P.out_root(a.farm) / "_derived_manifest.json"}')
|
|
|
- except Exception as _e:
|
|
|
- print(f' ⚠ 来源自登记失败 ({type(_e).__name__}: {_e}) —— 台账会漏记这些件, 但不影响产物本身')
|
|
|
-
|
|
|
- print(f'\n完成: 用时 {time.time() - t0:.0f}s; 清单 → {m5 / "vib_raw_manifest.json"}')
|
|
|
- print(f' 新窗 {final.name} 已进分析集: windcms 的谱图/标量/报告页会自动收录它 '
|
|
|
- f'(消费者读取时现扫 m5\\windows\\w????\\index.parquet, 无需改代码)')
|
|
|
- print(' windscada 融合面读 handoff (m5_cms_tcm/handoff_vibration_v2.json) —— 那是振动线出件, 本次不动它')
|
|
|
- print(' 看效果: python scripts/windcms.py serve --port 8020 → http://127.0.0.1:8020/')
|
|
|
- return 0
|
|
|
-
|
|
|
-
|
|
|
-if __name__ == '__main__':
|
|
|
- sys.exit(main())
|
|
|
+#!/usr/bin/env python3
|
|
|
+# -*- coding: utf-8 -*-
|
|
|
+"""振动侧一键摄入: data/raw/<场站>/{windcms,m5_cms_tcm} 的原始件 → 窗索引/谱库 → (可选)CMS 报告与知识库.
|
|
|
+
|
|
|
+## 为什么需要它 (2026-09-12 用户令)
|
|
|
+
|
|
|
+用户令: 「数据层里的「CMS 振动评估报告」应遵循 `<安装目录>\\data\\raw\\如东\\windcms`、
|
|
|
+「振动线 handoff」应遵循 `<安装目录>\\data\\raw\\如东\\m5_cms_tcm` 存放; 修改系统支持振动数据参与
|
|
|
+运行、重算」。此前振动侧只有**产物**(`outputs/<场>/m5_cms_tcm` 90 件 + `outputs/<场>/windcms` 53 件),
|
|
|
+源件不在 data/raw, 生成端脚本 (振动线分支的 rudong_tcm_*.py 七件) 也没随包 ⇒ 从零重算时振动侧是断的。
|
|
|
+
|
|
|
+两个事实把这件事限定得很清楚:
|
|
|
+ · `src/windcms/pipeline.py` 认的输入就是"含 `*_decode.json` 的目录"(Brande TCM Enterprise 导出),
|
|
|
+ 这一步**必须能跑**, 否则数据进了 data/raw 也只是躺着;
|
|
|
+ · 同分支的六层链 (oem_scan / energy_share / model_run / fusion) 脚本仍未随包 —— 本脚本**不假装**
|
|
|
+ 能把那几步跑出来: 它只做"索引 + 谱 + (可选)报告/知识库", 其余在报告里如实写"缺失"。
|
|
|
+
|
|
|
+## 这一步跑完, 振动数据就真的参与了运行
|
|
|
+
|
|
|
+ data/raw/<场站>/windcms/.../*_decode.json
|
|
|
+ → outputs/<场>/m5_cms_tcm/windows/<窗>/index.parquet (54 列, 与包内 tcm_index.parquet 同构)
|
|
|
+ → outputs/<场>/m5_cms_tcm/windows/<窗>/spectra/*.npz + spectra_meta.parquet
|
|
|
+ 消费者 (无需改一行代码, 窗是自动发现的):
|
|
|
+ src/windcms/data.py::windows() → 新窗进分析集
|
|
|
+ src/windcms/data.py::load_scalars() → 标量进 CMS 报告/页面
|
|
|
+ src/windcms/data.py::spectrum() → 谱图能取到 (此前 spectra 库整个缺失, 谱图是死的)
|
|
|
+ scripts/windcms.py report / kb → 报告与知识库重生成
|
|
|
+
|
|
|
+## ★ 报告默认**跑**了 (2026-09-18 用户令后口径翻转)
|
|
|
+
|
|
|
+2026-09-12 的旧口径是"默认不跑": 那时拿**随包快照**当基准, 重生成会让 `windcms/report.md`
|
|
|
+从 16,857 B 掉到 467 B (−97%), 因为 `data.load_model()` 要读的 `m5/model_run_l6.parquet` 与
|
|
|
+`m5/fusion_38.csv` 属六层链 `model_run`/`fusion` 两步, 那两步脚本没随包。
|
|
|
+
|
|
|
+用户令 2026-09-18「`/detail` 各页面不能用旧版产出补, 必须基于输入数据重算」把这条基准**换掉了**:
|
|
|
+交付包里根本没有旧产物可比, 缺的那个"完整快照"也不该再出现在页面上。于是:
|
|
|
+
|
|
|
+ · 报告里少的只有**融合级**那一列(如实写 `—`, 因为 model_run/fusion 两步确实没跑),
|
|
|
+ 其余判级(状态等级/证据状态/行动等级/逐台页/总览)全部由 `data/raw/<场站>/windcms` 的
|
|
|
+ `*_decode.json` 现算 —— 实测 2026-09-18: 38 台判级出来(报警 3 / 优秀 35), `/api/vibcms` 立即可用;
|
|
|
+ · 故本脚本**默认**带报告与知识库 (`--no-report` 可关), `rebuild_all.py` 的 ④b 同步。
|
|
|
+
|
|
|
+## 用法
|
|
|
+
|
|
|
+ python scripts/vib_raw_build.py --farm rudong # 摄入(索引+谱) + 报告/知识库
|
|
|
+ python scripts/vib_raw_build.py --window w0316 # 指定窗名
|
|
|
+ python scripts/vib_raw_build.py --jobs 12 --meas FFT_ # 并行/测量过滤
|
|
|
+ python scripts/vib_raw_build.py --no-report # 只摄入, 不动报告/知识库
|
|
|
+ python scripts/vib_raw_build.py --dry-run # 只报会做什么
|
|
|
+ python scripts/vib_raw_build.py --publish-only # ★恢复: 索引/谱已算完但发布被占用过 ⇒ 直接发布
|
|
|
+
|
|
|
+## ★ 发布(改名)那一步为什么要特殊处理 (2026-09-19 远端实逮)
|
|
|
+
|
|
|
+服务器上索引(719 s)+谱(894 s) 落进 `windows\\_staging_ingest`(1710 件 / 3.3 GB)后,
|
|
|
+`staging.rename(final)` 回 `PermissionError: [WinError 5] 拒绝访问` ⇒ ④b 失败 ⇒ 整条重算链停在第 4 步,
|
|
|
+页面表现为「需要关注 0 台 + 融合面缺件」。Windows 的目录 rename 要求树里没有打开的文件句柄,
|
|
|
+而杀软/索引器正在扫刚写的这批小文件 ⇒ 故本脚本的发布改为 `_move_tree()`:
|
|
|
+整目录 rename 退避重试 → 逐件搬(每件重试, rename 不行就 copy+删) → 仍锁住的如实报出来。
|
|
|
+已经算完的暂存窗可用 `--publish-only` 直接发布,不必重跑那 25 分钟的索引/谱。
|
|
|
+"""
|
|
|
+from __future__ import annotations
|
|
|
+
|
|
|
+import argparse
|
|
|
+import json
|
|
|
+import os
|
|
|
+import pathlib
|
|
|
+import shutil
|
|
|
+import subprocess
|
|
|
+import sys
|
|
|
+import time
|
|
|
+
|
|
|
+ROOT = pathlib.Path(__file__).resolve().parents[1]
|
|
|
+sys.path.insert(0, str(ROOT))
|
|
|
+
|
|
|
+# ★GBK 控制台/日志重定向不再因 '²' '⇒' 这类字符抛 UnicodeEncodeError (2026-09-19 远端实逮:
|
|
|
+# `baseline_38_build.py` 打印 'm/s²' 时在 cp936 下 rc=1, 整条恢复链少一件产物)。
|
|
|
+# 只把**不可编码字符替换掉**, 不改流编码 —— 中文照旧可读。
|
|
|
+try:
|
|
|
+ sys.stdout.reconfigure(errors='replace')
|
|
|
+except Exception:
|
|
|
+ pass
|
|
|
+from src.console import soft # noqa: E402
|
|
|
+
|
|
|
+soft()
|
|
|
+PY = sys.executable
|
|
|
+
|
|
|
+
|
|
|
+def human(n):
|
|
|
+ for u in ('B', 'KB', 'MB', 'GB'):
|
|
|
+ if n < 1024 or u == 'GB':
|
|
|
+ return f'{n:.1f} {u}'
|
|
|
+ n /= 1024
|
|
|
+
|
|
|
+
|
|
|
+def find_roots(station: pathlib.Path):
|
|
|
+ """→ 振动侧原始件根列表 (含 *_decode.json 的目录或导出 zip)。"""
|
|
|
+ roots = []
|
|
|
+ for name in ('windcms', 'm5_cms_tcm'):
|
|
|
+ d = station / name
|
|
|
+ if not d.is_dir():
|
|
|
+ continue
|
|
|
+ for z in sorted(d.rglob('*.zip')):
|
|
|
+ if z.name.upper().startswith('CMS') or 'decode' in z.name.lower():
|
|
|
+ roots.append(z)
|
|
|
+ for sub in sorted({p.parent for p in d.rglob('*_decode.json')}):
|
|
|
+ # 取**最上层**那个含 decode 文件的目录 (measurement/), 避免每台机组一个 root
|
|
|
+ top = sub
|
|
|
+ while top.parent != d and any(top.parent.rglob('*_decode.json')):
|
|
|
+ top = top.parent
|
|
|
+ if top not in roots:
|
|
|
+ roots.append(top)
|
|
|
+ return roots
|
|
|
+
|
|
|
+
|
|
|
+def run(step, cmd, log):
|
|
|
+ t0 = time.time()
|
|
|
+ print(f'[{step}] {" ".join(str(c) for c in cmd)}', flush=True)
|
|
|
+ p = subprocess.run([str(c) for c in cmd], cwd=str(ROOT), capture_output=True, text=True,
|
|
|
+ encoding='utf-8', errors='replace')
|
|
|
+ tail = (p.stdout or '') + (p.stderr or '')
|
|
|
+ rec = dict(step=step, cmd=' '.join(str(c) for c in cmd), rc=p.returncode,
|
|
|
+ seconds=round(time.time() - t0, 1), tail=tail[-1200:])
|
|
|
+ log.append(rec)
|
|
|
+ print(f' rc={p.returncode} {rec["seconds"]}s', flush=True)
|
|
|
+ if p.returncode != 0:
|
|
|
+ print(tail[-2000:], flush=True)
|
|
|
+ raise SystemExit(f'步骤 {step} 失败 (rc={p.returncode})')
|
|
|
+ return rec
|
|
|
+
|
|
|
+
|
|
|
+def _move_tree(src, dst, log=print):
|
|
|
+ """把整个目录搬到 `dst`: 先整目录 rename(退避重试), 再逐件搬, 每件都重试。
|
|
|
+
|
|
|
+ ★**2026-09-19 远端实逮, 整条重算链就断在这里**(服务器 `D:\\产品\\app`):
|
|
|
+ 索引(719 s)+谱(894 s) 刚落进 `windows\\_staging_ingest`(1710 件 / 3.3 GB),
|
|
|
+ 紧接的 `staging.rename(final)` 回 `PermissionError: [WinError 5] 拒绝访问` ⇒ ④b 失败 ⇒
|
|
|
+ 链在第 4 步停下 ⇒ 融合面 handoff / 变桨面 / 本体 / 契约全没产出,页面上就是
|
|
|
+ 「需要关注 0 台 + 融合面缺件」。
|
|
|
+ 成因不是权限:Windows 的**目录** rename 要求树里没有任何被打开的句柄,而刚写完这么多小文件时
|
|
|
+ 杀软/桌面搜索/备份正在逐个扫它们(持句柄)⇒ 目录不许改名。等一会儿句柄就放开,所以:
|
|
|
+ ① 整目录 rename + 退避重试(2/4/8/15/20/30 s)
|
|
|
+ ② 还不行就**逐件搬**(每件自己重试;rename 不行就 copy+删 —— 读句柄通常允许共享读)
|
|
|
+ 逐件路径下,个别仍被锁住的文件会在最后如实报出来(不静默丢件)。
|
|
|
+ """
|
|
|
+ import shutil as _sh
|
|
|
+ import time as _t
|
|
|
+ waits = [2, 4, 8, 15, 20, 30]
|
|
|
+ last = None
|
|
|
+ for w in [0] + waits:
|
|
|
+ if w:
|
|
|
+ _t.sleep(w)
|
|
|
+ try:
|
|
|
+ src.rename(dst)
|
|
|
+ return 'rename'
|
|
|
+ except OSError as e: # PermissionError / WinError 5 / 1224 …
|
|
|
+ last = e
|
|
|
+ log(f' [!] 目录改名被占({type(last).__name__}: {last})→ 转逐件搬(杀软/索引在扫刚写的文件,'
|
|
|
+ f'等一会儿就好;逐件搬对读句柄友好)')
|
|
|
+ dst.mkdir(parents=True, exist_ok=True)
|
|
|
+ stuck, moved = [], 0
|
|
|
+ for p in sorted(src.rglob('*')):
|
|
|
+ rel = p.relative_to(src)
|
|
|
+ target = dst / rel
|
|
|
+ if p.is_dir():
|
|
|
+ target.mkdir(parents=True, exist_ok=True)
|
|
|
+ continue
|
|
|
+ target.parent.mkdir(parents=True, exist_ok=True)
|
|
|
+ ok = False
|
|
|
+ for w in [0, 1, 2, 5, 10, 20]:
|
|
|
+ if w:
|
|
|
+ _t.sleep(w)
|
|
|
+ try:
|
|
|
+ p.rename(target)
|
|
|
+ ok = True
|
|
|
+ break
|
|
|
+ except OSError:
|
|
|
+ pass
|
|
|
+ try: # rename 不行就复制+删(读一般被允许)
|
|
|
+ _sh.copy2(p, target)
|
|
|
+ p.unlink()
|
|
|
+ ok = True
|
|
|
+ break
|
|
|
+ except OSError:
|
|
|
+ pass
|
|
|
+ if ok:
|
|
|
+ moved += 1
|
|
|
+ else:
|
|
|
+ stuck.append(str(rel))
|
|
|
+ for d in sorted([x for x in src.rglob('*') if x.is_dir()], reverse=True):
|
|
|
+ try:
|
|
|
+ d.rmdir()
|
|
|
+ except OSError:
|
|
|
+ pass
|
|
|
+ try:
|
|
|
+ src.rmdir()
|
|
|
+ except OSError:
|
|
|
+ pass
|
|
|
+ if stuck:
|
|
|
+ raise RuntimeError(
|
|
|
+ f'暂存窗里还有 {len(stuck)} 件被别的进程占着、搬不动(已搬 {moved} 件): {stuck[:8]}'
|
|
|
+ f'{" …" if len(stuck) > 8 else ""} —— 把占用方(杀软实时扫描/备份/资源管理器预览)停一下,'
|
|
|
+ f'或稍后重跑 `python scripts/vib_raw_build.py --publish-only`(不必重跑索引/谱)。')
|
|
|
+ log(f' 逐件搬完成: {moved} 件')
|
|
|
+ return 'per-file'
|
|
|
+
|
|
|
+
|
|
|
+def _rels_of_window(win: pathlib.Path, m5: pathlib.Path) -> dict:
|
|
|
+ """某窗的产物 → {相对产物仓的路径: 构建器说明}; 供正常摄入与 --register-only 共用 (单一实现)。"""
|
|
|
+ rels = {f'm5_cms_tcm/windows/{win.name}/index.parquet':
|
|
|
+ 'scripts/rudong_tcm_index.py (54 列, 与包内 tcm_index.parquet 同列名列序)',
|
|
|
+ f'm5_cms_tcm/windows/{win.name}/spectra_meta.parquet':
|
|
|
+ 'scripts/rudong_tcm_spectra.py',
|
|
|
+ 'm5_cms_tcm/vib_raw_manifest.json': 'scripts/vib_raw_build.py'}
|
|
|
+ sp = win / 'spectra'
|
|
|
+ if sp.is_dir():
|
|
|
+ for p in sp.rglob('*.npz'):
|
|
|
+ rels[f'm5_cms_tcm/windows/{win.name}/spectra/{p.relative_to(sp).as_posix()}'] = \
|
|
|
+ 'scripts/rudong_tcm_spectra.py (npz 分片: DataSets.DataSet.Values 空格串 → float 数组)'
|
|
|
+ # 谱库目录里那份 meta 是同一内容的两条读取路径之一, 也登记
|
|
|
+ rels[f'm5_cms_tcm/windows/{win.name}/spectra/spectra_meta.parquet'] = \
|
|
|
+ 'scripts/rudong_tcm_spectra.py (与窗根同名件同内容: 兼顾两种读取约定)'
|
|
|
+ return rels
|
|
|
+
|
|
|
+
|
|
|
+def main() -> int:
|
|
|
+ ap = argparse.ArgumentParser()
|
|
|
+ ap.add_argument('--farm', default='rudong')
|
|
|
+ ap.add_argument('--window', default=None, help='窗名 (默认按数据起始日推 wMMDD)')
|
|
|
+ ap.add_argument('--jobs', type=int, default=min(8, (os.cpu_count() or 4)))
|
|
|
+ ap.add_argument('--meas', default='FFT_', help="谱摄入的测量过滤 (默认 FFT_; ALL=全转)")
|
|
|
+ ap.add_argument('--skip-spectra', action='store_true', help='只做索引 (谱很占盘)')
|
|
|
+ ap.add_argument('--no-report', dest='with_report', action='store_false',
|
|
|
+ help='只摄入(索引+谱), 不跑 windcms.py report/kb')
|
|
|
+ ap.set_defaults(with_report=True) # ★2026-09-18 口径翻转: 见文件头"报告默认跑了"
|
|
|
+ ap.add_argument('--manifest-only', action='store_true', help='只补写 vib_raw_manifest.json(依已存在的窗; 不重跑摄入)')
|
|
|
+ ap.add_argument('--publish-only', action='store_true',
|
|
|
+ help='恢复模式: 上次摄入已算完索引/谱但发布(改名)被占用而失败 ⇒ 直接发布暂存窗并接着跑'
|
|
|
+ '报告链, 不重跑索引/谱 (见 _move_tree 的说明)')
|
|
|
+ ap.add_argument('--register-only', action='store_true',
|
|
|
+ help='不摄入, 只把**已存在**的窗体件补进来源自登记 (幂等修复: 例如先前的摄入跑在自登记'
|
|
|
+ '功能之前, 台账就漏记了那些件)')
|
|
|
+ ap.add_argument('--limit', type=int, default=0, help='冒烟: 只摄入前 N 个文件')
|
|
|
+ ap.add_argument('--dry-run', action='store_true')
|
|
|
+ a = ap.parse_args()
|
|
|
+
|
|
|
+ from src.windscada.config import farm, raw_station_dir
|
|
|
+ from src import paths as P
|
|
|
+ cfg = farm(a.farm)
|
|
|
+ station = pathlib.Path(raw_station_dir(a.farm))
|
|
|
+ m5 = P.m5(a.farm)
|
|
|
+ roots = find_roots(station)
|
|
|
+ print(f'场站原始件目录: {station}')
|
|
|
+ print(f'振动侧产物目录: {m5}')
|
|
|
+
|
|
|
+ if a.register_only:
|
|
|
+ # 只补登记: 找已存在的窗 (--window 指定, 否则取名字最大的那个), 逐件登记。
|
|
|
+ # 放在"源件存在性检查"之前 —— 补登记不需要重新读原始件。
|
|
|
+ cands = sorted(p for p in (m5 / 'windows').glob('w[0-9][0-9][0-9][0-9]')
|
|
|
+ if (p / 'index.parquet').exists())
|
|
|
+ if a.window:
|
|
|
+ cands = [p for p in cands if p.name == a.window]
|
|
|
+ if not cands:
|
|
|
+ print('没有已存在的窗可登记 (先正常跑一次摄入)')
|
|
|
+ return 0
|
|
|
+ win = cands[-1]
|
|
|
+ rels = _rels_of_window(win, m5)
|
|
|
+ from src.derived_manifest import record as _rec
|
|
|
+ p = _rec(P.out_root(a.farm), rels, by='scripts/vib_raw_build.py')
|
|
|
+ print(f'已登记 {len(rels)} 件 (窗 {win.name}) → {p}')
|
|
|
+ return 0
|
|
|
+
|
|
|
+ if a.manifest_only:
|
|
|
+ # 2026-09-19: 摄入已跑完、清单却没写成的情形(实测: ④b 在 report 步崩掉 → 本脚本在写清单**之前**
|
|
|
+ # 就退出 ⇒ `m5_cms_tcm/vib_raw_manifest.json` 缺失, 页面归口审计按"溯源缺失"报 rc=6、整条链报失败)。
|
|
|
+ # 本模式**只补写清单**, 不重跑摄入(重跑会再花 ~22 分钟并另立一个 _reimport 窗)。
|
|
|
+ cands = sorted(p for p in (m5 / 'windows').glob('w[0-9][0-9][0-9][0-9]')
|
|
|
+ if (p / 'index.parquet').exists())
|
|
|
+ if a.window:
|
|
|
+ cands = [p for p in cands if p.name == a.window]
|
|
|
+ if not cands:
|
|
|
+ print('没有已存在的窗可写清单 (先正常跑一次摄入)')
|
|
|
+ return 0
|
|
|
+ win = cands[-1]
|
|
|
+ import pandas as _pd
|
|
|
+ ix = _pd.read_parquet(win / 'index.parquet',
|
|
|
+ columns=[c for c in ('turbine', 'sensor_name', 'meas_name', 'trigger_time')
|
|
|
+ if c in _pd.read_parquet(win / 'index.parquet').columns])
|
|
|
+ sm = win / 'spectra_meta.parquet'
|
|
|
+ n_spec = int(len(_pd.read_parquet(sm, columns=['shard']))) if sm.is_file() else 0
|
|
|
+ ts = _pd.to_datetime(ix['trigger_time'], errors='coerce') if 'trigger_time' in ix else None
|
|
|
+ man = dict(farm=a.farm, window=win.name, window_dir=str(win.relative_to(ROOT)),
|
|
|
+ sources=[str(r.relative_to(ROOT)) if str(r).startswith(str(ROOT)) else str(r) for r in roots],
|
|
|
+ sources_note='路径相对<安装目录>; ★本清单由 --manifest-only 依**已存在的窗**补写 '
|
|
|
+ '(2026-09-19), 不是摄入当次的原始记录 —— steps 留空, 时间取窗首末',
|
|
|
+ rows=int(len(ix)), spectra=n_spec,
|
|
|
+ time_min=(str(ts.min()) if ts is not None and len(ts) else ''), time_max=(str(ts.max()) if ts is not None and len(ts) else ''),
|
|
|
+ turbines=int(ix.turbine.nunique()) if 'turbine' in ix else 0,
|
|
|
+ sensors=int(ix.sensor_name.nunique()) if 'sensor_name' in ix else 0,
|
|
|
+ meas_names=int(ix.meas_name.nunique()) if 'meas_name' in ix else 0,
|
|
|
+ steps=[], finished=time.strftime('%Y-%m-%d %H:%M:%S'), seconds=None,
|
|
|
+ missing_chain=[], missing_note='')
|
|
|
+ (m5 / 'vib_raw_manifest.json').write_text(json.dumps(man, ensure_ascii=False, indent=1), encoding='utf-8')
|
|
|
+ print(f'已补写 {P.rel(m5 / "vib_raw_manifest.json")}: 窗 {win.name} · {len(ix)} 行 · {n_spec} 谱')
|
|
|
+ try:
|
|
|
+ from src.derived_manifest import record as _rec
|
|
|
+ rels = _rels_of_window(win, m5)
|
|
|
+ _rec(P.out_root(a.farm), rels, by='scripts/vib_raw_build.py (manifest-only)')
|
|
|
+ print(f' 已自登记 {len(rels)} 件')
|
|
|
+ except Exception as _e:
|
|
|
+ print(f' [i] 自登记跳过: {type(_e).__name__}: {_e}')
|
|
|
+ return 0
|
|
|
+
|
|
|
+ if not roots:
|
|
|
+ print('没找到振动原始件 (*_decode.json 目录或 CMS 导出 zip)。'
|
|
|
+ '先用 scripts/place_raw_data.py --scope vib 从现场包落位。')
|
|
|
+ return 0
|
|
|
+ for r in roots:
|
|
|
+ n = len(list(r.rglob('*_decode.json'))) if r.is_dir() else '(zip)'
|
|
|
+ sz = sum(p.stat().st_size for p in r.rglob('*_decode.json')) if r.is_dir() else r.stat().st_size
|
|
|
+ print(f' 源件根: {r} {n} 个 decode 文件, {human(sz)}')
|
|
|
+ if a.dry_run:
|
|
|
+ print('(dry-run, 未摄入)')
|
|
|
+ return 0
|
|
|
+
|
|
|
+ windows_dir = m5 / 'windows'
|
|
|
+ staging = windows_dir / '_staging_ingest' # 'test'/'_reimport' 之外的名字, 会被 data.windows() 看见 → 立即改名
|
|
|
+ log = []
|
|
|
+ t0 = time.time()
|
|
|
+ idx = staging / 'index.parquet'
|
|
|
+ # ★`root` 必须在两个分支之前定好 (2026-09-19 实逮: 加 --publish-only 时把它挪进了 else 分支,
|
|
|
+ # 于是"只发布"这条路走到 spectra 那步就 UnboundLocalError —— 远端第一次恢复就撞上, rc=1)
|
|
|
+ root = roots[0] if len(roots) == 1 else station / 'windcms' # 多个根时交给各自目录 (index 支持 rglob)
|
|
|
+ if a.publish_only:
|
|
|
+ # ★恢复通道 (2026-09-19 远端实逮后的补): 上次摄入把索引+谱**已经算完**了, 只是最后的改名被
|
|
|
+ # 占住失败 ⇒ 别让运维白等 20~30 分钟重跑索引/谱, 直接把暂存窗发布出去, 然后照常接着跑报告链。
|
|
|
+ if not idx.is_file():
|
|
|
+ print(f'[X] 没有可发布的暂存窗: {idx} 不存在 —— 正常跑一次摄入 (python scripts/vib_raw_build.py)')
|
|
|
+ return 4
|
|
|
+ print(f'恢复模式 --publish-only: 直接发布已算好的暂存窗 {staging}(不重跑索引/谱)')
|
|
|
+ else:
|
|
|
+ if staging.exists():
|
|
|
+ shutil.rmtree(staging)
|
|
|
+ staging.mkdir(parents=True, exist_ok=True)
|
|
|
+ run('index', [PY, str(ROOT / 'scripts/rudong_tcm_index.py'), '--root', str(root), '--out', str(idx),
|
|
|
+ '--jobs', str(a.jobs)] + (['--limit', str(a.limit)] if a.limit else []), log)
|
|
|
+
|
|
|
+ n_spectra = 0
|
|
|
+ if not a.skip_spectra and not a.publish_only:
|
|
|
+ run('spectra', [PY, str(ROOT / 'scripts/rudong_tcm_spectra.py'), '--root', str(root),
|
|
|
+ '--out', str(staging / 'spectra'), '--meas', a.meas, '--jobs', str(a.jobs)]
|
|
|
+ + (['--limit', str(a.limit)] if a.limit else []), log)
|
|
|
+ sm = staging / 'spectra' / 'spectra_meta.parquet'
|
|
|
+ if sm.exists():
|
|
|
+ import pandas as pd
|
|
|
+ n_spectra = len(pd.read_parquet(sm))
|
|
|
+
|
|
|
+ # 窗名: 用数据自身的起始日 (wMMDD), 与既有 w0127/w0707/w0811 同口径
|
|
|
+ import pandas as pd
|
|
|
+ d = pd.read_parquet(idx)
|
|
|
+ tt = pd.to_datetime(d.trigger_time, errors='coerce')
|
|
|
+ tmin, tmax = tt.min(), tt.max()
|
|
|
+ win = a.window or f'w{tmin:%m%d}'
|
|
|
+ final = windows_dir / win
|
|
|
+ if final.exists():
|
|
|
+ final = windows_dir / f'{win}_reimport_{time.strftime("%m%d%H%M")}'
|
|
|
+ print(f'⚠ 窗 {win} 已存在 → 本次摄入落到 {final.name} (data.EXCLUDE_DEFAULT 会排除 _reimport 窗, 不污染生产集)')
|
|
|
+ print(f'\n发布窗目录: {final.name}({_move_tree(staging, final, print)} 方式)'
|
|
|
+ f' 数据窗 {tmin} → {tmax} 行 {len(d)} 谱 {n_spectra}')
|
|
|
+
|
|
|
+ if a.with_report:
|
|
|
+ # ★2026-09-19 补链序: 六层链的业务顺序是 oem_scan → energy_share → **model_run → fusion** → report,
|
|
|
+ # 报告要读 model_run_l6.parquet 与 fusion_38.csv 才有"融合级"那一列。原先这里直接跑 report,
|
|
|
+ # 于是清过产物的机器上报告永远缺融合级(整列 '—')。现在按序补齐(两步都是**按口径重建**, 无标准答案
|
|
|
+ # 对拍 —— 口径与来路写在各自产物与台账里, 见 docs/振动六层链_接口规格与缺口_v0.1.md §3)。
|
|
|
+ run('model_run', [PY, str(ROOT / 'scripts/rudong_model_run.py'), '--window', final.name], log)
|
|
|
+ # 六层链的 energy_share 步(2026-09-19 补上: 此前该脚本未随包 ⇒ 链上少一环):
|
|
|
+ # 线能量占比 = ±2bin 带内能量 ÷ 全谱能量; 只报占比不定级。
|
|
|
+ run('energy_share', [PY, str(ROOT / 'scripts/rudong_line_energy_share.py'), '--window', final.name], log)
|
|
|
+ run('fusion', [PY, str(ROOT / 'scripts/rudong_fusion_run.py'), '--window', final.name], log)
|
|
|
+ # CMS 掩码阈值件(观澜自算: 中位+3/6×MAD 自适应基线; 现场导出件里没有厂商掩码定义)
|
|
|
+ run('mask_thresholds', [PY, str(ROOT / 'scripts/tcm_mask_thresholds_build.py'), '--window', final.name], log)
|
|
|
+ run('report', [PY, str(ROOT / 'scripts/windcms.py'), 'report', '--farm', a.farm], log)
|
|
|
+ # ★2026-09-19 用户令「所有的计算均要形成观澜的源代码,确保别的电脑装完系统运行正常」:
|
|
|
+ # 融合面 handoff 此前是**振动线人工交证件**、全库 0 处写入方 ⇒ 新机器上 `/detail/v2` 的
|
|
|
+ # 「需要关注/全场状态」永远空着。现在由观澜自算生成同结构件(脚本在包内; 盘上若有正本则不覆盖)。
|
|
|
+ run('fusion_handoff', [PY, str(ROOT / 'scripts/rudong_fusion_handoff.py')], log)
|
|
|
+ run('kb', [PY, str(ROOT / 'scripts/windcms.py'), 'kb', '--farm', a.farm], log)
|
|
|
+ # 六层链 fusion 步里**唯一已被标准答案验证过**的那一件 (fleet_scalar_z): 窗索引 → 同工况族中位
|
|
|
+ # → 稳健 z。上面那次 fusion 已经一并出它 (同一条命令: --window final.name)。
|
|
|
+ else:
|
|
|
+ print('(按 --no-report 跳过 CMS 报告/知识库重生成 —— 页面上"振动评估"一栏会显示为无产物)')
|
|
|
+
|
|
|
+ man = dict(farm=a.farm, window=final.name, window_dir=str(final.relative_to(ROOT)),
|
|
|
+ # 源件路径写**相对安装根**的形态: 这文件会随包分发, 绝对路径换机后就是死链
|
|
|
+ # (check_transferable.py 会把 outputs 下的绝对路径算作"机器相关路径")
|
|
|
+ sources=[str(r.relative_to(ROOT)) if str(r).startswith(str(ROOT)) else str(r) for r in roots],
|
|
|
+ sources_note='路径相对<安装目录>', rows=int(len(d)), spectra=int(n_spectra),
|
|
|
+ time_min=str(tmin), time_max=str(tmax),
|
|
|
+ turbines=int(d.turbine.nunique()), sensors=int(d.sensor_name.nunique()),
|
|
|
+ meas_names=int(d.meas_name.nunique()),
|
|
|
+ steps=[{k: v for k, v in r.items() if k != 'tail'} for r in log],
|
|
|
+ finished=time.strftime('%Y-%m-%d %H:%M:%S'), seconds=round(time.time() - t0, 1),
|
|
|
+ missing_chain=['rudong_tcm_oem_scan', 'rudong_line_energy_share', 'rudong_model_run',
|
|
|
+ 'rudong_fusion_run'],
|
|
|
+ missing_note='六层链的四步脚本未随包 (振动线分支 claude/vibration-data-diagnosis-32b69e); '
|
|
|
+ '本脚本只做索引/谱/报告/知识库, 不冒充跑过那四步')
|
|
|
+ (m5 / 'vib_raw_manifest.json').write_text(json.dumps(man, ensure_ascii=False, indent=1), encoding='utf-8')
|
|
|
+
|
|
|
+ # 来源自登记: 让 _provenance.json 把这些件记成 raw-derived (见 src/derived_manifest.py 的说明)
|
|
|
+ try:
|
|
|
+ from src.derived_manifest import record as _rec
|
|
|
+ rels = _rels_of_window(final, m5)
|
|
|
+ if a.with_report:
|
|
|
+ for p in (P.cms(a.farm)).glob('报告_CMS振动状态评估报告_*.md'):
|
|
|
+ # 只登记"今天生成的"这一件 (随包/厂家转录的那些不冒充自算)
|
|
|
+ if time.strftime('%Y-%m-%d') in p.name:
|
|
|
+ rels[f'windcms/{p.name}'] = 'scripts/windcms.py report (数据窗: ' + final.name + ')'
|
|
|
+ _rec(P.out_root(a.farm), rels, by='scripts/vib_raw_build.py')
|
|
|
+ print(f' 自登记 {len(rels)} 件产物 → {P.out_root(a.farm) / "_derived_manifest.json"}')
|
|
|
+ except Exception as _e:
|
|
|
+ print(f' ⚠ 来源自登记失败 ({type(_e).__name__}: {_e}) —— 台账会漏记这些件, 但不影响产物本身')
|
|
|
+
|
|
|
+ print(f'\n完成: 用时 {time.time() - t0:.0f}s; 清单 → {m5 / "vib_raw_manifest.json"}')
|
|
|
+ print(f' 新窗 {final.name} 已进分析集: windcms 的谱图/标量/报告页会自动收录它 '
|
|
|
+ f'(消费者读取时现扫 m5\\windows\\w????\\index.parquet, 无需改代码)')
|
|
|
+ print(' windscada 融合面读 handoff (m5_cms_tcm/handoff_vibration_v2.json) —— 那是振动线出件, 本次不动它')
|
|
|
+ print(' 看效果: python scripts/windcms.py serve --port 8020 → http://127.0.0.1:8020/')
|
|
|
+ return 0
|
|
|
+
|
|
|
+
|
|
|
+if __name__ == '__main__':
|
|
|
+ sys.exit(main())
|