| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383 |
- #!/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 <dir> --out <m5>/windows/<w>/index.parquet
- ② rudong_tcm_spectra.py --root <dir> --out <m5>/windows/<w>/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": {"<ISO时间戳>": [{"Record": {...}}, ...]}}}
- `Record` 下才有真数据: `Site` / `Location` / `ConfigurationSettings` / `Sensor` / `Measurement`
- (后者再带 `Conditions` 与 `DataSets`)。文件布局实测两种, 本脚本都认:
- measurement/<YYYY>/<MM>/<WTGxx>/<WTGxx>_<uuid>_decode.json (本包; --root 指到 measurement 或更上层都行)
- <WTGxx>/<YYYY>/<MM>/<WTGxx>_<uuid>_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 (`<out>.parts/p<i>.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 <dir> --out <index.parquet> --jobs 12 --turbines WTG01,WTG02
- python scripts/rudong_tcm_index.py --root <dir> --out <index.parquet> --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())
|