rudong_tcm_index.py 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383
  1. #!/usr/bin/env python3
  2. # -*- coding: utf-8 -*-
  3. """TCM 解码 JSON → 窗索引 index.parquet (2026-09-12 补; `src/windcms/pipeline.py` 的 tcm_decoded_json 第一步).
  4. ## 它是谁、为什么在这里
  5. `src/windcms/pipeline.py::ingest()` 对"含 `*_decode.json` 的目录"跑两步:
  6. ① rudong_tcm_index.py --root <dir> --out <m5>/windows/<w>/index.parquet
  7. ② rudong_tcm_spectra.py --root <dir> --out <m5>/windows/<w>/spectra
  8. 这两个脚本在 v0.2.0 里**缺失**(振动线分支 claude/vibration-data-diagnosis-32b69e 上的东西没随包),
  9. 于是"振动数据参与重算"这条路是断的: 只有产物 `tcm_index.parquet`(w0127 汇总) 在, 没有生成端。
  10. 2026-09-12 用户补入 CMS 原始导出 (`CMS_RuDong_CGN_202603-04.zip`, 25,679 件 Brande TCM 导出) 后,
  11. 按消费端契约把这两步补齐。
  12. ## 输入
  13. Brande TCM Enterprise 导出 (西门子机组自带 M-system), 每个文件是一次 API 响应:
  14. {"expiry", "buildId", "method", "controller", "serial", "requestInfo",
  15. "body": {"body": {"<ISO时间戳>": [{"Record": {...}}, ...]}}}
  16. `Record` 下才有真数据: `Site` / `Location` / `ConfigurationSettings` / `Sensor` / `Measurement`
  17. (后者再带 `Conditions` 与 `DataSets`)。文件布局实测两种, 本脚本都认:
  18. measurement/<YYYY>/<MM>/<WTGxx>/<WTGxx>_<uuid>_decode.json (本包; --root 指到 measurement 或更上层都行)
  19. <WTGxx>/<YYYY>/<MM>/<WTGxx>_<uuid>_decode.json (w0127 那批的布局)
  20. ## 输出契约 (消费者 = src/windcms/data.py)
  21. `_read_index()` 只挑 `turbine, sensor_name, meas_name, trigger_time, rpm, condition_key, alarm_type,
  22. ds_size, scalar_value, overload`; `load_alarms` 读 `turbine, alarm_type`; 其余列是谱/工况/报警阈值元数据,
  23. 供风电场页面与后续六层链用。**列名与列义必须与包内 `outputs/<场>/m5_cms_tcm/tcm_index.parquet`
  24. (330,308 行 × 54 列, 2026-01-27~02-03 窗) 逐列对齐** —— 那份是同一摄取逻辑的产物, 是唯一的格式基准;
  25. 本脚本的 54 列与它同名同义 (见 COLSPEC), 因此新旧窗可以拼接进入同一分析集。
  26. ## 并行与可重入
  27. 150 GB / 2.5 万件的单线程 json 解析要 ~1 小时, 机器 14 核 → 按机组切片并行 (`--jobs`, 默认 8)。
  28. 每个子进程写自己的分片 parquet (`<out>.parts/p<i>.parquet`), 父进程合并后删除分片。
  29. 单文件/单记录异常不中断整窗: 记进 `parse_error` 列 (响亮留痕), 文件级失败计数并在末尾汇总报告。
  30. 用法:
  31. python scripts/rudong_tcm_index.py --root data/raw/如东/windcms/CMS_RuDong_CGN_202603-04/measurement \
  32. --out outputs/rudong/m5_cms_tcm/windows/w0317/index.parquet
  33. python scripts/rudong_tcm_index.py --root <dir> --out <index.parquet> --jobs 12 --turbines WTG01,WTG02
  34. python scripts/rudong_tcm_index.py --root <dir> --out <index.parquet> --limit 200 --dry-run
  35. """
  36. from __future__ import annotations
  37. import argparse
  38. import json
  39. import math
  40. import os
  41. import pathlib
  42. import shutil
  43. import sys
  44. import time
  45. import zipfile
  46. ROOT = pathlib.Path(__file__).resolve().parents[1]
  47. sys.path.insert(0, str(ROOT))
  48. from src.console import soft # noqa: E402
  49. soft()
  50. # 厂内 JSON 层 → 索引列 (缺键给 None; 表里写的就是 JSON 里的键名, 便于逐列对拍)
  51. REC_MAP = {
  52. 'serial': ('Location', 'SerialNumber'),
  53. 'turbine': ('Location', 'LocationName'),
  54. 'config_name': ('ConfigurationSettings', 'ConfigurationName'),
  55. 'config_time': ('ConfigurationSettings', 'ConfigurationTime'),
  56. 'turbine_type': ('ConfigurationSettings', 'TurbineType'),
  57. 'recording_time_s': ('ConfigurationSettings', 'RecordingTime_s'),
  58. 'monitoring_cycle_min': ('ConfigurationSettings', 'MonitoringCycle-Minutes'),
  59. 'sensor_name': ('Sensor', 'SensorName'),
  60. 'sensor_type': ('Sensor', 'SensorType'),
  61. 'sensor_addr': ('Sensor', 'SensorAddress'),
  62. 'sensor_sn': ('Sensor', 'Sensor_Sn'),
  63. 'sensitivity_mvpg': ('Sensor', 'Sensitivity_mVpG'),
  64. 'sensor_unit': ('Sensor', 'SensorUnit'),
  65. 'hw_serial': ('Sensor', 'HwSerial'),
  66. 'meas_type': ('Measurement', 'MeasurementType'),
  67. 'meas_name': ('Measurement', 'MeasurementName'),
  68. 'meas_source': ('Measurement', 'Source'),
  69. 'category': ('Measurement', 'Category'),
  70. 'meas_key': ('Measurement', 'MeasurementKey'),
  71. 'trigger_time': ('Measurement', 'Trigger_Time'),
  72. 'trigger_ms': ('Measurement', 'Trigger_Time_ms'),
  73. 'duration_s': ('Measurement', 'Measurement_Duration-s'),
  74. 'rpm': ('Measurement', 'RPM'),
  75. 'nominal_freq_hz': ('Measurement', 'NominalFrequency-Hz'),
  76. 'min_freq_hz': ('Measurement', 'MinimumFrequency-Hz'),
  77. 'max_freq_hz': ('Measurement', 'MaximumFrequency-Hz'),
  78. 'lower_freq_hz': ('Measurement', 'LowerFrequency-Hz'),
  79. 'upper_freq_hz': ('Measurement', 'UpperFrequency-Hz'),
  80. 'overload': ('Measurement', 'Overload'),
  81. 'n_averages': ('Measurement', 'NumberOfAverages'),
  82. 'bandwidth_hz': ('Measurement', 'Bandwidth-Hz'),
  83. 'lines': ('Measurement', 'Lines'),
  84. 'integration': ('Measurement', 'Integration'),
  85. 'condition_name': ('Measurement', 'Conditions', 'ConditionsName'),
  86. 'condition_key': ('Measurement', 'ConditionKey'),
  87. 'alarm_type': ('Measurement', 'AlarmType'),
  88. 'red_alarm': ('Measurement', 'RedAlarm'),
  89. 'yellow_alarm': ('Measurement', 'YellowAlarm'),
  90. 'blue_alarm': ('Measurement', 'BlueAlarm'),
  91. 'red_hys': ('Measurement', 'RedHysAlarm'),
  92. 'yellow_hys': ('Measurement', 'YellowHysAlarm'),
  93. 'trend_hys': ('Measurement', 'TrendHysAlarm'),
  94. 'fault_freq': ('Measurement', 'FaultFreq'),
  95. }
  96. # DataSets 下的列。★层级别踩错 (2026-09-12 对拍逮到): `Size`/`Dimension`/`Values` 在
  97. # `Measurement.DataSets.DataSet` 里, 而 X/Y 轴四件在**上一层** `Measurement.DataSets` 里 ——
  98. # 一开始全按 DataSet 取, 结果 x_offset/x_delta/x_unit/y_unit 四列整列为空,
  99. # 而 "x_delta == 带宽/lines" 的内部一致性检查比例是 0.000 (本该 1.000), 就是这里露的马脚。
  100. DS_MAP = {
  101. 'ds_dim': 'Dimension',
  102. 'ds_size': 'Size',
  103. }
  104. AXIS_MAP = {
  105. 'x_offset': 'X-axisOffset',
  106. 'x_delta': 'X-axisDelta',
  107. 'x_unit': 'X-axisUnit',
  108. 'y_unit': 'Y-axisUnit',
  109. }
  110. TXT_COLS = ('file', 'serial', 'turbine', 'ts_key', 'config_name', 'config_time', 'turbine_type',
  111. 'sensor_name', 'sensor_type', 'sensor_addr', 'sensor_sn', 'sensor_unit', 'hw_serial',
  112. 'meas_type', 'meas_name', 'meas_source', 'category', 'meas_key', 'trigger_time',
  113. 'integration', 'condition_name', 'condition_key', 'alarm_type', 'red_hys', 'yellow_hys',
  114. 'trend_hys', 'x_unit', 'y_unit', 'parse_error')
  115. NUM_COLS = ('rec_i', 'recording_time_s', 'monitoring_cycle_min', 'sensitivity_mvpg', 'trigger_ms',
  116. 'duration_s', 'rpm', 'nominal_freq_hz', 'min_freq_hz', 'max_freq_hz', 'lower_freq_hz',
  117. 'upper_freq_hz', 'overload', 'n_averages', 'bandwidth_hz', 'lines', 'fault_freq',
  118. 'red_alarm', 'yellow_alarm', 'blue_alarm', 'ds_dim', 'ds_size', 'x_offset', 'x_delta',
  119. 'scalar_value')
  120. # 54 列的**列序**逐字抄自包内 `outputs/<场>/m5_cms_tcm/tcm_index.parquet` (振动线 v0.2.0 产物,
  121. # 330,308 行) —— 消费者按列名取数, 但"列序也要一致"是为了让新旧窗在人工对拍/并排打印时能直接比。
  122. # (文本列与数值列是交错的, 所以不能用 TXT_COLS + NUM_COLS 拼。)
  123. COLSPEC = ['file', 'serial', 'turbine', 'ts_key', 'rec_i', 'config_name', 'config_time', 'turbine_type',
  124. 'recording_time_s', 'monitoring_cycle_min', 'sensor_name', 'sensor_type', 'sensor_addr',
  125. 'sensor_sn', 'sensitivity_mvpg', 'sensor_unit', 'hw_serial', 'meas_type', 'meas_name',
  126. 'meas_source', 'category', 'meas_key', 'trigger_time', 'trigger_ms', 'duration_s', 'rpm',
  127. 'nominal_freq_hz', 'min_freq_hz', 'max_freq_hz', 'lower_freq_hz', 'upper_freq_hz', 'overload',
  128. 'n_averages', 'bandwidth_hz', 'lines', 'integration', 'condition_name', 'condition_key',
  129. 'alarm_type', 'red_alarm', 'yellow_alarm', 'blue_alarm', 'red_hys', 'yellow_hys', 'trend_hys',
  130. 'fault_freq', 'ds_dim', 'ds_size', 'x_offset', 'x_delta', 'x_unit', 'y_unit', 'scalar_value',
  131. 'parse_error']
  132. def _dig(d, path):
  133. """按键路径取 (任一环缺失返回 None, 不抛)。"""
  134. for k in path:
  135. if not isinstance(d, dict):
  136. return None
  137. d = d.get(k)
  138. return d
  139. def _f(v):
  140. """→ float; 空/非数 → NaN (索引列必须能进 pandas float 列, 字符串混进去会让整列变 object)。"""
  141. if v is None or v == '':
  142. return math.nan
  143. try:
  144. return float(v)
  145. except (TypeError, ValueError):
  146. return math.nan
  147. def _s(v):
  148. return None if v is None else str(v)
  149. def _datasets(m):
  150. """Measurement.DataSets → (DataSets 层, DataSet 记录)。DataSet 可能是 dict 也可能是 list。"""
  151. dss = m.get('DataSets') or {}
  152. if not isinstance(dss, dict):
  153. return {}, {}
  154. ds = dss.get('DataSet')
  155. if isinstance(ds, list):
  156. ds = ds[0] if ds else {}
  157. return dss, (ds if isinstance(ds, dict) else {})
  158. def rows_of_doc(doc, relpath, ts_key, recs, rows):
  159. """一个 (文件, 时间戳) 下的一批记录 → 追加进 rows。"""
  160. for rec_i, item in enumerate(recs):
  161. rec = item.get('Record') if isinstance(item, dict) else None
  162. if not isinstance(rec, dict):
  163. rows.append(dict(file=relpath, ts_key=ts_key, rec_i=rec_i,
  164. parse_error='记录非 Record 结构: ' + str(type(item).__name__)))
  165. continue
  166. try:
  167. m = rec.get('Measurement') or {}
  168. dss, ds = _datasets(m)
  169. vals = ds.get('Values')
  170. size = _f(ds.get('Size'))
  171. row = dict(file=relpath, ts_key=ts_key, rec_i=rec_i)
  172. for col, path in REC_MAP.items():
  173. row[col] = _dig(rec, path)
  174. for col, key in DS_MAP.items():
  175. row[col] = ds.get(key)
  176. for col, key in AXIS_MAP.items(): # X/Y 轴在 DataSets 层, 不在 DataSet 里
  177. row[col] = dss.get(key)
  178. # 标量 (Size==1): Values 字符串本身就是标量值; 谱/波形 (Size>1) 的 Values 是长串, 值不进索引
  179. row['scalar_value'] = _f(vals) if (size == 1 and isinstance(vals, (str, int, float))) else None
  180. row['parse_error'] = None
  181. rows.append(row)
  182. except Exception as exc: # 单记录异常不许断整窗
  183. rows.append(dict(file=relpath, ts_key=ts_key, rec_i=rec_i,
  184. parse_error=f'{type(exc).__name__}: {exc}'))
  185. def rows_of_file(raw: bytes, relpath: str, rows: list):
  186. doc = json.loads(raw.decode('utf-8', 'replace'))
  187. body = ((doc.get('body') or {}).get('body')) or {}
  188. if not isinstance(body, dict):
  189. rows.append(dict(file=relpath, parse_error='body.body 非 dict'))
  190. return
  191. for ts_key, recs in body.items():
  192. if isinstance(recs, list):
  193. rows_of_doc(doc, relpath, ts_key, recs, rows)
  194. else:
  195. rows.append(dict(file=relpath, ts_key=ts_key, parse_error='记录集非 list'))
  196. def relpath_of(name: str) -> str:
  197. """包内成员名 → `file` 列。去掉 `measurement/` 这一层 (包本身的组织层, 不是数据层)。"""
  198. n = name.replace('\\', '/')
  199. if n.startswith('measurement/'):
  200. n = n[len('measurement/'):]
  201. return n
  202. def discover(root: pathlib.Path, turbines, limit):
  203. """→ [(成员名/相对路径, 完整路径|None, zip 路径|None)]。目录与 zip 都支持。"""
  204. out = []
  205. if root.is_file() and root.suffix.lower() == '.zip':
  206. with zipfile.ZipFile(root) as zf:
  207. for i in zf.infolist():
  208. if i.is_dir() or not i.filename.endswith('_decode.json'):
  209. continue
  210. out.append((i.filename, None, str(root)))
  211. else:
  212. for p in sorted(root.rglob('*_decode.json')):
  213. out.append((str(p.relative_to(root)).replace('\\', '/'), str(p), None))
  214. if turbines:
  215. want = {t.upper() for t in turbines}
  216. out = [x for x in out if any(t in x[0].upper() for t in want)]
  217. if limit:
  218. out = out[:limit]
  219. return out
  220. def flatten(items, rows):
  221. """逐文件解析 (目录: 直接读; zip: 按需打开一次)。"""
  222. cur_zip = None
  223. cur_path = None
  224. try:
  225. for name, full, zpath in items:
  226. rel = relpath_of(name)
  227. try:
  228. if zpath:
  229. if cur_zip is None or cur_path != zpath:
  230. if cur_zip:
  231. cur_zip.close()
  232. cur_zip = zipfile.ZipFile(zpath)
  233. cur_path = zpath
  234. raw = cur_zip.read(name)
  235. else:
  236. raw = pathlib.Path(full).read_bytes()
  237. rows_of_file(raw, rel, rows)
  238. except Exception as exc:
  239. rows.append(dict(file=rel, parse_error=f'文件级失败 {type(exc).__name__}: {exc}'))
  240. finally:
  241. if cur_zip:
  242. cur_zip.close()
  243. def worker(payload):
  244. """子进程: 解析分片 → 写分片 parquet → 回 (分片路径, 行数, 文件数, 失败数)。"""
  245. part_path, items = payload
  246. import pandas as pd
  247. rows = []
  248. flatten(items, rows)
  249. df = pd.DataFrame(rows)
  250. for c in COLSPEC:
  251. if c not in df.columns:
  252. df[c] = None
  253. df = df[list(COLSPEC)]
  254. for c in NUM_COLS:
  255. # 与 shipped 一致的 float64 (整列同型; 有 NaN 的列本就会升为 float, 没有 NaN 的列不升则出现 int64)
  256. df[c] = pd.to_numeric(df[c], errors='coerce').astype('float64')
  257. for c in TXT_COLS:
  258. df[c] = df[c].astype('str') # pandas 3 的 'str' dtype (shipped 同款)
  259. bad = int(df['parse_error'].notna().sum()) if 'parse_error' in df else 0
  260. df.to_parquet(part_path, index=False)
  261. return part_path, len(df), len(items), bad
  262. def main() -> int:
  263. ap = argparse.ArgumentParser()
  264. ap.add_argument('--root', required=True, help='含 *_decode.json 的目录, 或导出 zip')
  265. ap.add_argument('--out', required=True, help='index.parquet 输出路径')
  266. ap.add_argument('--jobs', type=int, default=min(8, (os.cpu_count() or 4)),
  267. help='并行进程数 (按机组切片; 默认 %(default)s)')
  268. ap.add_argument('--turbines', default=None, help='只处理这些机组 (逗号分隔, 如 WTG01,WTG02)')
  269. ap.add_argument('--limit', type=int, default=0, help='只处理前 N 个文件 (冒烟; 默认 0=全部)')
  270. ap.add_argument('--dry-run', action='store_true', help='只报将处理多少文件, 不解析')
  271. a = ap.parse_args()
  272. root = pathlib.Path(a.root)
  273. if not root.exists():
  274. raise SystemExit(f'输入不存在: {root}')
  275. out = pathlib.Path(a.out)
  276. turbines = [x.strip() for x in a.turbines.split(',')] if a.turbines else None
  277. t0 = time.time()
  278. items = discover(root, turbines, a.limit)
  279. if not items:
  280. raise SystemExit(f'没找到 *_decode.json: {root}')
  281. size = sum((pathlib.Path(f).stat().st_size if f else 0) for _, f, _ in items)
  282. print(f'输入: {root}')
  283. print(f'发现 {len(items)} 个 *_decode.json 目录内合计 {size / 1073741824:.2f} GB(仅目录模式可量)')
  284. if a.dry_run:
  285. print('(dry-run, 未解析)')
  286. return 0
  287. import pandas as pd
  288. import concurrent.futures as cf
  289. # 切片: 优先按机组 (一个子进程只碰自己那几台 → 负载均匀且写入互不干扰)
  290. groups = {}
  291. for it in items:
  292. key = next((seg for seg in it[0].replace('\\', '/').split('/') if seg.upper().startswith('WTG')), '_other')
  293. groups.setdefault(key, []).append(it)
  294. chunks = []
  295. n = max(1, a.jobs)
  296. per = math.ceil(len(items) / n)
  297. cur = []
  298. for k in sorted(groups):
  299. cur.extend(groups[k])
  300. if len(cur) >= per:
  301. chunks.append(cur)
  302. cur = []
  303. if cur:
  304. chunks.append(cur)
  305. parts_dir = out.parent / (out.stem + '.parts')
  306. if parts_dir.exists():
  307. shutil.rmtree(parts_dir)
  308. parts_dir.mkdir(parents=True, exist_ok=True)
  309. payloads = [(str(parts_dir / f'p{i:02d}.parquet'), ch) for i, ch in enumerate(chunks)]
  310. print(f'并行 {min(n, len(payloads))} 进程 × {len(payloads)} 片 → {out}')
  311. done_files = 0
  312. total_rows = 0
  313. total_bad = 0
  314. part_files = []
  315. with cf.ProcessPoolExecutor(max_workers=min(n, len(payloads))) as ex:
  316. for part, nrow, nfile, nbad in ex.map(worker, payloads):
  317. part_files.append(part)
  318. done_files += nfile
  319. total_rows += nrow
  320. total_bad += nbad
  321. print(f' [{done_files}/{len(items)} 文件] 累计 {total_rows} 行, 异常 {total_bad} '
  322. f'({time.time() - t0:.0f}s, {part})', flush=True)
  323. df = pd.concat([pd.read_parquet(p) for p in sorted(part_files)], ignore_index=True)
  324. out.parent.mkdir(parents=True, exist_ok=True)
  325. df.to_parquet(out, index=False)
  326. shutil.rmtree(parts_dir, ignore_errors=True)
  327. print(f'\n完成: {len(df)} 行 × {df.shape[1]} 列 → {out}')
  328. print(f' 机组 {df.turbine.nunique()} 台; 传感器 {df.sensor_name.nunique()} 种; 测量名 {df.meas_name.nunique()} 种')
  329. if 'trigger_time' in df:
  330. tt = pd.to_datetime(df.trigger_time, errors='coerce')
  331. print(f' 时间窗 {tt.min()} → {tt.max()}')
  332. print(f' 标量行(Size=1) {int((df.ds_size == 1).sum())}; 谱/波形行(Size>1) {int((df.ds_size > 1).sum())}; '
  333. f'解析异常 {int(df.parse_error.notna().sum())}')
  334. if total_bad:
  335. print(f' ⚠ 有 {total_bad} 条记录带 parse_error (见该列), 未静默丢弃')
  336. print(f' 用时 {time.time() - t0:.0f}s')
  337. return 0
  338. if __name__ == '__main__':
  339. sys.exit(main())