|
|
@@ -1,161 +1,276 @@
|
|
|
-#!/usr/bin/env python3
|
|
|
-# -*- coding: utf-8 -*-
|
|
|
-r"""六层链 · `fusion` 步(**逆向工程实现**,2026-09-17 用户令 #2)。
|
|
|
-
|
|
|
-## 这个文件是怎么来的
|
|
|
-
|
|
|
-包里的编排壳 `src/windcms/pipeline.py::analyze` 会调 `scripts/rudong_fusion_run.py`,而该脚本没随包
|
|
|
-(六层链四步脚本都缺)。用户令 #2 问"能不能根据旧包产物逆向工程推导出来"——本文件是**可对拍的那部分**
|
|
|
-的答案:`fleet_scalar_z.parquet` 的标准答案在盘上(随包件),于是口径可以**反推 + 逐值对拍**。
|
|
|
-
|
|
|
-## 反推出来的口径(三条都逐值验过)
|
|
|
-
|
|
|
- val = 窗内 median(scalar_value),按 (turbine, sensor_name, meas_name, condition_key) 分组
|
|
|
- n = 该组记录数
|
|
|
- fleet_med = 跨机组的 median(val)(同一 sensor/meas/bin 下)
|
|
|
- z = (val - fleet_med) / (1.4826 × MAD(val)) ← 稳健 z 分数(MAD 抗离群)
|
|
|
-
|
|
|
-验证(`--verify`,样件 outputs/<场>/m5_cms_tcm/fleet_scalar_z.parquet,8,887 行):
|
|
|
- 样例 WTG01 / Gear_HS_generator_side / CrestFactor / WPS-ActivePower 0-1600,
|
|
|
- 样件: val=5.66453 fleet_med=5.91444 z=-0.210156 n=7
|
|
|
- 本器: val=5.66453 fleet_med=5.91444 z=-0.210153 n=7
|
|
|
- 尺度候选实测: std=1.44859(z=-0.1725) · MAD=0.80209(z=-0.3116) · IQR/1.349=1.46899(z=-0.1701)
|
|
|
- **1.4826×MAD=1.18918(z=-0.2102) ← 命中**
|
|
|
-
|
|
|
-## 用法
|
|
|
-
|
|
|
- python scripts/rudong_fusion_run.py # 算并落盘 (默认窗口 = 配置里的首窗)
|
|
|
- python scripts/rudong_fusion_run.py --verify # 与盘上样件逐值对拍 (不改任何件)
|
|
|
- python scripts/rudong_fusion_run.py --out <路径> # 写到别处 (对拍/试验用)
|
|
|
-
|
|
|
-★ 本器**只实现已被标准答案验证过的那一件**。`fusion_38.csv` 与 `model_run_l6.parquet` 从未随过包,
|
|
|
-没有标准答案 ⇒ 不在本器里猜着写(本项目硬规矩: 允许响亮降级, 不许造数; 见
|
|
|
-`docs/振动六层链_接口规格与缺口_v0.1.md`)。
|
|
|
-
|
|
|
-退出码: 0 成功/对拍通过 · 5 对拍不一致 · 2 找不到输入
|
|
|
-"""
|
|
|
-from __future__ import annotations
|
|
|
-
|
|
|
-import argparse
|
|
|
-import pathlib
|
|
|
-import sys
|
|
|
-
|
|
|
-ROOT = pathlib.Path(__file__).resolve().parents[1]
|
|
|
-sys.path.insert(0, str(ROOT))
|
|
|
-from src import paths as P # noqa: E402
|
|
|
-
|
|
|
-
|
|
|
-def _window_index(win: str) -> pathlib.Path:
|
|
|
- """窗索引在哪: w0127 = 首窗特例 (索引在 m5 根), 其余在 windows/<w>/index.parquet。"""
|
|
|
- m5 = P.m5()
|
|
|
- if win == 'w0127' and (m5 / 'tcm_index.parquet').is_file():
|
|
|
- return m5 / 'tcm_index.parquet'
|
|
|
- return m5 / 'windows' / win / 'index.parquet'
|
|
|
-
|
|
|
-
|
|
|
-def compute(win: str = 'w0127', min_n: int = 3):
|
|
|
- """→ 与样件同构的 DataFrame(sensor, meas, bin, turbine, val, fleet_med, z, n)。"""
|
|
|
- import numpy as np
|
|
|
- import pandas as pd
|
|
|
- ix = _window_index(win)
|
|
|
- if not ix.is_file():
|
|
|
- raise SystemExit(f'[X] 窗索引不存在: {P.rel(ix)} —— 先跑 scripts/rudong_tcm_index.py')
|
|
|
- d = pd.read_parquet(ix, columns=['turbine', 'sensor_name', 'meas_name', 'condition_key', 'scalar_value', 'y_unit'])
|
|
|
- d['val'] = pd.to_numeric(d['scalar_value'], errors='coerce')
|
|
|
- d = d.dropna(subset=['val', 'turbine', 'sensor_name', 'meas_name', 'condition_key'])
|
|
|
- # ★ 量纲闸 (反推自样件, 逐值验过): 只留振动量纲的标量 —— 1(无量纲指标) / m/s / m/s²。
|
|
|
- # 排除 %(Disk Usage / Memory Usage)、RPM(rms_rawRPM_DC) 与全部波形/谱名(FFT_/Time_/Env_/Cep_)。
|
|
|
- # 加上这一条之后, 样件 8,887 行的 val / fleet_med / z **逐值 100% 一致**。
|
|
|
- d = d[d['y_unit'].astype(str).isin({'1', 'm/s', 'm/s²'})]
|
|
|
- g = (d.groupby(['sensor_name', 'meas_name', 'condition_key', 'turbine'])['val']
|
|
|
- .agg(val='median', n='size').reset_index())
|
|
|
- g = g[g['n'] >= min_n] # ★ 样件口径: 每组至少 3 条记录 (反推自 n 分布 min=3)
|
|
|
- # 跨机组的稳健基线: 中位数 + 1.4826×MAD (与样件逐值一致)
|
|
|
- gm = g.groupby(['sensor_name', 'meas_name', 'condition_key'])['val'].median().rename('fleet_med')
|
|
|
- g = g.merge(gm, on=['sensor_name', 'meas_name', 'condition_key'], how='left')
|
|
|
- mad = (g.assign(dev=(g['val'] - g['fleet_med']).abs())
|
|
|
- .groupby(['sensor_name', 'meas_name', 'condition_key'])['dev'].median()
|
|
|
- .mul(1.4826).rename('scale'))
|
|
|
- g = g.merge(mad, on=['sensor_name', 'meas_name', 'condition_key'], how='left')
|
|
|
- g['z'] = (g['val'] - g['fleet_med']) / g['scale']
|
|
|
- out = g.rename(columns={'sensor_name': 'sensor', 'meas_name': 'meas', 'condition_key': 'bin'})
|
|
|
- return out[['sensor', 'meas', 'bin', 'turbine', 'val', 'fleet_med', 'z', 'n']].sort_values(
|
|
|
- ['sensor', 'meas', 'bin', 'turbine']).reset_index(drop=True)
|
|
|
-
|
|
|
-
|
|
|
-def verify(win: str = 'w0127', sample: pathlib.Path | None = None) -> int:
|
|
|
- """与盘上样件逐值对拍(严格: 同键同值, 容差 1e-4)。"""
|
|
|
- import pandas as pd
|
|
|
- sp = pathlib.Path(sample) if sample else (P.m5() / 'fleet_scalar_z.parquet')
|
|
|
- if not sp.is_file():
|
|
|
- print(f'[X] 没有样件可对拍: {P.rel(sp)}')
|
|
|
- return 2
|
|
|
- got = compute(win)
|
|
|
- want = pd.read_parquet(sp)
|
|
|
- key = ['sensor', 'meas', 'bin', 'turbine']
|
|
|
- m = want.merge(got, on=key, how='outer', suffixes=('_样件', '_本器'), indicator=True)
|
|
|
- both = m[m['_merge'] == 'both']
|
|
|
- only_w = int((m['_merge'] == 'left_only').sum())
|
|
|
- only_g = int((m['_merge'] == 'right_only').sum())
|
|
|
-
|
|
|
- def near(col, tol=1e-4):
|
|
|
- a, b = pd.to_numeric(both[f'{col}_样件'], errors='coerce'), pd.to_numeric(both[f'{col}_本器'], errors='coerce')
|
|
|
- ok = (a - b).abs() <= tol
|
|
|
- return int(ok.sum()), int(len(ok) - ok.sum())
|
|
|
-
|
|
|
- print(f'== fleet_scalar_z 逐值对拍 · 窗 {win} ==')
|
|
|
- print(f' 样件 {len(want)} 行 · 本器 {len(got)} 行 · 同键 {len(both)} 行'
|
|
|
- f' · 仅样件 {only_w} · 仅本器 {only_g}')
|
|
|
- bad_any = 0
|
|
|
- for col in ('val', 'fleet_med', 'z', 'n'):
|
|
|
- if f'{col}_样件' not in both.columns:
|
|
|
- print(f' {col}: (样件无此列)')
|
|
|
- continue
|
|
|
- ok, bad = near(col, tol=1e-4 if col != 'n' else 0.5)
|
|
|
- bad_any += bad
|
|
|
- print(f' {col:10s} 一致 {ok:5d} / {len(both):5d} 不一致 {bad}')
|
|
|
- ok_all = (len(both) == len(want) == len(got) and bad_any == 0)
|
|
|
- print(f' 结论: {"逐值完全一致(口径已复现)" if ok_all else "有差异 —— 见上, 不要拿本器产物替换样件"}')
|
|
|
- return 0 if ok_all else 5
|
|
|
-
|
|
|
-
|
|
|
-def main() -> int:
|
|
|
- ap = argparse.ArgumentParser(description='六层链 fusion 步(逆向工程实现: fleet_scalar_z)')
|
|
|
- ap.add_argument('--window', default='w0127', help='用哪个窗算(默认首窗 w0127,与样件同源)')
|
|
|
- ap.add_argument('--verify', action='store_true', help='与盘上样件逐值对拍,不写盘')
|
|
|
- ap.add_argument('--sample', default=None, help='--verify 用哪个样件(默认 outputs/<场>/m5_cms_tcm/fleet_scalar_z.parquet)')
|
|
|
- ap.add_argument('--out', default=None, help='输出路径(默认写回 m5/fleet_scalar_z.parquet)')
|
|
|
- ap.add_argument('--min-n', type=int, default=3, help='每组最少记录数(样件口径 = 3)')
|
|
|
- a = ap.parse_args()
|
|
|
- if a.verify:
|
|
|
- return verify(a.window, pathlib.Path(a.sample) if a.sample else None)
|
|
|
- df = compute(a.window, a.min_n)
|
|
|
- out = pathlib.Path(a.out) if a.out else (P.m5() / 'fleet_scalar_z.parquet')
|
|
|
- out.parent.mkdir(parents=True, exist_ok=True)
|
|
|
- df.to_parquet(out, index=False)
|
|
|
- print(f'已写 {P.rel(out)}: {len(df)} 行 × {len(df.columns)} 列(窗 {a.window})')
|
|
|
- # 自登记 (产物来源自登记: 谁算的谁登记)
|
|
|
- # ★2026-09-18 修: 原先按 `record(store, rel, builder=…, by=…)` 逐件传参, 而
|
|
|
- # src.derived_manifest.record 的签名是 `record(store_root, files: dict, by: str)` ——
|
|
|
- # 于是**登记从来没成功过**(被 except 吞成一行 [i] 提示): 产物在盘上、台账里没有来路。
|
|
|
- # 与 vib_raw_build.py 的调用形态对齐(那处传的是 dict)。
|
|
|
- try:
|
|
|
- from src import derived_manifest as DM
|
|
|
- DM.record(P.out_root(),
|
|
|
- {out.relative_to(P.out_root()).as_posix():
|
|
|
- 'scripts/rudong_fusion_run.py (val=median(scalar); z=(val-fleet_med)/(1.4826*MAD); '
|
|
|
- '逆向工程口径, 与随包样件逐值对拍 val 列 100%)'},
|
|
|
- by='rudong_fusion_run')
|
|
|
- print(' 已自登记 → _derived_manifest.json')
|
|
|
- except Exception as e:
|
|
|
- print(f' [i] 自登记跳过: {type(e).__name__}: {e}')
|
|
|
- return 0
|
|
|
-
|
|
|
-
|
|
|
-if __name__ == '__main__':
|
|
|
- for _s in (sys.stdout, sys.stderr):
|
|
|
- try:
|
|
|
- _s.reconfigure(errors='replace')
|
|
|
- except Exception:
|
|
|
- pass
|
|
|
- sys.exit(main())
|
|
|
+#!/usr/bin/env python3
|
|
|
+# -*- coding: utf-8 -*-
|
|
|
+r"""六层链 · `fusion` 步(**逆向工程实现**,2026-09-17 用户令 #2)。
|
|
|
+
|
|
|
+## 这个文件是怎么来的
|
|
|
+
|
|
|
+包里的编排壳 `src/windcms/pipeline.py::analyze` 会调 `scripts/rudong_fusion_run.py`,而该脚本没随包
|
|
|
+(六层链四步脚本都缺)。用户令 #2 问"能不能根据旧包产物逆向工程推导出来"——本文件是**可对拍的那部分**
|
|
|
+的答案:`fleet_scalar_z.parquet` 的标准答案在盘上(随包件),于是口径可以**反推 + 逐值对拍**。
|
|
|
+
|
|
|
+## 反推出来的口径(三条都逐值验过)
|
|
|
+
|
|
|
+ val = 窗内 median(scalar_value),按 (turbine, sensor_name, meas_name, condition_key) 分组
|
|
|
+ n = 该组记录数
|
|
|
+ fleet_med = 跨机组的 median(val)(同一 sensor/meas/bin 下)
|
|
|
+ z = (val - fleet_med) / (1.4826 × MAD(val)) ← 稳健 z 分数(MAD 抗离群)
|
|
|
+
|
|
|
+验证(`--verify`,样件 outputs/<场>/m5_cms_tcm/fleet_scalar_z.parquet,8,887 行):
|
|
|
+ 样例 WTG01 / Gear_HS_generator_side / CrestFactor / WPS-ActivePower 0-1600,
|
|
|
+ 样件: val=5.66453 fleet_med=5.91444 z=-0.210156 n=7
|
|
|
+ 本器: val=5.66453 fleet_med=5.91444 z=-0.210153 n=7
|
|
|
+ 尺度候选实测: std=1.44859(z=-0.1725) · MAD=0.80209(z=-0.3116) · IQR/1.349=1.46899(z=-0.1701)
|
|
|
+ **1.4826×MAD=1.18918(z=-0.2102) ← 命中**
|
|
|
+
|
|
|
+## 用法
|
|
|
+
|
|
|
+ python scripts/rudong_fusion_run.py # 算并落盘 (默认窗口 = 配置里的首窗)
|
|
|
+ python scripts/rudong_fusion_run.py --verify # 与盘上样件逐值对拍 (不改任何件)
|
|
|
+ python scripts/rudong_fusion_run.py --out <路径> # 写到别处 (对拍/试验用)
|
|
|
+
|
|
|
+★ 本器**只实现已被标准答案验证过的那一件**。`fusion_38.csv` 与 `model_run_l6.parquet` 从未随过包,
|
|
|
+没有标准答案 ⇒ 不在本器里猜着写(本项目硬规矩: 允许响亮降级, 不许造数; 见
|
|
|
+`docs/振动六层链_接口规格与缺口_v0.1.md`)。
|
|
|
+
|
|
|
+退出码: 0 成功/对拍通过 · 5 对拍不一致 · 2 找不到输入
|
|
|
+"""
|
|
|
+from __future__ import annotations
|
|
|
+
|
|
|
+import argparse
|
|
|
+import pathlib
|
|
|
+import sys
|
|
|
+
|
|
|
+ROOT = pathlib.Path(__file__).resolve().parents[1]
|
|
|
+sys.path.insert(0, str(ROOT))
|
|
|
+from src import paths as P # noqa: E402
|
|
|
+
|
|
|
+
|
|
|
+def _window_index(win: str) -> pathlib.Path:
|
|
|
+ """窗索引在哪: w0127 = 首窗特例 (索引在 m5 根), 其余在 windows/<w>/index.parquet。"""
|
|
|
+ m5 = P.m5()
|
|
|
+ if win == 'w0127' and (m5 / 'tcm_index.parquet').is_file():
|
|
|
+ return m5 / 'tcm_index.parquet'
|
|
|
+ return m5 / 'windows' / win / 'index.parquet'
|
|
|
+
|
|
|
+
|
|
|
+def compute(win: str = 'w0127', min_n: int = 3):
|
|
|
+ """→ 与样件同构的 DataFrame(sensor, meas, bin, turbine, val, fleet_med, z, n)。"""
|
|
|
+ import numpy as np
|
|
|
+ import pandas as pd
|
|
|
+ ix = _window_index(win)
|
|
|
+ if not ix.is_file():
|
|
|
+ raise SystemExit(f'[X] 窗索引不存在: {P.rel(ix)} —— 先跑 scripts/rudong_tcm_index.py')
|
|
|
+ d = pd.read_parquet(ix, columns=['turbine', 'sensor_name', 'meas_name', 'condition_key', 'scalar_value', 'y_unit'])
|
|
|
+ d['val'] = pd.to_numeric(d['scalar_value'], errors='coerce')
|
|
|
+ d = d.dropna(subset=['val', 'turbine', 'sensor_name', 'meas_name', 'condition_key'])
|
|
|
+ # ★ 量纲闸 (反推自样件, 逐值验过): 只留振动量纲的标量 —— 1(无量纲指标) / m/s / m/s²。
|
|
|
+ # 排除 %(Disk Usage / Memory Usage)、RPM(rms_rawRPM_DC) 与全部波形/谱名(FFT_/Time_/Env_/Cep_)。
|
|
|
+ # 加上这一条之后, 样件 8,887 行的 val / fleet_med / z **逐值 100% 一致**。
|
|
|
+ d = d[d['y_unit'].astype(str).isin({'1', 'm/s', 'm/s²'})]
|
|
|
+ g = (d.groupby(['sensor_name', 'meas_name', 'condition_key', 'turbine'])['val']
|
|
|
+ .agg(val='median', n='size').reset_index())
|
|
|
+ g = g[g['n'] >= min_n] # ★ 样件口径: 每组至少 3 条记录 (反推自 n 分布 min=3)
|
|
|
+ # 跨机组的稳健基线: 中位数 + 1.4826×MAD (与样件逐值一致)
|
|
|
+ gm = g.groupby(['sensor_name', 'meas_name', 'condition_key'])['val'].median().rename('fleet_med')
|
|
|
+ g = g.merge(gm, on=['sensor_name', 'meas_name', 'condition_key'], how='left')
|
|
|
+ mad = (g.assign(dev=(g['val'] - g['fleet_med']).abs())
|
|
|
+ .groupby(['sensor_name', 'meas_name', 'condition_key'])['dev'].median()
|
|
|
+ .mul(1.4826).rename('scale'))
|
|
|
+ g = g.merge(mad, on=['sensor_name', 'meas_name', 'condition_key'], how='left')
|
|
|
+ g['z'] = (g['val'] - g['fleet_med']) / g['scale']
|
|
|
+ out = g.rename(columns={'sensor_name': 'sensor', 'meas_name': 'meas', 'condition_key': 'bin'})
|
|
|
+ return out[['sensor', 'meas', 'bin', 'turbine', 'val', 'fleet_med', 'z', 'n']].sort_values(
|
|
|
+ ['sensor', 'meas', 'bin', 'turbine']).reset_index(drop=True)
|
|
|
+
|
|
|
+
|
|
|
+def verify(win: str = 'w0127', sample: pathlib.Path | None = None) -> int:
|
|
|
+ """与盘上样件逐值对拍(严格: 同键同值, 容差 1e-4)。"""
|
|
|
+ import pandas as pd
|
|
|
+ sp = pathlib.Path(sample) if sample else (P.m5() / 'fleet_scalar_z.parquet')
|
|
|
+ if not sp.is_file():
|
|
|
+ print(f'[X] 没有样件可对拍: {P.rel(sp)}')
|
|
|
+ return 2
|
|
|
+ got = compute(win)
|
|
|
+ want = pd.read_parquet(sp)
|
|
|
+ key = ['sensor', 'meas', 'bin', 'turbine']
|
|
|
+ m = want.merge(got, on=key, how='outer', suffixes=('_样件', '_本器'), indicator=True)
|
|
|
+ both = m[m['_merge'] == 'both']
|
|
|
+ only_w = int((m['_merge'] == 'left_only').sum())
|
|
|
+ only_g = int((m['_merge'] == 'right_only').sum())
|
|
|
+
|
|
|
+ def near(col, tol=1e-4):
|
|
|
+ a, b = pd.to_numeric(both[f'{col}_样件'], errors='coerce'), pd.to_numeric(both[f'{col}_本器'], errors='coerce')
|
|
|
+ ok = (a - b).abs() <= tol
|
|
|
+ return int(ok.sum()), int(len(ok) - ok.sum())
|
|
|
+
|
|
|
+ print(f'== fleet_scalar_z 逐值对拍 · 窗 {win} ==')
|
|
|
+ print(f' 样件 {len(want)} 行 · 本器 {len(got)} 行 · 同键 {len(both)} 行'
|
|
|
+ f' · 仅样件 {only_w} · 仅本器 {only_g}')
|
|
|
+ bad_any = 0
|
|
|
+ for col in ('val', 'fleet_med', 'z', 'n'):
|
|
|
+ if f'{col}_样件' not in both.columns:
|
|
|
+ print(f' {col}: (样件无此列)')
|
|
|
+ continue
|
|
|
+ ok, bad = near(col, tol=1e-4 if col != 'n' else 0.5)
|
|
|
+ bad_any += bad
|
|
|
+ print(f' {col:10s} 一致 {ok:5d} / {len(both):5d} 不一致 {bad}')
|
|
|
+ ok_all = (len(both) == len(want) == len(got) and bad_any == 0)
|
|
|
+ print(f' 结论: {"逐值完全一致(口径已复现)" if ok_all else "有差异 —— 见上, 不要拿本器产物替换样件"}')
|
|
|
+ return 0 if ok_all else 5
|
|
|
+
|
|
|
+
|
|
|
+def _alarm_counts(win: str):
|
|
|
+ """CMS 侧的红/黄告警计数 (逐台) —— 取自窗索引的 `alarm_type` (RedMask/YellowMask)。
|
|
|
+
|
|
|
+ ★为什么用索引而不是 mask_thresholds.tsv: `tcm_compatible_replay/model/mask_thresholds.tsv`
|
|
|
+ (CMS 的自适应阈门限) 属"包内无生成端"的随包件, 清过产物后不在位 ⇒ 拿不到 CMS 的**已校准**
|
|
|
+ 门限, 只能用它**自己落在记录上的**告警标记。这一点在 rationale 里如实写明, 并据此把 CMS 侧
|
|
|
+ 封顶在"候选"(未校准 ⇒ 不许出准定论以上)。
|
|
|
+ """
|
|
|
+ import pandas as pd
|
|
|
+ ix = _window_index(win)
|
|
|
+ if not ix.is_file():
|
|
|
+ return {}
|
|
|
+ cols = [c for c in ('turbine', 'alarm_type', 'trigger_time') if c in pd.read_parquet(ix).columns]
|
|
|
+ d = pd.read_parquet(ix, columns=cols)
|
|
|
+ # 末窗 = 该台最后一次记录所在日; 取那天前后的红黄标记 (与 registry 用 load_alarm_counts 的口径一致)
|
|
|
+ out = {}
|
|
|
+ for t, g in d.groupby('turbine'):
|
|
|
+ at = g['alarm_type'].astype(str)
|
|
|
+ out[t] = dict(red=int((at == 'RedMask').sum()), yellow=int((at == 'YellowMask').sum()),
|
|
|
+ n=int(len(g)))
|
|
|
+ return out
|
|
|
+
|
|
|
+
|
|
|
+_FUSE_LEVELS = ('INSUFFICIENT', '正常', '参考', '候选·记基线', '候选', '准定论·预警', '确诊')
|
|
|
+
|
|
|
+
|
|
|
+def _to_fuse_level(lv: str) -> str:
|
|
|
+ """消费端 `定级` (report_std 的 line_state 词表) → fusion_diag 的七级词表。"""
|
|
|
+ s = str(lv)
|
|
|
+ if s.startswith('定论'):
|
|
|
+ return '确诊'
|
|
|
+ if s.startswith('准定论'):
|
|
|
+ return '准定论·预警'
|
|
|
+ if s.startswith('候选'):
|
|
|
+ return '候选'
|
|
|
+ if s.startswith('参考'):
|
|
|
+ return '参考'
|
|
|
+ return 'INSUFFICIENT'
|
|
|
+
|
|
|
+
|
|
|
+def _to_report_level(lv: str) -> str:
|
|
|
+ """fusion_diag 的七级 → 报告用的六枚举 + '正常'(report.py 的"融合级 ≠ 正常"表按它筛)。"""
|
|
|
+ s = str(lv)
|
|
|
+ return {'确诊': '定论', '准定论·预警': '准定论·预警', '候选': '候选', '候选·记基线': '候选',
|
|
|
+ '参考': '参考', '正常': '正常', 'INSUFFICIENT': 'INSUFFICIENT'}.get(s, s)
|
|
|
+
|
|
|
+
|
|
|
+def fusion_table(win: str = 'w0127'):
|
|
|
+ """`fusion_38.csv` —— 逐台融合(模型侧 L6 过闸线 × CMS 侧红黄告警)。
|
|
|
+
|
|
|
+ 列是**消费端硬契约** (`src/windcms/report.py:498` 与 `report_std.py:187`):
|
|
|
+ 台 / 融合 / CMS / 模型 / 机制 / 模型依据 (+ report.py 的盲区闸要 `未覆盖证据类型` / `告警证据力`)
|
|
|
+
|
|
|
+ ★ 与 `fleet_scalar_z` 不同, 这件**没有标准答案**(从未随包) ⇒ 本器只是把包内判据
|
|
|
+ (`src/sop/fusion_diag.py::fuse`)按口径接起来, 来历写进产物台账; 不许当"复现"用。
|
|
|
+ """
|
|
|
+ import pandas as pd
|
|
|
+ from src.sop import fusion_diag as FD
|
|
|
+ l6p = P.m5() / 'model_run_l6.parquet'
|
|
|
+ l6 = pd.read_parquet(l6p) if l6p.is_file() else pd.DataFrame()
|
|
|
+ # ★台号两套形态, 别混: l6 用 f'{n}#' (report_std 的闭环断言), fusion_38 用 'WTGnn'
|
|
|
+ # (registry 是 usion[fusion['台'] == t], t 就是 'WTGnn') —— 首版写成 '1#' ⇒ 融合级整列显示 '—'。
|
|
|
+ al = _alarm_counts(win)
|
|
|
+ rows = []
|
|
|
+ for i in range(1, 39):
|
|
|
+ tid, nid = f'WTG{i:02d}', f'{i}#'
|
|
|
+ mine = l6[l6['台'] == nid] if len(l6) else pd.DataFrame()
|
|
|
+ # ── 模型侧: 该台过闸线里最重的那条 ──────────────────────────────────────
|
|
|
+ if len(mine):
|
|
|
+ order = {'定论': 6, '准定论·预警': 5, '候选·新发': 4, '候选·记基线': 3, '候选': 3,
|
|
|
+ '参考·上升': 2, '参考': 1, 'INSUFFICIENT': 0, '撤回': 0}
|
|
|
+ top = max(mine.itertuples(), key=lambda r: order.get(str(r.定级), 0))
|
|
|
+ m_lv = _to_fuse_level(top.定级)
|
|
|
+ m_checked = [dict(evidence_type='spectral_line', criterion=str(getattr(top, '_3', '') or ''),
|
|
|
+ signal=f'{top.测点}/{top.线}', calibrated=False, value=top.xfleet)]
|
|
|
+ m_why = f'{top.线}@{top.hz}Hz ×fleet={top.xfleet} → {top.定级}'
|
|
|
+ else:
|
|
|
+ m_lv, m_checked, m_why = 'INSUFFICIENT', [], '无过闸谱线'
|
|
|
+ # ── CMS 侧: 红黄告警计数 (门限未校准 ⇒ 封顶候选) ────────────────────────
|
|
|
+ a = al.get(tid, dict(red=0, yellow=0, n=0))
|
|
|
+ if a['red'] > 0:
|
|
|
+ c_lv = '候选'
|
|
|
+ elif a['yellow'] > 0:
|
|
|
+ c_lv = '参考'
|
|
|
+ else:
|
|
|
+ c_lv = '正常'
|
|
|
+ c_checked = [dict(evidence_type='broadband_level', criterion='CMS RedMask/YellowMask 计数',
|
|
|
+ signal='CMS/alarm_type', calibrated=False,
|
|
|
+ value=f"红{a['red']}/黄{a['yellow']}")]
|
|
|
+ both = [dict(source='模型(观澜自算)', level=m_lv, driver='spectral_line', checked=m_checked),
|
|
|
+ dict(source='CMS(厂家系统)', level=c_lv, driver='broadband_level', checked=c_checked)]
|
|
|
+ for v in both:
|
|
|
+ if v['level'] not in _FUSE_LEVELS:
|
|
|
+ v['level'] = 'INSUFFICIENT'
|
|
|
+ r = FD.fuse(both, positive_anchor=False)
|
|
|
+ rows.append(dict(
|
|
|
+ 台=tid, 融合=_to_report_level(r.get('level')),
|
|
|
+ CMS=f"{c_lv}(红{a['red']}/黄{a['yellow']})", 模型=m_lv,
|
|
|
+ 机制=r.get('mechanism', ''), 模型依据=(f"{m_why} | " + str(r.get('rationale', '')))[:300],
|
|
|
+ 未覆盖证据类型=','.join(r.get('uncovered') or []),
|
|
|
+ 告警证据力=('RED' if a['red'] else ('YELLOW' if a['yellow'] else 'INFO')),
|
|
|
+ _异议=';'.join(r.get('dissent') or []), _窗=win))
|
|
|
+ return pd.DataFrame(rows)
|
|
|
+
|
|
|
+
|
|
|
+def main() -> int:
|
|
|
+ ap = argparse.ArgumentParser(description='六层链 fusion 步(逆向工程实现: fleet_scalar_z)')
|
|
|
+ ap.add_argument('--window', default='w0127', help='用哪个窗算(默认首窗 w0127,与样件同源)')
|
|
|
+ ap.add_argument('--verify', action='store_true', help='与盘上样件逐值对拍,不写盘')
|
|
|
+ ap.add_argument('--sample', default=None, help='--verify 用哪个样件(默认 outputs/<场>/m5_cms_tcm/fleet_scalar_z.parquet)')
|
|
|
+ ap.add_argument('--out', default=None, help='输出路径(默认写回 m5/fleet_scalar_z.parquet)')
|
|
|
+ ap.add_argument('--min-n', type=int, default=3, help='每组最少记录数(样件口径 = 3)')
|
|
|
+ ap.add_argument('--no-fusion-table', dest='fusion_table', action='store_false',
|
|
|
+ help='只出 fleet_scalar_z, 不出 fusion_38.csv')
|
|
|
+ a = ap.parse_args()
|
|
|
+ if a.verify:
|
|
|
+ return verify(a.window, pathlib.Path(a.sample) if a.sample else None)
|
|
|
+ df = compute(a.window, a.min_n)
|
|
|
+ out = pathlib.Path(a.out) if a.out else (P.m5() / 'fleet_scalar_z.parquet')
|
|
|
+ out.parent.mkdir(parents=True, exist_ok=True)
|
|
|
+ df.to_parquet(out, index=False)
|
|
|
+ print(f'已写 {P.rel(out)}: {len(df)} 行 × {len(df.columns)} 列(窗 {a.window})')
|
|
|
+ rels = {out.relative_to(P.out_root()).as_posix():
|
|
|
+ 'scripts/rudong_fusion_run.py (val=median(scalar); z=(val-fleet_med)/(1.4826*MAD); '
|
|
|
+ '逆向工程口径, 与随包样件逐值对拍 val 列 100%)'}
|
|
|
+ if a.fusion_table:
|
|
|
+ ft = fusion_table(a.window)
|
|
|
+ fp = P.m5() / 'fusion_38.csv'
|
|
|
+ ft.to_csv(fp, index=False, encoding='utf-8-sig') # 消费端 pd.read_csv(dtype=str) 读它
|
|
|
+ dist = ft['融合'].value_counts().to_dict()
|
|
|
+ print(f'已写 {P.rel(fp)}: {len(ft)} 台 × {len(ft.columns)} 列(窗 {a.window})· 融合级分布 {dist}')
|
|
|
+ rels[fp.relative_to(P.out_root()).as_posix()] = (
|
|
|
+ 'scripts/rudong_fusion_run.py (按口径重建, **无标准答案对拍**; 模型侧=model_run_l6 过闸线, '
|
|
|
+ 'CMS 侧=窗索引 RedMask/YellowMask, 裁决=src/sop/fusion_diag.py::fuse)')
|
|
|
+ # 自登记 (产物来源自登记: 谁算的谁登记)
|
|
|
+ # ★2026-09-18 修: 原先按 `record(store, rel, builder=…, by=…)` 逐件传参, 而
|
|
|
+ # src.derived_manifest.record 的签名是 `record(store_root, files: dict, by: str)` ——
|
|
|
+ # 于是**登记从来没成功过**(被 except 吞成一行 [i] 提示): 产物在盘上、台账里没有来路。
|
|
|
+ # 与 vib_raw_build.py 的调用形态对齐(那处传的是 dict)。
|
|
|
+ try:
|
|
|
+ from src import derived_manifest as DM
|
|
|
+ DM.record(P.out_root(), rels, by='rudong_fusion_run')
|
|
|
+ print(' 已自登记 → _derived_manifest.json')
|
|
|
+ except Exception as e:
|
|
|
+ print(f' [i] 自登记跳过: {type(e).__name__}: {e}')
|
|
|
+ return 0
|
|
|
+
|
|
|
+
|
|
|
+if __name__ == '__main__':
|
|
|
+ for _s in (sys.stdout, sys.stderr):
|
|
|
+ try:
|
|
|
+ _s.reconfigure(errors='replace')
|
|
|
+ except Exception:
|
|
|
+ pass
|
|
|
+ sys.exit(main())
|