rudong_fusion_run.py 8.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161
  1. #!/usr/bin/env python3
  2. # -*- coding: utf-8 -*-
  3. r"""六层链 · `fusion` 步(**逆向工程实现**,2026-09-17 用户令 #2)。
  4. ## 这个文件是怎么来的
  5. 包里的编排壳 `src/windcms/pipeline.py::analyze` 会调 `scripts/rudong_fusion_run.py`,而该脚本没随包
  6. (六层链四步脚本都缺)。用户令 #2 问"能不能根据旧包产物逆向工程推导出来"——本文件是**可对拍的那部分**
  7. 的答案:`fleet_scalar_z.parquet` 的标准答案在盘上(随包件),于是口径可以**反推 + 逐值对拍**。
  8. ## 反推出来的口径(三条都逐值验过)
  9. val = 窗内 median(scalar_value),按 (turbine, sensor_name, meas_name, condition_key) 分组
  10. n = 该组记录数
  11. fleet_med = 跨机组的 median(val)(同一 sensor/meas/bin 下)
  12. z = (val - fleet_med) / (1.4826 × MAD(val)) ← 稳健 z 分数(MAD 抗离群)
  13. 验证(`--verify`,样件 outputs/<场>/m5_cms_tcm/fleet_scalar_z.parquet,8,887 行):
  14. 样例 WTG01 / Gear_HS_generator_side / CrestFactor / WPS-ActivePower 0-1600,
  15. 样件: val=5.66453 fleet_med=5.91444 z=-0.210156 n=7
  16. 本器: val=5.66453 fleet_med=5.91444 z=-0.210153 n=7
  17. 尺度候选实测: std=1.44859(z=-0.1725) · MAD=0.80209(z=-0.3116) · IQR/1.349=1.46899(z=-0.1701)
  18. **1.4826×MAD=1.18918(z=-0.2102) ← 命中**
  19. ## 用法
  20. python scripts/rudong_fusion_run.py # 算并落盘 (默认窗口 = 配置里的首窗)
  21. python scripts/rudong_fusion_run.py --verify # 与盘上样件逐值对拍 (不改任何件)
  22. python scripts/rudong_fusion_run.py --out <路径> # 写到别处 (对拍/试验用)
  23. ★ 本器**只实现已被标准答案验证过的那一件**。`fusion_38.csv` 与 `model_run_l6.parquet` 从未随过包,
  24. 没有标准答案 ⇒ 不在本器里猜着写(本项目硬规矩: 允许响亮降级, 不许造数; 见
  25. `docs/振动六层链_接口规格与缺口_v0.1.md`)。
  26. 退出码: 0 成功/对拍通过 · 5 对拍不一致 · 2 找不到输入
  27. """
  28. from __future__ import annotations
  29. import argparse
  30. import pathlib
  31. import sys
  32. ROOT = pathlib.Path(__file__).resolve().parents[1]
  33. sys.path.insert(0, str(ROOT))
  34. from src import paths as P # noqa: E402
  35. def _window_index(win: str) -> pathlib.Path:
  36. """窗索引在哪: w0127 = 首窗特例 (索引在 m5 根), 其余在 windows/<w>/index.parquet。"""
  37. m5 = P.m5()
  38. if win == 'w0127' and (m5 / 'tcm_index.parquet').is_file():
  39. return m5 / 'tcm_index.parquet'
  40. return m5 / 'windows' / win / 'index.parquet'
  41. def compute(win: str = 'w0127', min_n: int = 3):
  42. """→ 与样件同构的 DataFrame(sensor, meas, bin, turbine, val, fleet_med, z, n)。"""
  43. import numpy as np
  44. import pandas as pd
  45. ix = _window_index(win)
  46. if not ix.is_file():
  47. raise SystemExit(f'[X] 窗索引不存在: {P.rel(ix)} —— 先跑 scripts/rudong_tcm_index.py')
  48. d = pd.read_parquet(ix, columns=['turbine', 'sensor_name', 'meas_name', 'condition_key', 'scalar_value', 'y_unit'])
  49. d['val'] = pd.to_numeric(d['scalar_value'], errors='coerce')
  50. d = d.dropna(subset=['val', 'turbine', 'sensor_name', 'meas_name', 'condition_key'])
  51. # ★ 量纲闸 (反推自样件, 逐值验过): 只留振动量纲的标量 —— 1(无量纲指标) / m/s / m/s²。
  52. # 排除 %(Disk Usage / Memory Usage)、RPM(rms_rawRPM_DC) 与全部波形/谱名(FFT_/Time_/Env_/Cep_)。
  53. # 加上这一条之后, 样件 8,887 行的 val / fleet_med / z **逐值 100% 一致**。
  54. d = d[d['y_unit'].astype(str).isin({'1', 'm/s', 'm/s²'})]
  55. g = (d.groupby(['sensor_name', 'meas_name', 'condition_key', 'turbine'])['val']
  56. .agg(val='median', n='size').reset_index())
  57. g = g[g['n'] >= min_n] # ★ 样件口径: 每组至少 3 条记录 (反推自 n 分布 min=3)
  58. # 跨机组的稳健基线: 中位数 + 1.4826×MAD (与样件逐值一致)
  59. gm = g.groupby(['sensor_name', 'meas_name', 'condition_key'])['val'].median().rename('fleet_med')
  60. g = g.merge(gm, on=['sensor_name', 'meas_name', 'condition_key'], how='left')
  61. mad = (g.assign(dev=(g['val'] - g['fleet_med']).abs())
  62. .groupby(['sensor_name', 'meas_name', 'condition_key'])['dev'].median()
  63. .mul(1.4826).rename('scale'))
  64. g = g.merge(mad, on=['sensor_name', 'meas_name', 'condition_key'], how='left')
  65. g['z'] = (g['val'] - g['fleet_med']) / g['scale']
  66. out = g.rename(columns={'sensor_name': 'sensor', 'meas_name': 'meas', 'condition_key': 'bin'})
  67. return out[['sensor', 'meas', 'bin', 'turbine', 'val', 'fleet_med', 'z', 'n']].sort_values(
  68. ['sensor', 'meas', 'bin', 'turbine']).reset_index(drop=True)
  69. def verify(win: str = 'w0127', sample: pathlib.Path | None = None) -> int:
  70. """与盘上样件逐值对拍(严格: 同键同值, 容差 1e-4)。"""
  71. import pandas as pd
  72. sp = pathlib.Path(sample) if sample else (P.m5() / 'fleet_scalar_z.parquet')
  73. if not sp.is_file():
  74. print(f'[X] 没有样件可对拍: {P.rel(sp)}')
  75. return 2
  76. got = compute(win)
  77. want = pd.read_parquet(sp)
  78. key = ['sensor', 'meas', 'bin', 'turbine']
  79. m = want.merge(got, on=key, how='outer', suffixes=('_样件', '_本器'), indicator=True)
  80. both = m[m['_merge'] == 'both']
  81. only_w = int((m['_merge'] == 'left_only').sum())
  82. only_g = int((m['_merge'] == 'right_only').sum())
  83. def near(col, tol=1e-4):
  84. a, b = pd.to_numeric(both[f'{col}_样件'], errors='coerce'), pd.to_numeric(both[f'{col}_本器'], errors='coerce')
  85. ok = (a - b).abs() <= tol
  86. return int(ok.sum()), int(len(ok) - ok.sum())
  87. print(f'== fleet_scalar_z 逐值对拍 · 窗 {win} ==')
  88. print(f' 样件 {len(want)} 行 · 本器 {len(got)} 行 · 同键 {len(both)} 行'
  89. f' · 仅样件 {only_w} · 仅本器 {only_g}')
  90. bad_any = 0
  91. for col in ('val', 'fleet_med', 'z', 'n'):
  92. if f'{col}_样件' not in both.columns:
  93. print(f' {col}: (样件无此列)')
  94. continue
  95. ok, bad = near(col, tol=1e-4 if col != 'n' else 0.5)
  96. bad_any += bad
  97. print(f' {col:10s} 一致 {ok:5d} / {len(both):5d} 不一致 {bad}')
  98. ok_all = (len(both) == len(want) == len(got) and bad_any == 0)
  99. print(f' 结论: {"逐值完全一致(口径已复现)" if ok_all else "有差异 —— 见上, 不要拿本器产物替换样件"}')
  100. return 0 if ok_all else 5
  101. def main() -> int:
  102. ap = argparse.ArgumentParser(description='六层链 fusion 步(逆向工程实现: fleet_scalar_z)')
  103. ap.add_argument('--window', default='w0127', help='用哪个窗算(默认首窗 w0127,与样件同源)')
  104. ap.add_argument('--verify', action='store_true', help='与盘上样件逐值对拍,不写盘')
  105. ap.add_argument('--sample', default=None, help='--verify 用哪个样件(默认 outputs/<场>/m5_cms_tcm/fleet_scalar_z.parquet)')
  106. ap.add_argument('--out', default=None, help='输出路径(默认写回 m5/fleet_scalar_z.parquet)')
  107. ap.add_argument('--min-n', type=int, default=3, help='每组最少记录数(样件口径 = 3)')
  108. a = ap.parse_args()
  109. if a.verify:
  110. return verify(a.window, pathlib.Path(a.sample) if a.sample else None)
  111. df = compute(a.window, a.min_n)
  112. out = pathlib.Path(a.out) if a.out else (P.m5() / 'fleet_scalar_z.parquet')
  113. out.parent.mkdir(parents=True, exist_ok=True)
  114. df.to_parquet(out, index=False)
  115. print(f'已写 {P.rel(out)}: {len(df)} 行 × {len(df.columns)} 列(窗 {a.window})')
  116. # 自登记 (产物来源自登记: 谁算的谁登记)
  117. # ★2026-09-18 修: 原先按 `record(store, rel, builder=…, by=…)` 逐件传参, 而
  118. # src.derived_manifest.record 的签名是 `record(store_root, files: dict, by: str)` ——
  119. # 于是**登记从来没成功过**(被 except 吞成一行 [i] 提示): 产物在盘上、台账里没有来路。
  120. # 与 vib_raw_build.py 的调用形态对齐(那处传的是 dict)。
  121. try:
  122. from src import derived_manifest as DM
  123. DM.record(P.out_root(),
  124. {out.relative_to(P.out_root()).as_posix():
  125. 'scripts/rudong_fusion_run.py (val=median(scalar); z=(val-fleet_med)/(1.4826*MAD); '
  126. '逆向工程口径, 与随包样件逐值对拍 val 列 100%)'},
  127. by='rudong_fusion_run')
  128. print(' 已自登记 → _derived_manifest.json')
  129. except Exception as e:
  130. print(f' [i] 自登记跳过: {type(e).__name__}: {e}')
  131. return 0
  132. if __name__ == '__main__':
  133. for _s in (sys.stdout, sys.stderr):
  134. try:
  135. _s.reconfigure(errors='replace')
  136. except Exception:
  137. pass
  138. sys.exit(main())