pipeline.py 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182
  1. # -*- coding: utf-8 -*-
  2. """一条命令从原始数据到报告: 导入 (rar/7z | 解码 JSON 目录 | 通用波形目录) → 六层链 (扫线→能量占比→六层→融合) → 报告 → 知识库.
  3. 确定性管线处理数据; 模型 (千问) 只负责触发/读产物/推理/叙述. 每步计时与返回码落 out/analyze_log.json, 失败响亮."""
  4. import os, sys, json, time, subprocess
  5. from pathlib import Path
  6. from src.windscada.config import ROOT
  7. from src.windcms import data
  8. # 六层链脚本需要真正的 Python 解释器 (含 pandas/scipy/pyarrow); 冻结版 (PyInstaller) 的 sys.executable 不是 python → 用 WINDCMS_SCRIPT_PY 或回退 venv
  9. PY = os.environ.get('WINDCMS_SCRIPT_PY') or sys.executable # 冻结包也用自己的解释器 (原写死 mac 路径)
  10. def _run(step, cmd, env, log, cwd=ROOT, timeout=7200):
  11. t0 = time.time()
  12. p = subprocess.run(cmd, cwd=str(cwd), env=env, capture_output=True, text=True, errors='replace', timeout=timeout)
  13. # ★ 2026-09-17: `text=True` 下 stdout/stderr 可能是 None(极少数平台/编码路径) ⇒ 直接相加会抛
  14. # "unsupported operand type(s) for +: 'NoneType' and 'str'", 于是**真正的失败原因被这句
  15. # 无意义的异常顶掉**(实测: oem_scan 步明确报了 rc=3 与缺什么口径, 经壳调用后看到的是这句)。
  16. out = (p.stdout or '') + (p.stderr or '')
  17. rec = dict(step=step, cmd=' '.join(map(str, cmd)), rc=p.returncode, seconds=round(time.time() - t0, 1),
  18. tail=out[-1500:])
  19. log.append(rec)
  20. print(f'[{step}] rc={p.returncode} {rec["seconds"]}s', flush=True)
  21. if p.returncode != 0:
  22. raise RuntimeError(f'步骤 {step} 失败 (rc={p.returncode}):\n{rec["tail"]}')
  23. return rec
  24. def detect_input(path):
  25. p = Path(path)
  26. if p.is_file() and p.suffix.lower() in ('.rar', '.7z', '.zip'):
  27. return 'tcm_archive'
  28. if p.is_dir():
  29. if any(p.rglob('*_decode.json')):
  30. return 'tcm_decoded_json'
  31. if any(p.glob('*.npz')) or any(p.glob('*.csv')):
  32. return 'waveforms'
  33. return 'unknown'
  34. def window_year_gate(m5_root, window, max_age_days=400):
  35. """B-1 年份闸 (2026-08-26 w1226 事故): wMMDD 窗名无年份, 历史包混入现役窗名会被字典序 sorted()[-1]
  36. 当末窗 → 冲击分型基线倒挂全场 INSUFFICIENT. 现役窗名 (w+4位数字) 只准装最近 ~13 个月数据;
  37. 更早的历史回填必须显式 histYYYY_ 前缀 (不匹配 w[0-9]{4} glob, 不进现役分析集). 缺时间列/缺 index 响亮 raise, 不静默放行."""
  38. import re
  39. import pandas as pd
  40. if not re.fullmatch(r'w\d{4}', str(window)):
  41. return None # hist*/测试窗等显式非现役名, 不适用
  42. ix = Path(m5_root) / 'windows' / str(window) / 'index.parquet'
  43. if not ix.exists():
  44. raise RuntimeError(f'年份闸: {ix} 不存在 — ingest 未落 index, 不可判年份 (禁静默放行)')
  45. import pyarrow.parquet as pq
  46. cols = pq.read_schema(str(ix)).names
  47. tcol = next((c for c in ('trigger_time', 'time', 'timestamp') if c in cols), None)
  48. if tcol is None:
  49. raise RuntimeError(f'年份闸: {ix} 无时间列 (找过 trigger_time/time/timestamp, 实有 {cols[:8]}) — 不可判年份 (禁静默放行)')
  50. tmax = pd.to_datetime(pd.read_parquet(ix, columns=[tcol])[tcol]).max()
  51. age = (pd.Timestamp.now() - tmax).days
  52. if age > max_age_days:
  53. raise RuntimeError(
  54. f'年份闸: 窗 {window} 数据止于 {tmax:%Y-%m-%d} (距今 {age} 天 > {max_age_days}) — wMMDD 窗名无年份, '
  55. f'历史数据会被字典序当成末窗 (w1226 事故 2026-08-26). 历史回填请改名 hist{tmax.year}_{window} '
  56. f'(mv {ix.parent} 后重跑), 现役窗名只装最近 13 个月数据')
  57. return tmax
  58. def _zip_has_decoded_json(path) -> bool:
  59. """压缩包里是否已经是**解码后**的 `*_decode.json` (2026-09-16)。
  60. 为什么要有这个判断: `detect_input()` 把任何 .zip/.rar/.7z 都归为 `tcm_archive`, 而那条路要调
  61. `scripts/rudong_tcm_ingest_raw.py` (解原始 base64+XML 导出) —— 该脚本**没随包** ⇒ 现场直接
  62. 把 CMS 导出的 zip 指过来必然崩在 `FileNotFoundError`。但其实"zip 里就是 `*_decode.json`"时
  63. 根本不需要它: 本包的 `rudong_tcm_index.py` / `rudong_tcm_spectra.py` **支持把 zip 当 --root**
  64. (直接按成员读), 走它们即可。这里只做"是不是解码后包"的判定, 不假装支持原始导出。
  65. """
  66. import zipfile
  67. try:
  68. with zipfile.ZipFile(path) as zf:
  69. return any(not i.is_dir() and i.filename.endswith('_decode.json') for i in zf.infolist())
  70. except Exception:
  71. return False
  72. def ingest(cfg, path, window, log, env):
  73. kind = detect_input(path)
  74. if kind == 'tcm_archive' and _zip_has_decoded_json(path):
  75. # 解码后的导出包 (zip 形态): 直接走索引/谱, 不必经缺失的 rudong_tcm_ingest_raw.py
  76. print(f'[ingest] {Path(path).name}: 内含 *_decode.json → 按"已解码导出"处理 (zip 当 root)', flush=True)
  77. kind = 'tcm_decoded_json'
  78. if kind == 'tcm_archive':
  79. raw = ROOT / 'scripts/rudong_tcm_ingest_raw.py'
  80. if not raw.is_file():
  81. raise RuntimeError(
  82. f'输入 {path} 是**原始** TCM 导出包 (rar/7z, 需先解 base64+XML), 而解码脚本 '
  83. f'{raw.name} 未随包 (振动线分支产物)。两条可走的路: '
  84. f'① 现场若能直接给"含 *_decode.json 的导出"(zip 或目录), 把它放 data/raw/<场>/windcms/ 下, '
  85. f'用 scripts/vib_raw_build.py 摄入 (本包支持 zip 当 root); '
  86. f'② 向振动线索取 scripts/rudong_tcm_ingest_raw.py 后再跑本步。')
  87. _run('ingest:tcm_archive', [PY, str(raw), '--rar', str(path), '--window', window], env, log)
  88. elif kind == 'tcm_decoded_json':
  89. out = cfg['m5'] / 'windows' / window
  90. out.mkdir(parents=True, exist_ok=True)
  91. _run('ingest:index', [PY, str(ROOT / 'scripts/rudong_tcm_index.py'), '--root', str(path), '--out', str(out / 'index.parquet')], env, log)
  92. _run('ingest:spectra', [PY, str(ROOT / 'scripts/rudong_tcm_spectra.py'), '--root', str(path), '--out', str(out / 'spectra')], env, log)
  93. elif kind == 'waveforms':
  94. from src.windcms import ingest_waveforms
  95. t0 = time.time()
  96. n = ingest_waveforms.ingest_dir(cfg, Path(path), window)
  97. log.append(dict(step='ingest:waveforms', cmd=f'ingest_waveforms.ingest_dir({path}, {window})', rc=0, seconds=round(time.time() - t0, 1), tail=f'{n} records'))
  98. else:
  99. raise RuntimeError(f'无法识别输入格式: {path} (支持 TCM rar/7z/zip 导出包 · 含 *_decode.json 的目录 · 含 npz/csv 波形的目录)')
  100. window_year_gate(cfg['m5'], window)
  101. return kind
  102. def analyze(cfg, input_path=None, window=None, steps=None, bands='ALL'):
  103. """全链. steps 子集可选: ingest, oem_scan, energy_share, model_run, fusion, report, kb."""
  104. out = cfg['out']
  105. out.mkdir(parents=True, exist_ok=True)
  106. log = []
  107. env = dict(os.environ)
  108. env.setdefault('MPLBACKEND', 'Agg')
  109. t_all = time.time()
  110. try:
  111. if input_path:
  112. if not window:
  113. raise RuntimeError('导入需要 --window 名 (如 w0905)')
  114. ingest(cfg, input_path, window, log, env)
  115. wins = list(data.windows(cfg))
  116. env['M5_WINDOWS'] = ','.join(wins)
  117. env['SCAN_BANDS'] = bands
  118. want = set(steps or ['oem_scan', 'energy_share', 'model_run', 'fusion', 'report', 'kb'])
  119. # ★ 2026-09-17: 缺脚本要给"人看得懂、且知道去哪儿要料"的错, 而不是让 python 抛一句
  120. # `can't open file 'scripts/rudong_model_run.py'`(实测就是这句 —— 现场看到它没法判断
  121. # 这是"包缺件"还是"我敲错了")。四步脚本属振动线分支产物, 未随包; 缺了就指名要说清楚。
  122. def _step(name: str, script: str):
  123. p = ROOT / 'scripts' / script
  124. if not p.is_file():
  125. raise RuntimeError(
  126. f'六层链的 {name} 步脚本未随包: {p.relative_to(ROOT).as_posix()}(振动线分支产物)。'
  127. f'调用约定见 docs/振动六层链_接口规格与缺口_v0.1.md §1, 取料单见 docs/向振动线取料单_v0.1.md; '
  128. f'自查: python scripts/chain_gap_check.py。'
  129. f'(本包已逆向实现 fusion 步的一件产物: scripts/rudong_fusion_run.py 可重出 fleet_scalar_z)')
  130. return [PY, str(p)]
  131. if 'oem_scan' in want:
  132. _run('oem_scan', _step('oem_scan', 'rudong_tcm_oem_scan.py'), env, log)
  133. if 'energy_share' in want:
  134. _run('energy_share', _step('energy_share', 'rudong_line_energy_share.py'), env, log)
  135. if 'model_run' in want:
  136. _run('model_run', _step('model_run', 'rudong_model_run.py'), env, log)
  137. if 'fusion' in want:
  138. _run('fusion', _step('fusion', 'rudong_fusion_run.py'), env, log)
  139. if 'report' in want:
  140. # 标量缓存按窗集合命名, 新窗自动失效; 谱元数据进程内缓存由新进程重建
  141. _run('report', [PY, str(ROOT / 'scripts/windcms.py'), 'report', '--farm', cfg.get('key', 'rudong')], env, log)
  142. if 'kb' in want:
  143. _run('kb', [PY, str(ROOT / 'scripts/windcms.py'), 'kb', '--farm', cfg.get('key', 'rudong')], env, log)
  144. status = 'ok'
  145. except Exception as e:
  146. log.append(dict(step='FAILED', cmd='', rc=1, seconds=0, tail=str(e)[-1500:]))
  147. status = 'failed'
  148. summary = dict(status=status, windows=list(data.windows(cfg)), seconds=round(time.time() - t_all, 1), finished=time.strftime('%Y-%m-%d %H:%M:%S'), steps=log)
  149. # ★ 2026-09-17 用户令 2: 运行日志只在 logs/ —— 原先写 `out/analyze_log.json`(产物仓),
  150. # 于是产物仓里躺着运行日志(反向审计把它当产物、log_audit 也要报了)。改落
  151. # `logs/build/<场>/windcms/analyze_log.json`; 产物仓里的旧件已被清理。
  152. # 日志里有中文命令回显: 缺 encoding 会按系统 locale(cp936) 写 → 中文 Windows 崩
  153. from src import logfile as _lf
  154. logp = _lf.audit_log('windcms_analyze.json') # 审计/报告类 JSON 按口径落 logs/audit/
  155. logp.parent.mkdir(parents=True, exist_ok=True)
  156. with open(logp, 'w', encoding='utf-8', newline='\n') as f:
  157. json.dump(summary, f, ensure_ascii=False, indent=1)
  158. print(f'analyze {status}: {summary["seconds"]}s → {logp}', flush=True)
  159. if status != 'ok':
  160. # ★ 2026-09-17: 失败时把**原因**打到控制台, 而不是只留一句 "analyze failed" 让人去翻 JSON。
  161. # 实测场景: 四步脚本未随包时, 原来只在日志里能看到 python 的 FileNotFoundError,
  162. # 现场看到的是"重算失败"四个字, 无从判断是包缺件还是操作有误。
  163. bad = [s for s in log if s.get('step') == 'FAILED'] + [s for s in log if s.get('rc') not in (0, None)]
  164. for s in bad[:2]:
  165. print(' [X] ' + str(s.get('step')) + ': ' + str(s.get('tail', '')).strip()[:400], flush=True)
  166. print(f' 详情: {logp}', flush=True)
  167. return summary