#!/usr/bin/env python3 # -*- coding: utf-8 -*- """TCM 解码 JSON → 谱库 (npz 分片 + spectra_meta.parquet) —— pipeline.py 的 tcm_decoded_json 第二步. ## 契约 (消费者 = src/windcms/data.py::spectra_meta / spectrum) spectra_meta.parquet 每行一条谱: turbine, sensor, meas_name, trigger_time, shard, shard_row, x_offset, x_delta (+ rpm/y_unit/alarm_type/n_points 等附列) /.npz key='values' = 二维数组 (n 条 × n 点); 第 shard_row 行是该条谱 取数: x = x_offset + arange(len(v)) * x_delta, v = npz['values'][shard_row] `data.spectra_meta()` 找 meta 的两条路径 (都写, 免得换 `--out` 就失联): /spectra/spectra_meta.parquet ← 单窗 (本脚本 --out 指 windows//spectra) /spectra_meta.parquet ← 合并库 (--out 指 m5/spectra 时) ## 输入 同上一步 (Brande TCM 导出 JSON)。**谱在 `Record.Measurement.DataSets.DataSet.Values` 里, 是空格 分隔的数值字符串** (不是数组! 实测 6401 点 ~ 60 KB 字符串/条), Size = Lines + 1, X-axisDelta = 带宽/Lines, X-axisUnit='Hz'。标量记录 Size==1 (其 Values 就是标量值, 归索引脚本处理, 这里跳过)。 ## 只转需要的测量 原始导出里绝大部分字节是 `Time_*` 波形 (65536/200000 点), 而页面/判据用的是 `FFT_*` 谱。 默认 `--meas FFT_` 只转谱; 要全转用 `--meas ALL` (磁盘会显著变大: 实测单个窗的谱 npz 约 11 GB 量级)。 ## 落盘与并行 按 (测量名, 点数) 分组装分片, 每片 `--shard-rows` 条 (默认 256) 写一个 npz (float32, 见 --dtype)。 并行按机组切片, 子进程写自己的 `p/` 命名空间 (shard 是相对 store 的路径, 带子目录合法), 末尾父进程合并各分片的 meta。分片原子落盘 (先 .tmp 再改名), 中断不会留半截 npz。 用法: python scripts/rudong_tcm_spectra.py --root --out outputs/rudong/m5_cms_tcm/windows/w0317/spectra python scripts/rudong_tcm_spectra.py --root --out --meas FFT_ --jobs 8 --turbines WTG09 """ from __future__ import annotations import argparse import json import math import os import pathlib import shutil import sys import time import zipfile ROOT = pathlib.Path(__file__).resolve().parents[1] sys.path.insert(0, str(ROOT)) from src.console import soft # noqa: E402 soft() def _f(v): if v is None or v == '': return math.nan try: return float(v) except (TypeError, ValueError): return math.nan def _datasets(m): """Measurement.DataSets → (DataSets 层, DataSet 记录)。X/Y 轴在 DataSets 层, Size/Values 在 DataSet 里。""" dss = m.get('DataSets') or {} if not isinstance(dss, dict): return {}, {} ds = dss.get('DataSet') if isinstance(ds, list): ds = ds[0] if ds else {} return dss, (ds if isinstance(ds, dict) else {}) def values_to_array(vals): """Values (空格分隔字符串 / list) → float 数组; 不合法返回 None。""" import numpy as np if isinstance(vals, list): try: return np.asarray(vals, dtype=float) except (TypeError, ValueError): return None if not isinstance(vals, str): return None try: arr = np.fromstring(vals, sep=' ') # C 级解析, 6401 点 ~ 0.2 ms except Exception: return None if arr.size == 0: try: arr = np.asarray(vals.split(), dtype=float) except (TypeError, ValueError): return None return arr if arr.size else None class ShardWriter: """按 (meas_name, 点数) 攒批写 npz。""" def __init__(self, store: pathlib.Path, ns: str, shard_rows: int, dtype): self.store = store self.ns = ns self.shard_rows = shard_rows self.dtype = dtype self.buf = {} # key → list[np.ndarray] self.idx = {} # key → 该 key 已写出片数 self.meta = [] self.shards = 0 def add(self, key, turbine, sensor, meas, trigger, x_off, x_delta, extra, arr): rows = self.buf.setdefault(key, []) rows.append(arr) row_in_shard = len(rows) - 1 shard = f'{self.ns}/{self._stem(key, self.idx.get(key, 0))}' self.meta.append(dict(turbine=turbine, sensor=sensor, meas_name=meas, trigger_time=trigger, shard=shard, shard_row=row_in_shard, x_offset=x_off, x_delta=x_delta, n_points=int(arr.size), **extra)) if len(rows) >= self.shard_rows: self.flush_key(key) def _stem(self, key, i): meas, npts = key safe = ''.join(c if (c.isalnum() or c in '-_.') else '_' for c in meas) return f'{safe}_n{npts}_{i:04d}.npz' def flush_key(self, key): import numpy as np rows = self.buf.get(key) if not rows: return i = self.idx.get(key, 0) rel = f'{self.ns}/{self._stem(key, i)}' dst = self.store / rel dst.parent.mkdir(parents=True, exist_ok=True) # 不同测量/点数的分片点数一致 (同 key 同长度), 可直接 stack arr = np.stack(rows).astype(self.dtype, copy=False) tmp = dst.with_suffix('.npz.tmp') with open(tmp, 'wb') as fh: np.savez_compressed(fh, values=arr) os.replace(tmp, dst) self.shards += 1 self.idx[key] = i + 1 self.buf[key] = [] def close(self): for key in list(self.buf): self.flush_key(key) return self.meta, self.shards def rows_of_file(raw: bytes, meas_prefix: str, writer: ShardWriter, stats: dict): doc = json.loads(raw.decode('utf-8', 'replace')) body = ((doc.get('body') or {}).get('body')) or {} if not isinstance(body, dict): stats['bad'] += 1 return for _ts, recs in body.items(): if not isinstance(recs, list): continue for item in recs: rec = item.get('Record') if isinstance(item, dict) else None if not isinstance(rec, dict): stats['bad'] += 1 continue try: m = rec.get('Measurement') or {} meas = str(m.get('MeasurementName') or '?') if meas_prefix != 'ALL' and not meas.startswith(meas_prefix): stats['skipped'] += 1 continue dss, ds = _datasets(m) size = _f(ds.get('Size')) if not (size and size > 1): stats['skipped'] += 1 # 标量行归索引脚本 continue arr = values_to_array(ds.get('Values')) if arr is None or arr.size <= 1: stats['bad'] += 1 continue loc = (rec.get('Location') or {}).get('LocationName') or '?' sens = (rec.get('Sensor') or {}).get('SensorName') or '?' extra = dict(rpm=_f(m.get('RPM')), y_unit=dss.get('Y-axisUnit'), x_unit=dss.get('X-axisUnit'), alarm_type=m.get('AlarmType'), condition_key=m.get('ConditionKey'), ds_size=size, meas_key=m.get('MeasurementKey')) writer.add((meas, int(arr.size)), loc, sens, meas, m.get('Trigger_Time'), _f(dss.get('X-axisOffset')), _f(dss.get('X-axisDelta')), extra, arr) stats['rows'] += 1 except Exception: stats['bad'] += 1 def worker(payload): part_path, items, store, ns, shard_rows, meas, dtype = payload import pandas as pd writer = ShardWriter(pathlib.Path(store), ns, shard_rows, dtype) stats = dict(rows=0, bad=0, skipped=0) cur_zip = None cur_path = None try: for name, full, zpath in items: try: if zpath: if cur_zip is None or cur_path != zpath: if cur_zip: cur_zip.close() cur_zip = zipfile.ZipFile(zpath) cur_path = zpath raw = cur_zip.read(name) else: raw = pathlib.Path(full).read_bytes() rows_of_file(raw, meas, writer, stats) except Exception: stats['bad'] += 1 finally: if cur_zip: cur_zip.close() meta, shards = writer.close() pd.DataFrame(meta).to_parquet(part_path, index=False) row = dict(part=part_path, rows=stats['rows'], bad=stats['bad'], skipped=stats['skipped'], shards=shards, files=len(items)) return row def discover(root: pathlib.Path, turbines, limit): out = [] if root.is_file() and root.suffix.lower() == '.zip': with zipfile.ZipFile(root) as zf: for i in zf.infolist(): if not i.is_dir() and i.filename.endswith('_decode.json'): out.append((i.filename, None, str(root))) else: for p in sorted(root.rglob('*_decode.json')): out.append((str(p.relative_to(root)).replace('\\', '/'), str(p), None)) if turbines: want = {t.upper() for t in turbines} out = [x for x in out if any(t in x[0].upper() for t in want)] return out[:limit] if limit else out def main() -> int: ap = argparse.ArgumentParser() ap.add_argument('--root', required=True, help='含 *_decode.json 的目录, 或导出 zip') ap.add_argument('--out', required=True, help='谱库目录 (通常 /spectra)') ap.add_argument('--meas', default='FFT_', help="只转这些测量 (前缀匹配); 'ALL' = 全转 (默认 %(default)s)") ap.add_argument('--jobs', type=int, default=min(8, (os.cpu_count() or 4)), help='并行进程数') ap.add_argument('--shard-rows', type=int, default=256, help='每个 npz 装多少条谱 (默认 %(default)s)') ap.add_argument('--dtype', default='float32', choices=('float32', 'float64'), help='谱点存储精度 (默认 float32; 判据标量不从这里算, 谱图显示足够)') ap.add_argument('--turbines', default=None, help='只处理这些机组 (逗号分隔)') ap.add_argument('--limit', type=int, default=0, help='只处理前 N 个文件 (冒烟)') ap.add_argument('--dry-run', action='store_true') a = ap.parse_args() root = pathlib.Path(a.root) if not root.exists(): raise SystemExit(f'输入不存在: {root}') out = pathlib.Path(a.out) turbines = [x.strip() for x in a.turbines.split(',')] if a.turbines else None import numpy as np dtype = np.float32 if a.dtype == 'float32' else np.float64 t0 = time.time() items = discover(root, turbines, a.limit) if not items: raise SystemExit(f'没找到 *_decode.json: {root}') print(f'输入: {root}\n发现 {len(items)} 个 *_decode.json; 测量过滤 {a.meas}; 谱库 → {out}') if a.dry_run: print('(dry-run, 未解析)') return 0 import pandas as pd import concurrent.futures as cf groups = {} for it in items: key = next((seg for seg in it[0].replace('\\', '/').split('/') if seg.upper().startswith('WTG')), '_other') groups.setdefault(key, []).append(it) n = max(1, a.jobs) per = math.ceil(len(items) / n) chunks, cur = [], [] for k in sorted(groups): cur.extend(groups[k]) if len(cur) >= per: chunks.append(cur) cur = [] if cur: chunks.append(cur) parts_dir = out.parent / (out.name + '.meta.parts') if parts_dir.exists(): shutil.rmtree(parts_dir) parts_dir.mkdir(parents=True, exist_ok=True) payloads = [(str(parts_dir / f'p{i:02d}.parquet'), ch, str(out), f'p{i:02d}', a.shard_rows, a.meas, dtype) for i, ch in enumerate(chunks)] print(f'并行 {min(n, len(payloads))} 进程 × {len(payloads)} 片') rows_total = bad = skipped = shards = files_done = 0 parts = [] with cf.ProcessPoolExecutor(max_workers=min(n, len(payloads))) as ex: for r in ex.map(worker, payloads): parts.append(r['part']) rows_total += r['rows'] bad += r['bad'] skipped += r['skipped'] shards += r['shards'] files_done += r['files'] print(f" [{files_done}/{len(items)} 文件] 谱 {rows_total} 条, 分片 {shards}, " f"跳过(非目标测量/标量) {skipped}, 异常 {bad} ({time.time() - t0:.0f}s)", flush=True) metas = [pd.read_parquet(p) for p in sorted(parts) if pathlib.Path(p).stat().st_size > 0] meta = pd.concat(metas, ignore_index=True) if metas else pd.DataFrame( columns=['turbine', 'sensor', 'meas_name', 'trigger_time', 'shard', 'shard_row', 'x_offset', 'x_delta']) for p in parts: pathlib.Path(p).unlink(missing_ok=True) parts_dir.rmdir() out.mkdir(parents=True, exist_ok=True) meta.to_parquet(out / 'spectra_meta.parquet', index=False) # 单窗约定: data.spectra_meta() 对 /spectra 优先读 /spectra_meta.parquet if out.name == 'spectra': meta.to_parquet(out.parent / 'spectra_meta.parquet', index=False) print(f'\n完成: {len(meta)} 条谱 → {out} (+ {out.parent / "spectra_meta.parquet" if out.name == "spectra" else "—"})') print(f' 分片文件 {shards} 个; 机组 {meta.turbine.nunique() if len(meta) else 0} 台; ' f'测量 {meta.meas_name.nunique() if len(meta) else 0} 种') if len(meta): print(' 逐测量条数:\n' + meta.meas_name.value_counts().head(20).to_string()) print(f' 用时 {time.time() - t0:.0f}s') return 0 if __name__ == '__main__': sys.exit(main())