vib_raw_build.py 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271
  1. #!/usr/bin/env python3
  2. # -*- coding: utf-8 -*-
  3. """振动侧一键摄入: data/raw/<场站>/{windcms,m5_cms_tcm} 的原始件 → 窗索引/谱库 → (可选)CMS 报告与知识库.
  4. ## 为什么需要它 (2026-09-12 用户令)
  5. 用户令: 「数据层里的「CMS 振动评估报告」应遵循 `<安装目录>\\data\\raw\\如东\\windcms`、
  6. 「振动线 handoff」应遵循 `<安装目录>\\data\\raw\\如东\\m5_cms_tcm` 存放; 修改系统支持振动数据参与
  7. 运行、重算」。此前振动侧只有**产物**(`outputs/<场>/m5_cms_tcm` 90 件 + `outputs/<场>/windcms` 53 件),
  8. 源件不在 data/raw, 生成端脚本 (振动线分支的 rudong_tcm_*.py 七件) 也没随包 ⇒ 从零重算时振动侧是断的。
  9. 两个事实把这件事限定得很清楚:
  10. · `src/windcms/pipeline.py` 认的输入就是"含 `*_decode.json` 的目录"(Brande TCM Enterprise 导出),
  11. 这一步**必须能跑**, 否则数据进了 data/raw 也只是躺着;
  12. · 同分支的六层链 (oem_scan / energy_share / model_run / fusion) 脚本仍未随包 —— 本脚本**不假装**
  13. 能把那几步跑出来: 它只做"索引 + 谱 + (可选)报告/知识库", 其余在报告里如实写"缺失"。
  14. ## 这一步跑完, 振动数据就真的参与了运行
  15. data/raw/<场站>/windcms/.../*_decode.json
  16. → outputs/<场>/m5_cms_tcm/windows/<窗>/index.parquet (54 列, 与包内 tcm_index.parquet 同构)
  17. → outputs/<场>/m5_cms_tcm/windows/<窗>/spectra/*.npz + spectra_meta.parquet
  18. 消费者 (无需改一行代码, 窗是自动发现的):
  19. src/windcms/data.py::windows() → 新窗进分析集
  20. src/windcms/data.py::load_scalars() → 标量进 CMS 报告/页面
  21. src/windcms/data.py::spectrum() → 谱图能取到 (此前 spectra 库整个缺失, 谱图是死的)
  22. scripts/windcms.py report / kb → 报告与知识库重生成
  23. ## ★ 为什么"重生成 CMS 报告"默认**不跑** (2026-09-12 实测)
  24. 跑一次 `scripts/windcms.py report` 在本包会**让产物变差**, 不是变好:
  25. | 件 | 随包快照 | 本包重生成后 | 变化 |
  26. |---|---|---|---|
  27. | `windcms/report.md` | 16,857 B (含逐台融合级表 + L4 过闸谱线) | **467 B** (融合级表空、L4 写"无") | **−97.2%** |
  28. | `windcms/overview.html` | 646,218 B | 548,466 B | −15.1% |
  29. | `windcms/index_eng.html` | 443,841 B | 408,112 B | −8.0% |
  30. 根因不在数据, 在**缺件**: `src/windcms/data.py::load_model()` 要读
  31. `m5/model_run_l6.parquet` 与 `m5/fusion_38.csv`, 而这两件属六层链的 `model_run` / `fusion` 两步 ——
  32. **那四步脚本没随包**(见 `missing_chain`), 产物也不在。于是重生成 = 用残缺输入覆盖完整快照。
  33. 所以: 摄入(索引/谱)默认跑, **报告/知识库默认不跑**; 要跑得显式 `--with-report`, 并且跑之前**先备份**
  34. `outputs/<场>/windcms/`。六层链补齐后这个默认值应当翻过来 (那时重生成才会≥快照)。
  35. ## 用法
  36. python scripts/vib_raw_build.py --farm rudong # 摄入(索引+谱), 不动报告
  37. python scripts/vib_raw_build.py --window w0316 # 指定窗名
  38. python scripts/vib_raw_build.py --jobs 12 --meas FFT_ # 并行/测量过滤
  39. python scripts/vib_raw_build.py --with-report # 额外重生成 CMS 报告/知识库(见上, 慎用)
  40. python scripts/vib_raw_build.py --dry-run # 只报会做什么
  41. """
  42. from __future__ import annotations
  43. import argparse
  44. import json
  45. import os
  46. import pathlib
  47. import shutil
  48. import subprocess
  49. import sys
  50. import time
  51. ROOT = pathlib.Path(__file__).resolve().parents[1]
  52. sys.path.insert(0, str(ROOT))
  53. from src.console import soft # noqa: E402
  54. soft()
  55. PY = sys.executable
  56. def human(n):
  57. for u in ('B', 'KB', 'MB', 'GB'):
  58. if n < 1024 or u == 'GB':
  59. return f'{n:.1f} {u}'
  60. n /= 1024
  61. def find_roots(station: pathlib.Path):
  62. """→ 振动侧原始件根列表 (含 *_decode.json 的目录或导出 zip)。"""
  63. roots = []
  64. for name in ('windcms', 'm5_cms_tcm'):
  65. d = station / name
  66. if not d.is_dir():
  67. continue
  68. for z in sorted(d.rglob('*.zip')):
  69. if z.name.upper().startswith('CMS') or 'decode' in z.name.lower():
  70. roots.append(z)
  71. for sub in sorted({p.parent for p in d.rglob('*_decode.json')}):
  72. # 取**最上层**那个含 decode 文件的目录 (measurement/), 避免每台机组一个 root
  73. top = sub
  74. while top.parent != d and any(top.parent.rglob('*_decode.json')):
  75. top = top.parent
  76. if top not in roots:
  77. roots.append(top)
  78. return roots
  79. def run(step, cmd, log):
  80. t0 = time.time()
  81. print(f'[{step}] {" ".join(str(c) for c in cmd)}', flush=True)
  82. p = subprocess.run([str(c) for c in cmd], cwd=str(ROOT), capture_output=True, text=True,
  83. encoding='utf-8', errors='replace')
  84. tail = (p.stdout or '') + (p.stderr or '')
  85. rec = dict(step=step, cmd=' '.join(str(c) for c in cmd), rc=p.returncode,
  86. seconds=round(time.time() - t0, 1), tail=tail[-1200:])
  87. log.append(rec)
  88. print(f' rc={p.returncode} {rec["seconds"]}s', flush=True)
  89. if p.returncode != 0:
  90. print(tail[-2000:], flush=True)
  91. raise SystemExit(f'步骤 {step} 失败 (rc={p.returncode})')
  92. return rec
  93. def _rels_of_window(win: pathlib.Path, m5: pathlib.Path) -> dict:
  94. """某窗的产物 → {相对产物仓的路径: 构建器说明}; 供正常摄入与 --register-only 共用 (单一实现)。"""
  95. rels = {f'm5_cms_tcm/windows/{win.name}/index.parquet':
  96. 'scripts/rudong_tcm_index.py (54 列, 与包内 tcm_index.parquet 同列名列序)',
  97. f'm5_cms_tcm/windows/{win.name}/spectra_meta.parquet':
  98. 'scripts/rudong_tcm_spectra.py',
  99. 'm5_cms_tcm/vib_raw_manifest.json': 'scripts/vib_raw_build.py'}
  100. sp = win / 'spectra'
  101. if sp.is_dir():
  102. for p in sp.rglob('*.npz'):
  103. rels[f'm5_cms_tcm/windows/{win.name}/spectra/{p.relative_to(sp).as_posix()}'] = \
  104. 'scripts/rudong_tcm_spectra.py (npz 分片: DataSets.DataSet.Values 空格串 → float 数组)'
  105. # 谱库目录里那份 meta 是同一内容的两条读取路径之一, 也登记
  106. rels[f'm5_cms_tcm/windows/{win.name}/spectra/spectra_meta.parquet'] = \
  107. 'scripts/rudong_tcm_spectra.py (与窗根同名件同内容: 兼顾两种读取约定)'
  108. return rels
  109. def main() -> int:
  110. ap = argparse.ArgumentParser()
  111. ap.add_argument('--farm', default='rudong')
  112. ap.add_argument('--window', default=None, help='窗名 (默认按数据起始日推 wMMDD)')
  113. ap.add_argument('--jobs', type=int, default=min(8, (os.cpu_count() or 4)))
  114. ap.add_argument('--meas', default='FFT_', help="谱摄入的测量过滤 (默认 FFT_; ALL=全转)")
  115. ap.add_argument('--skip-spectra', action='store_true', help='只做索引 (谱很占盘)')
  116. ap.add_argument('--with-report', action='store_true',
  117. help='额外跑 windcms.py report/kb —— ★默认关: 本包缺六层链的 model_run/fusion 产物, '
  118. '重生成会让 report.md/overview.html 掉内容 (见文件头实测表); 跑前先备份 windcms\\')
  119. ap.add_argument('--register-only', action='store_true',
  120. help='不摄入, 只把**已存在**的窗体件补进来源自登记 (幂等修复: 例如先前的摄入跑在自登记'
  121. '功能之前, 台账就漏记了那些件)')
  122. ap.add_argument('--limit', type=int, default=0, help='冒烟: 只摄入前 N 个文件')
  123. ap.add_argument('--dry-run', action='store_true')
  124. a = ap.parse_args()
  125. from src.windscada.config import farm, raw_station_dir
  126. from src import paths as P
  127. cfg = farm(a.farm)
  128. station = pathlib.Path(raw_station_dir(a.farm))
  129. m5 = P.m5(a.farm)
  130. roots = find_roots(station)
  131. print(f'场站原始件目录: {station}')
  132. print(f'振动侧产物目录: {m5}')
  133. if a.register_only:
  134. # 只补登记: 找已存在的窗 (--window 指定, 否则取名字最大的那个), 逐件登记。
  135. # 放在"源件存在性检查"之前 —— 补登记不需要重新读原始件。
  136. cands = sorted(p for p in (m5 / 'windows').glob('w[0-9][0-9][0-9][0-9]')
  137. if (p / 'index.parquet').exists())
  138. if a.window:
  139. cands = [p for p in cands if p.name == a.window]
  140. if not cands:
  141. print('没有已存在的窗可登记 (先正常跑一次摄入)')
  142. return 0
  143. win = cands[-1]
  144. rels = _rels_of_window(win, m5)
  145. from src.derived_manifest import record as _rec
  146. p = _rec(P.out_root(a.farm), rels, by='scripts/vib_raw_build.py')
  147. print(f'已登记 {len(rels)} 件 (窗 {win.name}) → {p}')
  148. return 0
  149. if not roots:
  150. print('没找到振动原始件 (*_decode.json 目录或 CMS 导出 zip)。'
  151. '先用 scripts/place_raw_data.py --scope vib 从现场包落位。')
  152. return 0
  153. for r in roots:
  154. n = len(list(r.rglob('*_decode.json'))) if r.is_dir() else '(zip)'
  155. sz = sum(p.stat().st_size for p in r.rglob('*_decode.json')) if r.is_dir() else r.stat().st_size
  156. print(f' 源件根: {r} {n} 个 decode 文件, {human(sz)}')
  157. if a.dry_run:
  158. print('(dry-run, 未摄入)')
  159. return 0
  160. windows_dir = m5 / 'windows'
  161. staging = windows_dir / '_staging_ingest' # 'test'/'_reimport' 之外的名字, 会被 data.windows() 看见 → 立即改名
  162. if staging.exists():
  163. shutil.rmtree(staging)
  164. staging.mkdir(parents=True, exist_ok=True)
  165. log = []
  166. t0 = time.time()
  167. root = roots[0] if len(roots) == 1 else station / 'windcms' # 多个根时交给各自目录 (index 支持 rglob)
  168. idx = staging / 'index.parquet'
  169. run('index', [PY, str(ROOT / 'scripts/rudong_tcm_index.py'), '--root', str(root), '--out', str(idx),
  170. '--jobs', str(a.jobs)] + (['--limit', str(a.limit)] if a.limit else []), log)
  171. n_spectra = 0
  172. if not a.skip_spectra:
  173. run('spectra', [PY, str(ROOT / 'scripts/rudong_tcm_spectra.py'), '--root', str(root),
  174. '--out', str(staging / 'spectra'), '--meas', a.meas, '--jobs', str(a.jobs)]
  175. + (['--limit', str(a.limit)] if a.limit else []), log)
  176. sm = staging / 'spectra' / 'spectra_meta.parquet'
  177. if sm.exists():
  178. import pandas as pd
  179. n_spectra = len(pd.read_parquet(sm))
  180. # 窗名: 用数据自身的起始日 (wMMDD), 与既有 w0127/w0707/w0811 同口径
  181. import pandas as pd
  182. d = pd.read_parquet(idx)
  183. tt = pd.to_datetime(d.trigger_time, errors='coerce')
  184. tmin, tmax = tt.min(), tt.max()
  185. win = a.window or f'w{tmin:%m%d}'
  186. final = windows_dir / win
  187. if final.exists():
  188. final = windows_dir / f'{win}_reimport_{time.strftime("%m%d%H%M")}'
  189. print(f'⚠ 窗 {win} 已存在 → 本次摄入落到 {final.name} (data.EXCLUDE_DEFAULT 会排除 _reimport 窗, 不污染生产集)')
  190. staging.rename(final)
  191. print(f'\n窗: {final.name} 数据窗 {tmin} → {tmax} 行 {len(d)} 谱 {n_spectra}')
  192. if a.with_report:
  193. run('report', [PY, str(ROOT / 'scripts/windcms.py'), 'report', '--farm', a.farm], log)
  194. run('kb', [PY, str(ROOT / 'scripts/windcms.py'), 'kb', '--farm', a.farm], log)
  195. else:
  196. print('(跳过 CMS 报告/知识库重生成 —— 本包缺 model_run/fusion 产物, 重生成会掉内容; '
  197. '要跑用 --with-report, 且先备份 windcms\\)')
  198. man = dict(farm=a.farm, window=final.name, window_dir=str(final.relative_to(ROOT)),
  199. # 源件路径写**相对安装根**的形态: 这文件会随包分发, 绝对路径换机后就是死链
  200. # (check_transferable.py 会把 outputs 下的绝对路径算作"机器相关路径")
  201. sources=[str(r.relative_to(ROOT)) if str(r).startswith(str(ROOT)) else str(r) for r in roots],
  202. sources_note='路径相对<安装目录>', rows=int(len(d)), spectra=int(n_spectra),
  203. time_min=str(tmin), time_max=str(tmax),
  204. turbines=int(d.turbine.nunique()), sensors=int(d.sensor_name.nunique()),
  205. meas_names=int(d.meas_name.nunique()),
  206. steps=[{k: v for k, v in r.items() if k != 'tail'} for r in log],
  207. finished=time.strftime('%Y-%m-%d %H:%M:%S'), seconds=round(time.time() - t0, 1),
  208. missing_chain=['rudong_tcm_oem_scan', 'rudong_line_energy_share', 'rudong_model_run',
  209. 'rudong_fusion_run'],
  210. missing_note='六层链的四步脚本未随包 (振动线分支 claude/vibration-data-diagnosis-32b69e); '
  211. '本脚本只做索引/谱/报告/知识库, 不冒充跑过那四步')
  212. (m5 / 'vib_raw_manifest.json').write_text(json.dumps(man, ensure_ascii=False, indent=1), encoding='utf-8')
  213. # 来源自登记: 让 _provenance.json 把这些件记成 raw-derived (见 src/derived_manifest.py 的说明)
  214. try:
  215. from src.derived_manifest import record as _rec
  216. rels = _rels_of_window(final, m5)
  217. if a.with_report:
  218. for p in (P.cms(a.farm)).glob('报告_CMS振动状态评估报告_*.md'):
  219. # 只登记"今天生成的"这一件 (随包/厂家转录的那些不冒充自算)
  220. if time.strftime('%Y-%m-%d') in p.name:
  221. rels[f'windcms/{p.name}'] = 'scripts/windcms.py report (数据窗: ' + final.name + ')'
  222. _rec(P.out_root(a.farm), rels, by='scripts/vib_raw_build.py')
  223. print(f' 自登记 {len(rels)} 件产物 → {P.out_root(a.farm) / "_derived_manifest.json"}')
  224. except Exception as _e:
  225. print(f' ⚠ 来源自登记失败 ({type(_e).__name__}: {_e}) —— 台账会漏记这些件, 但不影响产物本身')
  226. print(f'\n完成: 用时 {time.time() - t0:.0f}s; 清单 → {m5 / "vib_raw_manifest.json"}')
  227. print(f' 新窗 {final.name} 已进分析集: windcms 的谱图/标量/报告页会自动收录它 '
  228. f'(消费者读取时现扫 m5\\windows\\w????\\index.parquet, 无需改代码)')
  229. print(' windscada 融合面读 handoff (m5_cms_tcm/handoff_vibration_v2.json) —— 那是振动线出件, 本次不动它')
  230. print(' 看效果: python scripts/windcms.py serve --port 8020 → http://127.0.0.1:8020/')
  231. return 0
  232. if __name__ == '__main__':
  233. sys.exit(main())