|
|
@@ -0,0 +1,383 @@
|
|
|
+#!/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())
|