#!/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())