rudong_tcm_spectra.py 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334
  1. #!/usr/bin/env python3
  2. # -*- coding: utf-8 -*-
  3. """TCM 解码 JSON → 谱库 (npz 分片 + spectra_meta.parquet) —— pipeline.py 的 tcm_decoded_json 第二步.
  4. ## 契约 (消费者 = src/windcms/data.py::spectra_meta / spectrum)
  5. spectra_meta.parquet 每行一条谱: turbine, sensor, meas_name, trigger_time, shard, shard_row,
  6. x_offset, x_delta (+ rpm/y_unit/alarm_type/n_points 等附列)
  7. <store>/<shard>.npz key='values' = 二维数组 (n 条 × n 点); 第 shard_row 行是该条谱
  8. 取数: x = x_offset + arange(len(v)) * x_delta, v = npz['values'][shard_row]
  9. `data.spectra_meta()` 找 meta 的两条路径 (都写, 免得换 `--out` 就失联):
  10. <window>/spectra/spectra_meta.parquet ← 单窗 (本脚本 --out 指 windows/<w>/spectra)
  11. <store>/spectra_meta.parquet ← 合并库 (--out 指 m5/spectra 时)
  12. ## 输入
  13. 同上一步 (Brande TCM 导出 JSON)。**谱在 `Record.Measurement.DataSets.DataSet.Values` 里, 是空格
  14. 分隔的数值字符串** (不是数组! 实测 6401 点 ~ 60 KB 字符串/条), Size = Lines + 1, X-axisDelta =
  15. 带宽/Lines, X-axisUnit='Hz'。标量记录 Size==1 (其 Values 就是标量值, 归索引脚本处理, 这里跳过)。
  16. ## 只转需要的测量
  17. 原始导出里绝大部分字节是 `Time_*` 波形 (65536/200000 点), 而页面/判据用的是 `FFT_*` 谱。
  18. 默认 `--meas FFT_` 只转谱; 要全转用 `--meas ALL` (磁盘会显著变大: 实测单个窗的谱 npz 约 11 GB 量级)。
  19. ## 落盘与并行
  20. 按 (测量名, 点数) 分组装分片, 每片 `--shard-rows` 条 (默认 256) 写一个 npz (float32, 见 --dtype)。
  21. 并行按机组切片, 子进程写自己的 `p<NN>/` 命名空间 (shard 是相对 store 的路径, 带子目录合法),
  22. 末尾父进程合并各分片的 meta。分片原子落盘 (先 .tmp 再改名), 中断不会留半截 npz。
  23. 用法:
  24. python scripts/rudong_tcm_spectra.py --root <measurement 目录> --out outputs/rudong/m5_cms_tcm/windows/w0317/spectra
  25. python scripts/rudong_tcm_spectra.py --root <dir> --out <store> --meas FFT_ --jobs 8 --turbines WTG09
  26. """
  27. from __future__ import annotations
  28. import argparse
  29. import json
  30. import math
  31. import os
  32. import pathlib
  33. import shutil
  34. import sys
  35. import time
  36. import zipfile
  37. ROOT = pathlib.Path(__file__).resolve().parents[1]
  38. sys.path.insert(0, str(ROOT))
  39. from src.console import soft # noqa: E402
  40. soft()
  41. def _f(v):
  42. if v is None or v == '':
  43. return math.nan
  44. try:
  45. return float(v)
  46. except (TypeError, ValueError):
  47. return math.nan
  48. def _datasets(m):
  49. """Measurement.DataSets → (DataSets 层, DataSet 记录)。X/Y 轴在 DataSets 层, Size/Values 在 DataSet 里。"""
  50. dss = m.get('DataSets') or {}
  51. if not isinstance(dss, dict):
  52. return {}, {}
  53. ds = dss.get('DataSet')
  54. if isinstance(ds, list):
  55. ds = ds[0] if ds else {}
  56. return dss, (ds if isinstance(ds, dict) else {})
  57. def values_to_array(vals):
  58. """Values (空格分隔字符串 / list) → float 数组; 不合法返回 None。"""
  59. import numpy as np
  60. if isinstance(vals, list):
  61. try:
  62. return np.asarray(vals, dtype=float)
  63. except (TypeError, ValueError):
  64. return None
  65. if not isinstance(vals, str):
  66. return None
  67. try:
  68. arr = np.fromstring(vals, sep=' ') # C 级解析, 6401 点 ~ 0.2 ms
  69. except Exception:
  70. return None
  71. if arr.size == 0:
  72. try:
  73. arr = np.asarray(vals.split(), dtype=float)
  74. except (TypeError, ValueError):
  75. return None
  76. return arr if arr.size else None
  77. class ShardWriter:
  78. """按 (meas_name, 点数) 攒批写 npz。"""
  79. def __init__(self, store: pathlib.Path, ns: str, shard_rows: int, dtype):
  80. self.store = store
  81. self.ns = ns
  82. self.shard_rows = shard_rows
  83. self.dtype = dtype
  84. self.buf = {} # key → list[np.ndarray]
  85. self.idx = {} # key → 该 key 已写出片数
  86. self.meta = []
  87. self.shards = 0
  88. def add(self, key, turbine, sensor, meas, trigger, x_off, x_delta, extra, arr):
  89. rows = self.buf.setdefault(key, [])
  90. rows.append(arr)
  91. row_in_shard = len(rows) - 1
  92. shard = f'{self.ns}/{self._stem(key, self.idx.get(key, 0))}'
  93. self.meta.append(dict(turbine=turbine, sensor=sensor, meas_name=meas, trigger_time=trigger,
  94. shard=shard, shard_row=row_in_shard, x_offset=x_off, x_delta=x_delta,
  95. n_points=int(arr.size), **extra))
  96. if len(rows) >= self.shard_rows:
  97. self.flush_key(key)
  98. def _stem(self, key, i):
  99. meas, npts = key
  100. safe = ''.join(c if (c.isalnum() or c in '-_.') else '_' for c in meas)
  101. return f'{safe}_n{npts}_{i:04d}.npz'
  102. def flush_key(self, key):
  103. import numpy as np
  104. rows = self.buf.get(key)
  105. if not rows:
  106. return
  107. i = self.idx.get(key, 0)
  108. rel = f'{self.ns}/{self._stem(key, i)}'
  109. dst = self.store / rel
  110. dst.parent.mkdir(parents=True, exist_ok=True)
  111. # 不同测量/点数的分片点数一致 (同 key 同长度), 可直接 stack
  112. arr = np.stack(rows).astype(self.dtype, copy=False)
  113. tmp = dst.with_suffix('.npz.tmp')
  114. with open(tmp, 'wb') as fh:
  115. np.savez_compressed(fh, values=arr)
  116. os.replace(tmp, dst)
  117. self.shards += 1
  118. self.idx[key] = i + 1
  119. self.buf[key] = []
  120. def close(self):
  121. for key in list(self.buf):
  122. self.flush_key(key)
  123. return self.meta, self.shards
  124. def rows_of_file(raw: bytes, meas_prefix: str, writer: ShardWriter, stats: dict):
  125. doc = json.loads(raw.decode('utf-8', 'replace'))
  126. body = ((doc.get('body') or {}).get('body')) or {}
  127. if not isinstance(body, dict):
  128. stats['bad'] += 1
  129. return
  130. for _ts, recs in body.items():
  131. if not isinstance(recs, list):
  132. continue
  133. for item in recs:
  134. rec = item.get('Record') if isinstance(item, dict) else None
  135. if not isinstance(rec, dict):
  136. stats['bad'] += 1
  137. continue
  138. try:
  139. m = rec.get('Measurement') or {}
  140. meas = str(m.get('MeasurementName') or '?')
  141. if meas_prefix != 'ALL' and not meas.startswith(meas_prefix):
  142. stats['skipped'] += 1
  143. continue
  144. dss, ds = _datasets(m)
  145. size = _f(ds.get('Size'))
  146. if not (size and size > 1):
  147. stats['skipped'] += 1 # 标量行归索引脚本
  148. continue
  149. arr = values_to_array(ds.get('Values'))
  150. if arr is None or arr.size <= 1:
  151. stats['bad'] += 1
  152. continue
  153. loc = (rec.get('Location') or {}).get('LocationName') or '?'
  154. sens = (rec.get('Sensor') or {}).get('SensorName') or '?'
  155. extra = dict(rpm=_f(m.get('RPM')), y_unit=dss.get('Y-axisUnit'),
  156. x_unit=dss.get('X-axisUnit'), alarm_type=m.get('AlarmType'),
  157. condition_key=m.get('ConditionKey'), ds_size=size,
  158. meas_key=m.get('MeasurementKey'))
  159. writer.add((meas, int(arr.size)), loc, sens, meas, m.get('Trigger_Time'),
  160. _f(dss.get('X-axisOffset')), _f(dss.get('X-axisDelta')), extra, arr)
  161. stats['rows'] += 1
  162. except Exception:
  163. stats['bad'] += 1
  164. def worker(payload):
  165. part_path, items, store, ns, shard_rows, meas, dtype = payload
  166. import pandas as pd
  167. writer = ShardWriter(pathlib.Path(store), ns, shard_rows, dtype)
  168. stats = dict(rows=0, bad=0, skipped=0)
  169. cur_zip = None
  170. cur_path = None
  171. try:
  172. for name, full, zpath in items:
  173. try:
  174. if zpath:
  175. if cur_zip is None or cur_path != zpath:
  176. if cur_zip:
  177. cur_zip.close()
  178. cur_zip = zipfile.ZipFile(zpath)
  179. cur_path = zpath
  180. raw = cur_zip.read(name)
  181. else:
  182. raw = pathlib.Path(full).read_bytes()
  183. rows_of_file(raw, meas, writer, stats)
  184. except Exception:
  185. stats['bad'] += 1
  186. finally:
  187. if cur_zip:
  188. cur_zip.close()
  189. meta, shards = writer.close()
  190. pd.DataFrame(meta).to_parquet(part_path, index=False)
  191. row = dict(part=part_path, rows=stats['rows'], bad=stats['bad'], skipped=stats['skipped'],
  192. shards=shards, files=len(items))
  193. return row
  194. def discover(root: pathlib.Path, turbines, limit):
  195. out = []
  196. if root.is_file() and root.suffix.lower() == '.zip':
  197. with zipfile.ZipFile(root) as zf:
  198. for i in zf.infolist():
  199. if not i.is_dir() and i.filename.endswith('_decode.json'):
  200. out.append((i.filename, None, str(root)))
  201. else:
  202. for p in sorted(root.rglob('*_decode.json')):
  203. out.append((str(p.relative_to(root)).replace('\\', '/'), str(p), None))
  204. if turbines:
  205. want = {t.upper() for t in turbines}
  206. out = [x for x in out if any(t in x[0].upper() for t in want)]
  207. return out[:limit] if limit else out
  208. def main() -> int:
  209. ap = argparse.ArgumentParser()
  210. ap.add_argument('--root', required=True, help='含 *_decode.json 的目录, 或导出 zip')
  211. ap.add_argument('--out', required=True, help='谱库目录 (通常 <window>/spectra)')
  212. ap.add_argument('--meas', default='FFT_', help="只转这些测量 (前缀匹配); 'ALL' = 全转 (默认 %(default)s)")
  213. ap.add_argument('--jobs', type=int, default=min(8, (os.cpu_count() or 4)), help='并行进程数')
  214. ap.add_argument('--shard-rows', type=int, default=256, help='每个 npz 装多少条谱 (默认 %(default)s)')
  215. ap.add_argument('--dtype', default='float32', choices=('float32', 'float64'),
  216. help='谱点存储精度 (默认 float32; 判据标量不从这里算, 谱图显示足够)')
  217. ap.add_argument('--turbines', default=None, help='只处理这些机组 (逗号分隔)')
  218. ap.add_argument('--limit', type=int, default=0, help='只处理前 N 个文件 (冒烟)')
  219. ap.add_argument('--dry-run', action='store_true')
  220. a = ap.parse_args()
  221. root = pathlib.Path(a.root)
  222. if not root.exists():
  223. raise SystemExit(f'输入不存在: {root}')
  224. out = pathlib.Path(a.out)
  225. turbines = [x.strip() for x in a.turbines.split(',')] if a.turbines else None
  226. import numpy as np
  227. dtype = np.float32 if a.dtype == 'float32' else np.float64
  228. t0 = time.time()
  229. items = discover(root, turbines, a.limit)
  230. if not items:
  231. raise SystemExit(f'没找到 *_decode.json: {root}')
  232. print(f'输入: {root}\n发现 {len(items)} 个 *_decode.json; 测量过滤 {a.meas}; 谱库 → {out}')
  233. if a.dry_run:
  234. print('(dry-run, 未解析)')
  235. return 0
  236. import pandas as pd
  237. import concurrent.futures as cf
  238. groups = {}
  239. for it in items:
  240. key = next((seg for seg in it[0].replace('\\', '/').split('/') if seg.upper().startswith('WTG')), '_other')
  241. groups.setdefault(key, []).append(it)
  242. n = max(1, a.jobs)
  243. per = math.ceil(len(items) / n)
  244. chunks, cur = [], []
  245. for k in sorted(groups):
  246. cur.extend(groups[k])
  247. if len(cur) >= per:
  248. chunks.append(cur)
  249. cur = []
  250. if cur:
  251. chunks.append(cur)
  252. parts_dir = out.parent / (out.name + '.meta.parts')
  253. if parts_dir.exists():
  254. shutil.rmtree(parts_dir)
  255. parts_dir.mkdir(parents=True, exist_ok=True)
  256. payloads = [(str(parts_dir / f'p{i:02d}.parquet'), ch, str(out), f'p{i:02d}', a.shard_rows, a.meas, dtype)
  257. for i, ch in enumerate(chunks)]
  258. print(f'并行 {min(n, len(payloads))} 进程 × {len(payloads)} 片')
  259. rows_total = bad = skipped = shards = files_done = 0
  260. parts = []
  261. with cf.ProcessPoolExecutor(max_workers=min(n, len(payloads))) as ex:
  262. for r in ex.map(worker, payloads):
  263. parts.append(r['part'])
  264. rows_total += r['rows']
  265. bad += r['bad']
  266. skipped += r['skipped']
  267. shards += r['shards']
  268. files_done += r['files']
  269. print(f" [{files_done}/{len(items)} 文件] 谱 {rows_total} 条, 分片 {shards}, "
  270. f"跳过(非目标测量/标量) {skipped}, 异常 {bad} ({time.time() - t0:.0f}s)", flush=True)
  271. metas = [pd.read_parquet(p) for p in sorted(parts) if pathlib.Path(p).stat().st_size > 0]
  272. meta = pd.concat(metas, ignore_index=True) if metas else pd.DataFrame(
  273. columns=['turbine', 'sensor', 'meas_name', 'trigger_time', 'shard', 'shard_row', 'x_offset', 'x_delta'])
  274. for p in parts:
  275. pathlib.Path(p).unlink(missing_ok=True)
  276. parts_dir.rmdir()
  277. out.mkdir(parents=True, exist_ok=True)
  278. meta.to_parquet(out / 'spectra_meta.parquet', index=False)
  279. # 单窗约定: data.spectra_meta() 对 <window>/spectra 优先读 <window>/spectra_meta.parquet
  280. if out.name == 'spectra':
  281. meta.to_parquet(out.parent / 'spectra_meta.parquet', index=False)
  282. print(f'\n完成: {len(meta)} 条谱 → {out} (+ {out.parent / "spectra_meta.parquet" if out.name == "spectra" else "—"})')
  283. print(f' 分片文件 {shards} 个; 机组 {meta.turbine.nunique() if len(meta) else 0} 台; '
  284. f'测量 {meta.meas_name.nunique() if len(meta) else 0} 种')
  285. if len(meta):
  286. print(' 逐测量条数:\n' + meta.meas_name.value_counts().head(20).to_string())
  287. print(f' 用时 {time.time() - t0:.0f}s')
  288. return 0
  289. if __name__ == '__main__':
  290. sys.exit(main())