| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133 |
- # -*- coding: utf-8 -*-
- r"""★C02 接线(用户令 2026-10-06):把**逐窗 level 序列**从多窗 L6 结果里生成,并过原厂迟滞规则。
- 背景:`sop/fusion_diag.py::apply_hysteresis(window_levels, h=3)` 早已存在(厂商 Hysteresis=3 的通用化:
- 连续 h 窗同向才改级,否则保持),但**一直缺入参** —— 全窗裁决不产生"逐窗 level 序列"。
- 本模块补上生产端:`outputs/rudong/m5_cms_tcm/model_run_l6.parquet` 现在带 `_窗`(5 个周窗)⇒
- 同一 (台, 线) 跨窗的 `定级` 即序列。
- 口径(如实):
- - 级别序:`候选·新发`(3) > `候选`(2) > `参考`(1);未在该窗过闸 ⇒ 该窗**无记录**(不补 0,序列只含实际过闸的窗)。
- - 迟滞由 `fusion_diag.apply_hysteresis` 执行(**不改写它的语义**);本模块只负责喂序列与呈现。
- - 用途:**诊断/呈现**(显示锁存);未接入任何判级,除非用户下令。
- """
- from __future__ import annotations
- import pathlib
- RANK = {'候选·新发': 3, '候选': 2, '参考': 1}
- def window_level_sequences(path=None) -> dict:
- """返回 {(台, 线): [(窗, 定级, 级别int), ...]}(按窗升序;只含实际过闸的窗)。"""
- import pandas as pd
- p = pathlib.Path(path) if path else pathlib.Path(__file__).resolve().parents[3] / \
- 'outputs/rudong/m5_cms_tcm/model_run_l6.parquet'
- if not p.exists():
- return {}
- df = pd.read_parquet(p)
- if '定级' not in df.columns or '_窗' not in df.columns:
- return {}
- df = df.assign(lv=df['定级'].map(lambda x: RANK.get(str(x), 0)))
- out = {}
- for (t, line), g in df.groupby(['台', '线']):
- s = g.sort_values('_窗')
- out[(str(t), str(line))] = [(str(w), str(lv), int(i)) for w, lv, i in zip(s['_窗'], s['定级'], s['lv'])]
- return out
- def latched(path=None, h: int = 3) -> dict:
- """对每条序列执行原厂迟滞;返回 {(台, 线): {"raw": [...], "latched": [...], "changed": bool}}。"""
- from . import fusion_diag as FD
- out = {}
- for key, seq in window_level_sequences(path).items():
- levels = [x[2] for x in seq]
- try:
- res = FD.apply_hysteresis(levels, h=h)
- except Exception as ex: # 实现异常不得静默: 如实带出
- out[key] = {'err': type(ex).__name__ + ': ' + str(ex)[:80], 'raw': levels}
- continue
- lat = res[1] if isinstance(res, tuple) and len(res) > 1 else res
- out[key] = {'raw': levels, 'latched': list(lat) if hasattr(lat, '__iter__') else lat,
- 'changed': list(lat) != levels if hasattr(lat, '__iter__') else None}
- return out
- def summary(path=None, h: int = 3) -> dict:
- d = latched(path, h)
- n = len(d)
- ch = [k for k, v in d.items() if v.get('changed')]
- err = [k for k, v in d.items() if v.get('err')]
- return {'n_seq': n, 'n_changed': len(ch), 'changed': [('%s/%s' % k) for k in ch[:10]],
- 'n_err': len(err), 'h': h,
- '口径': '原厂迟滞 Hysteresis=3:连续 3 窗同向才改级;序列只含实际过闸的窗'}
- # ── ★§5-3(2026-10-06):原厂"逐测量迟滞"读取 + 状态机 ────────────────────────────────
- # 事实基础(本轮实测):raw decode JSON 内确有数值型 RedHysteresis/YellowHysteresis/BlueHysteresis/
- # TrendHysteresis;抽样 120 文件得 RedHysteresis ∈ {3: 3330, 100: 30} ⇒ **"恒为 3"为假**,
- # 与上游报告(`mask.Hysteresis` 恒 3 为假;`rms_200` 一类上限 100)一致。
- _OEM_KEYS = ('RedHysteresis', 'YellowHysteresis', 'BlueHysteresis', 'TrendHysteresis')
- def oem_hysteresis_scan(root=None, limit=120) -> dict:
- """抽样扫描 raw decode JSON,返回 {键: {值: 次数}} 与逐测量分布(只读,不落盘)。"""
- import collections
- import json
- import pathlib as _p
- import re as _re
- base = _p.Path(root) if root else _p.Path(__file__).resolve().parents[3] / \
- 'data/raw/如东/windcms/CMS_RuDong_CGN_202603-04/measurement'
- vals = {k: collections.Counter() for k in _OEM_KEYS}
- per_meas = {}
- n = 0
- if not base.exists():
- return {'err': 'raw 目录不存在', 'root': str(base)}
- for f in base.rglob('*_decode.json'):
- n += 1
- txt = f.read_text(encoding='utf-8', errors='replace')
- mname = None
- mm = _re.search(r'"name"\s*:\s*"([^"]{1,48})"', txt)
- if mm:
- mname = mm.group(1)
- for k in _OEM_KEYS:
- for m2 in _re.finditer(r'"%s"\s*:\s*"?([^",}\]]{1,12})"?' % k, txt):
- try:
- v = int(float(m2.group(1)))
- except Exception:
- v = m2.group(1)
- vals[k][v] += 1
- if mname:
- per_meas.setdefault(mname, {}).setdefault(k, set()).add(v)
- if n >= limit:
- break
- return {'n_files': n,
- 'dist': {k: dict(v.most_common(6)) for k, v in vals.items()},
- 'odd': {m: {k: sorted(v) for k, v in d.items() if any(x != 3 for x in v)}
- for m, d in per_meas.items() if any(any(x != 3 for x in v) for v in d.values())}}
- def oem_latch(levels, *, h=3, upper=None):
- """原厂迟滞状态机(§5-3 口径):以数值 level 序列驱动。
- · 起报:连续 h 窗同向抬升且达到上限 upper ⇒ 起报;到 0 ⇒ 解除。
- · **缺口不当反例**:序列里以 None 表示"该窗无记录"(缺口)⇒ **保持当前状态**,不重置、不推进。
- 返回 (state, trace):state 为当前状态(0/起报值),trace 为逐窗状态。
- """
- cur = 0
- up = 0
- trace = []
- for lv in levels:
- if lv is None: # 缺口:保持,不计入连续计数
- trace.append(cur)
- continue
- if upper is not None and lv >= upper:
- up += 1
- if up >= h:
- cur = lv
- elif lv <= 0:
- cur = 0
- up = 0
- else:
- up = 0
- trace.append(cur)
- return cur, trace
|