hysteresis_wire.py 6.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133
  1. # -*- coding: utf-8 -*-
  2. r"""★C02 接线(用户令 2026-10-06):把**逐窗 level 序列**从多窗 L6 结果里生成,并过原厂迟滞规则。
  3. 背景:`sop/fusion_diag.py::apply_hysteresis(window_levels, h=3)` 早已存在(厂商 Hysteresis=3 的通用化:
  4. 连续 h 窗同向才改级,否则保持),但**一直缺入参** —— 全窗裁决不产生"逐窗 level 序列"。
  5. 本模块补上生产端:`outputs/rudong/m5_cms_tcm/model_run_l6.parquet` 现在带 `_窗`(5 个周窗)⇒
  6. 同一 (台, 线) 跨窗的 `定级` 即序列。
  7. 口径(如实):
  8. - 级别序:`候选·新发`(3) > `候选`(2) > `参考`(1);未在该窗过闸 ⇒ 该窗**无记录**(不补 0,序列只含实际过闸的窗)。
  9. - 迟滞由 `fusion_diag.apply_hysteresis` 执行(**不改写它的语义**);本模块只负责喂序列与呈现。
  10. - 用途:**诊断/呈现**(显示锁存);未接入任何判级,除非用户下令。
  11. """
  12. from __future__ import annotations
  13. import pathlib
  14. RANK = {'候选·新发': 3, '候选': 2, '参考': 1}
  15. def window_level_sequences(path=None) -> dict:
  16. """返回 {(台, 线): [(窗, 定级, 级别int), ...]}(按窗升序;只含实际过闸的窗)。"""
  17. import pandas as pd
  18. p = pathlib.Path(path) if path else pathlib.Path(__file__).resolve().parents[3] / \
  19. 'outputs/rudong/m5_cms_tcm/model_run_l6.parquet'
  20. if not p.exists():
  21. return {}
  22. df = pd.read_parquet(p)
  23. if '定级' not in df.columns or '_窗' not in df.columns:
  24. return {}
  25. df = df.assign(lv=df['定级'].map(lambda x: RANK.get(str(x), 0)))
  26. out = {}
  27. for (t, line), g in df.groupby(['台', '线']):
  28. s = g.sort_values('_窗')
  29. out[(str(t), str(line))] = [(str(w), str(lv), int(i)) for w, lv, i in zip(s['_窗'], s['定级'], s['lv'])]
  30. return out
  31. def latched(path=None, h: int = 3) -> dict:
  32. """对每条序列执行原厂迟滞;返回 {(台, 线): {"raw": [...], "latched": [...], "changed": bool}}。"""
  33. from . import fusion_diag as FD
  34. out = {}
  35. for key, seq in window_level_sequences(path).items():
  36. levels = [x[2] for x in seq]
  37. try:
  38. res = FD.apply_hysteresis(levels, h=h)
  39. except Exception as ex: # 实现异常不得静默: 如实带出
  40. out[key] = {'err': type(ex).__name__ + ': ' + str(ex)[:80], 'raw': levels}
  41. continue
  42. lat = res[1] if isinstance(res, tuple) and len(res) > 1 else res
  43. out[key] = {'raw': levels, 'latched': list(lat) if hasattr(lat, '__iter__') else lat,
  44. 'changed': list(lat) != levels if hasattr(lat, '__iter__') else None}
  45. return out
  46. def summary(path=None, h: int = 3) -> dict:
  47. d = latched(path, h)
  48. n = len(d)
  49. ch = [k for k, v in d.items() if v.get('changed')]
  50. err = [k for k, v in d.items() if v.get('err')]
  51. return {'n_seq': n, 'n_changed': len(ch), 'changed': [('%s/%s' % k) for k in ch[:10]],
  52. 'n_err': len(err), 'h': h,
  53. '口径': '原厂迟滞 Hysteresis=3:连续 3 窗同向才改级;序列只含实际过闸的窗'}
  54. # ── ★§5-3(2026-10-06):原厂"逐测量迟滞"读取 + 状态机 ────────────────────────────────
  55. # 事实基础(本轮实测):raw decode JSON 内确有数值型 RedHysteresis/YellowHysteresis/BlueHysteresis/
  56. # TrendHysteresis;抽样 120 文件得 RedHysteresis ∈ {3: 3330, 100: 30} ⇒ **"恒为 3"为假**,
  57. # 与上游报告(`mask.Hysteresis` 恒 3 为假;`rms_200` 一类上限 100)一致。
  58. _OEM_KEYS = ('RedHysteresis', 'YellowHysteresis', 'BlueHysteresis', 'TrendHysteresis')
  59. def oem_hysteresis_scan(root=None, limit=120) -> dict:
  60. """抽样扫描 raw decode JSON,返回 {键: {值: 次数}} 与逐测量分布(只读,不落盘)。"""
  61. import collections
  62. import json
  63. import pathlib as _p
  64. import re as _re
  65. base = _p.Path(root) if root else _p.Path(__file__).resolve().parents[3] / \
  66. 'data/raw/如东/windcms/CMS_RuDong_CGN_202603-04/measurement'
  67. vals = {k: collections.Counter() for k in _OEM_KEYS}
  68. per_meas = {}
  69. n = 0
  70. if not base.exists():
  71. return {'err': 'raw 目录不存在', 'root': str(base)}
  72. for f in base.rglob('*_decode.json'):
  73. n += 1
  74. txt = f.read_text(encoding='utf-8', errors='replace')
  75. mname = None
  76. mm = _re.search(r'"name"\s*:\s*"([^"]{1,48})"', txt)
  77. if mm:
  78. mname = mm.group(1)
  79. for k in _OEM_KEYS:
  80. for m2 in _re.finditer(r'"%s"\s*:\s*"?([^",}\]]{1,12})"?' % k, txt):
  81. try:
  82. v = int(float(m2.group(1)))
  83. except Exception:
  84. v = m2.group(1)
  85. vals[k][v] += 1
  86. if mname:
  87. per_meas.setdefault(mname, {}).setdefault(k, set()).add(v)
  88. if n >= limit:
  89. break
  90. return {'n_files': n,
  91. 'dist': {k: dict(v.most_common(6)) for k, v in vals.items()},
  92. 'odd': {m: {k: sorted(v) for k, v in d.items() if any(x != 3 for x in v)}
  93. for m, d in per_meas.items() if any(any(x != 3 for x in v) for v in d.values())}}
  94. def oem_latch(levels, *, h=3, upper=None):
  95. """原厂迟滞状态机(§5-3 口径):以数值 level 序列驱动。
  96. · 起报:连续 h 窗同向抬升且达到上限 upper ⇒ 起报;到 0 ⇒ 解除。
  97. · **缺口不当反例**:序列里以 None 表示"该窗无记录"(缺口)⇒ **保持当前状态**,不重置、不推进。
  98. 返回 (state, trace):state 为当前状态(0/起报值),trace 为逐窗状态。
  99. """
  100. cur = 0
  101. up = 0
  102. trace = []
  103. for lv in levels:
  104. if lv is None: # 缺口:保持,不计入连续计数
  105. trace.append(cur)
  106. continue
  107. if upper is not None and lv >= upper:
  108. up += 1
  109. if up >= h:
  110. cur = lv
  111. elif lv <= 0:
  112. cur = 0
  113. up = 0
  114. else:
  115. up = 0
  116. trace.append(cur)
  117. return cur, trace