#!/usr/bin/env python3 # -*- coding: utf-8 -*- """TCM 解码 JSON → 窗索引 index.parquet (2026-09-12 补; `src/windcms/pipeline.py` 的 tcm_decoded_json 第一步). ## 它是谁、为什么在这里 `src/windcms/pipeline.py::ingest()` 对"含 `*_decode.json` 的目录"跑两步: ① rudong_tcm_index.py --root --out /windows//index.parquet ② rudong_tcm_spectra.py --root --out /windows//spectra 这两个脚本在 v0.2.0 里**缺失**(振动线分支 claude/vibration-data-diagnosis-32b69e 上的东西没随包), 于是"振动数据参与重算"这条路是断的: 只有产物 `tcm_index.parquet`(w0127 汇总) 在, 没有生成端。 2026-09-12 用户补入 CMS 原始导出 (`CMS_RuDong_CGN_202603-04.zip`, 25,679 件 Brande TCM 导出) 后, 按消费端契约把这两步补齐。 ## 输入 Brande TCM Enterprise 导出 (西门子机组自带 M-system), 每个文件是一次 API 响应: {"expiry", "buildId", "method", "controller", "serial", "requestInfo", "body": {"body": {"": [{"Record": {...}}, ...]}}} `Record` 下才有真数据: `Site` / `Location` / `ConfigurationSettings` / `Sensor` / `Measurement` (后者再带 `Conditions` 与 `DataSets`)。文件布局实测两种, 本脚本都认: measurement////__decode.json (本包; --root 指到 measurement 或更上层都行) ///__decode.json (w0127 那批的布局) ## 输出契约 (消费者 = src/windcms/data.py) `_read_index()` 只挑 `turbine, sensor_name, meas_name, trigger_time, rpm, condition_key, alarm_type, ds_size, scalar_value, overload`; `load_alarms` 读 `turbine, alarm_type`; 其余列是谱/工况/报警阈值元数据, 供风电场页面与后续六层链用。**列名与列义必须与包内 `outputs/<场>/m5_cms_tcm/tcm_index.parquet` (330,308 行 × 54 列, 2026-01-27~02-03 窗) 逐列对齐** —— 那份是同一摄取逻辑的产物, 是唯一的格式基准; 本脚本的 54 列与它同名同义 (见 COLSPEC), 因此新旧窗可以拼接进入同一分析集。 ## 并行与可重入 150 GB / 2.5 万件的单线程 json 解析要 ~1 小时, 机器 14 核 → 按机组切片并行 (`--jobs`, 默认 8)。 每个子进程写自己的分片 parquet (`.parts/p.parquet`), 父进程合并后删除分片。 单文件/单记录异常不中断整窗: 记进 `parse_error` 列 (响亮留痕), 文件级失败计数并在末尾汇总报告。 用法: python scripts/rudong_tcm_index.py --root data/raw/如东/windcms/CMS_RuDong_CGN_202603-04/measurement \ --out outputs/rudong/m5_cms_tcm/windows/w0317/index.parquet python scripts/rudong_tcm_index.py --root --out --jobs 12 --turbines WTG01,WTG02 python scripts/rudong_tcm_index.py --root --out --limit 200 --dry-run """ 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() # 厂内 JSON 层 → 索引列 (缺键给 None; 表里写的就是 JSON 里的键名, 便于逐列对拍) REC_MAP = { 'serial': ('Location', 'SerialNumber'), 'turbine': ('Location', 'LocationName'), 'config_name': ('ConfigurationSettings', 'ConfigurationName'), 'config_time': ('ConfigurationSettings', 'ConfigurationTime'), 'turbine_type': ('ConfigurationSettings', 'TurbineType'), 'recording_time_s': ('ConfigurationSettings', 'RecordingTime_s'), 'monitoring_cycle_min': ('ConfigurationSettings', 'MonitoringCycle-Minutes'), 'sensor_name': ('Sensor', 'SensorName'), 'sensor_type': ('Sensor', 'SensorType'), 'sensor_addr': ('Sensor', 'SensorAddress'), 'sensor_sn': ('Sensor', 'Sensor_Sn'), 'sensitivity_mvpg': ('Sensor', 'Sensitivity_mVpG'), 'sensor_unit': ('Sensor', 'SensorUnit'), 'hw_serial': ('Sensor', 'HwSerial'), 'meas_type': ('Measurement', 'MeasurementType'), 'meas_name': ('Measurement', 'MeasurementName'), 'meas_source': ('Measurement', 'Source'), 'category': ('Measurement', 'Category'), 'meas_key': ('Measurement', 'MeasurementKey'), 'trigger_time': ('Measurement', 'Trigger_Time'), 'trigger_ms': ('Measurement', 'Trigger_Time_ms'), 'duration_s': ('Measurement', 'Measurement_Duration-s'), 'rpm': ('Measurement', 'RPM'), 'nominal_freq_hz': ('Measurement', 'NominalFrequency-Hz'), 'min_freq_hz': ('Measurement', 'MinimumFrequency-Hz'), 'max_freq_hz': ('Measurement', 'MaximumFrequency-Hz'), 'lower_freq_hz': ('Measurement', 'LowerFrequency-Hz'), 'upper_freq_hz': ('Measurement', 'UpperFrequency-Hz'), 'overload': ('Measurement', 'Overload'), 'n_averages': ('Measurement', 'NumberOfAverages'), 'bandwidth_hz': ('Measurement', 'Bandwidth-Hz'), 'lines': ('Measurement', 'Lines'), 'integration': ('Measurement', 'Integration'), 'condition_name': ('Measurement', 'Conditions', 'ConditionsName'), 'condition_key': ('Measurement', 'ConditionKey'), 'alarm_type': ('Measurement', 'AlarmType'), 'red_alarm': ('Measurement', 'RedAlarm'), 'yellow_alarm': ('Measurement', 'YellowAlarm'), 'blue_alarm': ('Measurement', 'BlueAlarm'), 'red_hys': ('Measurement', 'RedHysAlarm'), 'yellow_hys': ('Measurement', 'YellowHysAlarm'), 'trend_hys': ('Measurement', 'TrendHysAlarm'), 'fault_freq': ('Measurement', 'FaultFreq'), } # DataSets 下的列。★层级别踩错 (2026-09-12 对拍逮到): `Size`/`Dimension`/`Values` 在 # `Measurement.DataSets.DataSet` 里, 而 X/Y 轴四件在**上一层** `Measurement.DataSets` 里 —— # 一开始全按 DataSet 取, 结果 x_offset/x_delta/x_unit/y_unit 四列整列为空, # 而 "x_delta == 带宽/lines" 的内部一致性检查比例是 0.000 (本该 1.000), 就是这里露的马脚。 DS_MAP = { 'ds_dim': 'Dimension', 'ds_size': 'Size', } AXIS_MAP = { 'x_offset': 'X-axisOffset', 'x_delta': 'X-axisDelta', 'x_unit': 'X-axisUnit', 'y_unit': 'Y-axisUnit', } TXT_COLS = ('file', 'serial', 'turbine', 'ts_key', 'config_name', 'config_time', 'turbine_type', 'sensor_name', 'sensor_type', 'sensor_addr', 'sensor_sn', 'sensor_unit', 'hw_serial', 'meas_type', 'meas_name', 'meas_source', 'category', 'meas_key', 'trigger_time', 'integration', 'condition_name', 'condition_key', 'alarm_type', 'red_hys', 'yellow_hys', 'trend_hys', 'x_unit', 'y_unit', 'parse_error') NUM_COLS = ('rec_i', 'recording_time_s', 'monitoring_cycle_min', 'sensitivity_mvpg', 'trigger_ms', 'duration_s', 'rpm', 'nominal_freq_hz', 'min_freq_hz', 'max_freq_hz', 'lower_freq_hz', 'upper_freq_hz', 'overload', 'n_averages', 'bandwidth_hz', 'lines', 'fault_freq', 'red_alarm', 'yellow_alarm', 'blue_alarm', 'ds_dim', 'ds_size', 'x_offset', 'x_delta', 'scalar_value') # 54 列的**列序**逐字抄自包内 `outputs/<场>/m5_cms_tcm/tcm_index.parquet` (振动线 v0.2.0 产物, # 330,308 行) —— 消费者按列名取数, 但"列序也要一致"是为了让新旧窗在人工对拍/并排打印时能直接比。 # (文本列与数值列是交错的, 所以不能用 TXT_COLS + NUM_COLS 拼。) COLSPEC = ['file', 'serial', 'turbine', 'ts_key', 'rec_i', 'config_name', 'config_time', 'turbine_type', 'recording_time_s', 'monitoring_cycle_min', 'sensor_name', 'sensor_type', 'sensor_addr', 'sensor_sn', 'sensitivity_mvpg', 'sensor_unit', 'hw_serial', 'meas_type', 'meas_name', 'meas_source', 'category', 'meas_key', 'trigger_time', 'trigger_ms', 'duration_s', 'rpm', 'nominal_freq_hz', 'min_freq_hz', 'max_freq_hz', 'lower_freq_hz', 'upper_freq_hz', 'overload', 'n_averages', 'bandwidth_hz', 'lines', 'integration', 'condition_name', 'condition_key', 'alarm_type', 'red_alarm', 'yellow_alarm', 'blue_alarm', 'red_hys', 'yellow_hys', 'trend_hys', 'fault_freq', 'ds_dim', 'ds_size', 'x_offset', 'x_delta', 'x_unit', 'y_unit', 'scalar_value', 'parse_error'] def _dig(d, path): """按键路径取 (任一环缺失返回 None, 不抛)。""" for k in path: if not isinstance(d, dict): return None d = d.get(k) return d def _f(v): """→ float; 空/非数 → NaN (索引列必须能进 pandas float 列, 字符串混进去会让整列变 object)。""" if v is None or v == '': return math.nan try: return float(v) except (TypeError, ValueError): return math.nan def _s(v): return None if v is None else str(v) def _datasets(m): """Measurement.DataSets → (DataSets 层, DataSet 记录)。DataSet 可能是 dict 也可能是 list。""" 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 rows_of_doc(doc, relpath, ts_key, recs, rows): """一个 (文件, 时间戳) 下的一批记录 → 追加进 rows。""" for rec_i, item in enumerate(recs): rec = item.get('Record') if isinstance(item, dict) else None if not isinstance(rec, dict): rows.append(dict(file=relpath, ts_key=ts_key, rec_i=rec_i, parse_error='记录非 Record 结构: ' + str(type(item).__name__))) continue try: m = rec.get('Measurement') or {} dss, ds = _datasets(m) vals = ds.get('Values') size = _f(ds.get('Size')) row = dict(file=relpath, ts_key=ts_key, rec_i=rec_i) for col, path in REC_MAP.items(): row[col] = _dig(rec, path) for col, key in DS_MAP.items(): row[col] = ds.get(key) for col, key in AXIS_MAP.items(): # X/Y 轴在 DataSets 层, 不在 DataSet 里 row[col] = dss.get(key) # 标量 (Size==1): Values 字符串本身就是标量值; 谱/波形 (Size>1) 的 Values 是长串, 值不进索引 row['scalar_value'] = _f(vals) if (size == 1 and isinstance(vals, (str, int, float))) else None row['parse_error'] = None rows.append(row) except Exception as exc: # 单记录异常不许断整窗 rows.append(dict(file=relpath, ts_key=ts_key, rec_i=rec_i, parse_error=f'{type(exc).__name__}: {exc}')) def rows_of_file(raw: bytes, relpath: str, rows: list): doc = json.loads(raw.decode('utf-8', 'replace')) body = ((doc.get('body') or {}).get('body')) or {} if not isinstance(body, dict): rows.append(dict(file=relpath, parse_error='body.body 非 dict')) return for ts_key, recs in body.items(): if isinstance(recs, list): rows_of_doc(doc, relpath, ts_key, recs, rows) else: rows.append(dict(file=relpath, ts_key=ts_key, parse_error='记录集非 list')) def relpath_of(name: str) -> str: """包内成员名 → `file` 列。去掉 `measurement/` 这一层 (包本身的组织层, 不是数据层)。""" n = name.replace('\\', '/') if n.startswith('measurement/'): n = n[len('measurement/'):] return n def discover(root: pathlib.Path, turbines, limit): """→ [(成员名/相对路径, 完整路径|None, zip 路径|None)]。目录与 zip 都支持。""" out = [] if root.is_file() and root.suffix.lower() == '.zip': with zipfile.ZipFile(root) as zf: for i in zf.infolist(): if i.is_dir() or not i.filename.endswith('_decode.json'): continue 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)] if limit: out = out[:limit] return out def flatten(items, rows): """逐文件解析 (目录: 直接读; zip: 按需打开一次)。""" cur_zip = None cur_path = None try: for name, full, zpath in items: rel = relpath_of(name) 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, rel, rows) except Exception as exc: rows.append(dict(file=rel, parse_error=f'文件级失败 {type(exc).__name__}: {exc}')) finally: if cur_zip: cur_zip.close() def worker(payload): """子进程: 解析分片 → 写分片 parquet → 回 (分片路径, 行数, 文件数, 失败数)。""" part_path, items = payload import pandas as pd rows = [] flatten(items, rows) df = pd.DataFrame(rows) for c in COLSPEC: if c not in df.columns: df[c] = None df = df[list(COLSPEC)] for c in NUM_COLS: # 与 shipped 一致的 float64 (整列同型; 有 NaN 的列本就会升为 float, 没有 NaN 的列不升则出现 int64) df[c] = pd.to_numeric(df[c], errors='coerce').astype('float64') for c in TXT_COLS: df[c] = df[c].astype('str') # pandas 3 的 'str' dtype (shipped 同款) bad = int(df['parse_error'].notna().sum()) if 'parse_error' in df else 0 df.to_parquet(part_path, index=False) return part_path, len(df), len(items), bad def main() -> int: ap = argparse.ArgumentParser() ap.add_argument('--root', required=True, help='含 *_decode.json 的目录, 或导出 zip') ap.add_argument('--out', required=True, help='index.parquet 输出路径') ap.add_argument('--jobs', type=int, default=min(8, (os.cpu_count() or 4)), help='并行进程数 (按机组切片; 默认 %(default)s)') ap.add_argument('--turbines', default=None, help='只处理这些机组 (逗号分隔, 如 WTG01,WTG02)') ap.add_argument('--limit', type=int, default=0, help='只处理前 N 个文件 (冒烟; 默认 0=全部)') ap.add_argument('--dry-run', action='store_true', help='只报将处理多少文件, 不解析') 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 t0 = time.time() items = discover(root, turbines, a.limit) if not items: raise SystemExit(f'没找到 *_decode.json: {root}') size = sum((pathlib.Path(f).stat().st_size if f else 0) for _, f, _ in items) print(f'输入: {root}') print(f'发现 {len(items)} 个 *_decode.json 目录内合计 {size / 1073741824:.2f} GB(仅目录模式可量)') 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) chunks = [] n = max(1, a.jobs) per = math.ceil(len(items) / n) 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.stem + '.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) for i, ch in enumerate(chunks)] print(f'并行 {min(n, len(payloads))} 进程 × {len(payloads)} 片 → {out}') done_files = 0 total_rows = 0 total_bad = 0 part_files = [] with cf.ProcessPoolExecutor(max_workers=min(n, len(payloads))) as ex: for part, nrow, nfile, nbad in ex.map(worker, payloads): part_files.append(part) done_files += nfile total_rows += nrow total_bad += nbad print(f' [{done_files}/{len(items)} 文件] 累计 {total_rows} 行, 异常 {total_bad} ' f'({time.time() - t0:.0f}s, {part})', flush=True) df = pd.concat([pd.read_parquet(p) for p in sorted(part_files)], ignore_index=True) out.parent.mkdir(parents=True, exist_ok=True) df.to_parquet(out, index=False) shutil.rmtree(parts_dir, ignore_errors=True) print(f'\n完成: {len(df)} 行 × {df.shape[1]} 列 → {out}') print(f' 机组 {df.turbine.nunique()} 台; 传感器 {df.sensor_name.nunique()} 种; 测量名 {df.meas_name.nunique()} 种') if 'trigger_time' in df: tt = pd.to_datetime(df.trigger_time, errors='coerce') print(f' 时间窗 {tt.min()} → {tt.max()}') print(f' 标量行(Size=1) {int((df.ds_size == 1).sum())}; 谱/波形行(Size>1) {int((df.ds_size > 1).sum())}; ' f'解析异常 {int(df.parse_error.notna().sum())}') if total_bad: print(f' ⚠ 有 {total_bad} 条记录带 parse_error (见该列), 未静默丢弃') print(f' 用时 {time.time() - t0:.0f}s') return 0 if __name__ == '__main__': sys.exit(main())