| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334 |
- #!/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 等附列)
- <store>/<shard>.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` 就失联):
- <window>/spectra/spectra_meta.parquet ← 单窗 (本脚本 --out 指 windows/<w>/spectra)
- <store>/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<NN>/` 命名空间 (shard 是相对 store 的路径, 带子目录合法),
- 末尾父进程合并各分片的 meta。分片原子落盘 (先 .tmp 再改名), 中断不会留半截 npz。
- 用法:
- python scripts/rudong_tcm_spectra.py --root <measurement 目录> --out outputs/rudong/m5_cms_tcm/windows/w0317/spectra
- python scripts/rudong_tcm_spectra.py --root <dir> --out <store> --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='谱库目录 (通常 <window>/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() 对 <window>/spectra 优先读 <window>/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())
|