|
|
@@ -47,6 +47,14 @@ import tempfile
|
|
|
CLS = {'10min': ('src_10min', 'scada_10min'), '1min': ('src_1min', 'scada_1min')}
|
|
|
TURBINE_COL = 'turbine' # 新库的机组列(csv_to_mdb.py 2026-09-19 起写)
|
|
|
LEGACY_TURBINE_COL = 'WTG' # 老库兜底:10min 源件自带 WTG 列
|
|
|
+# ★同台多件 (用户令 2026-09-20「真正把新月份的 10min 件吃进来」): 现场后来补的导出件命名与主件不同
|
|
|
+# (实测 `WTG01-B2.csv` 856 列 / 2026-08 数据),而原实现只认 `<台号>.csv` ⇒ 这批数据**从未被读过**,
|
|
|
+# 页面时间窗自然停在旧月份。现在本台取数 = 主件 + 同台补充件,按时间戳合并去重(主件优先)。
|
|
|
+SUPPLEMENT_GLOBS = ('{t}-*.csv', '{t}_*.csv')
|
|
|
+# 新旧导出的列名差: 新件把"类型前缀族"去掉了(`din_wtc_HydLevel_timeon` → `wtc_HydLevel_timeon`)。
|
|
|
+# 实测 WTG01 老件 598 列里 597 列能这样对上(唯一对不上的是 `WTG`,新件用 `StationId`)。
|
|
|
+PREFIXES = ('din_', 'dot_', 'prs_', 'tmp_', 'grd_', 'tur_', 'cnt_', 'flg_', 'int_', 'evt_')
|
|
|
+TIME_CANDIDATES = ('TimeStamp', 'timestamp', 'Time', 'time', 'time_stamp', 'occur_time')
|
|
|
|
|
|
_PS = r'''
|
|
|
$ErrorActionPreference = "Continue"
|
|
|
@@ -199,30 +207,117 @@ def source_of(cls: str, turbine: str, cfg=None) -> str | None:
|
|
|
return 'mdb' if turbine in mdb_turbines(cls, cfg, limit=1) else None
|
|
|
|
|
|
|
|
|
+def files_of(cls: str, turbine: str, cfg=None) -> list:
|
|
|
+ """本台这一次要读的 CSV 件: **主件 `<台号>.csv` 在前**, 后面是同台补充件(`<台号>-*.csv` 等)。
|
|
|
+
|
|
|
+ 用户令 2026-09-20: 现场补的件命名与主件不同(实测 `WTG01-B2.csv`),只认 `<台号>.csv` 会让
|
|
|
+ "新来的月份"静默漏读。这里如实列出这次会读的件(顺序=优先级),供取数/台账/页面说明来路。
|
|
|
+ """
|
|
|
+ csv_dir, _ = dirs_of(cls, cfg)
|
|
|
+ out = []
|
|
|
+ pri = csv_dir / f'{turbine}.csv'
|
|
|
+ if pri.is_file():
|
|
|
+ out.append(pri)
|
|
|
+ for pat in SUPPLEMENT_GLOBS:
|
|
|
+ out += sorted(csv_dir.glob(pat.format(t=turbine)))
|
|
|
+ seen, res = set(), []
|
|
|
+ for p in out:
|
|
|
+ if p.name not in seen:
|
|
|
+ seen.add(p.name)
|
|
|
+ res.append(p)
|
|
|
+ return res
|
|
|
+
|
|
|
+
|
|
|
+def _align_columns(df, want):
|
|
|
+ """按列名对齐新旧导出: 想要的列不在、但"去掉类型前缀"后在同表里 ⇒ 认回来。
|
|
|
+
|
|
|
+ 返回 (df, 认回的列数)。实测: 老 598 列 → 新 856 列, 597/598 靠这条规则对上。
|
|
|
+ """
|
|
|
+ if not want or df is None or not len(df.columns):
|
|
|
+ return df, 0
|
|
|
+ have = set(df.columns)
|
|
|
+ back = {}
|
|
|
+ for c in want:
|
|
|
+ if c in have:
|
|
|
+ continue
|
|
|
+ for pre in PREFIXES:
|
|
|
+ if c.startswith(pre) and c[len(pre):] in have:
|
|
|
+ back[c[len(pre):]] = c
|
|
|
+ break
|
|
|
+ return (df.rename(columns=back) if back else df), len(back)
|
|
|
+
|
|
|
+
|
|
|
+def _time_col(frames, first_cols):
|
|
|
+ """时间列: 先按已知名找(合并后的表列名可能来自任一形态),再退回主件的首列。"""
|
|
|
+ cols = set(first_cols)
|
|
|
+ for f in frames:
|
|
|
+ cols |= set(f.columns)
|
|
|
+ for c in TIME_CANDIDATES:
|
|
|
+ if c in cols:
|
|
|
+ return c
|
|
|
+ return first_cols[0] if first_cols else None
|
|
|
+
|
|
|
+
|
|
|
def load(cls: str, turbine: str, columns=None, cfg=None, month: str | None = None):
|
|
|
"""取一台的数据(**CSV 优先**;该台没有 CSV 时回落 MDB)。
|
|
|
|
|
|
columns = 白名单(两形态都按"实际有的列"取交集,缺列不报错 —— 与 `data.load_10min` 的既有口径一致)。
|
|
|
month = 'YYYY-MM',只取该月(MDB 形态天然是按月的;CSV 形态按首列前缀过滤)。
|
|
|
+ ★同台多件 (用户令 2026-09-20): `<台号>.csv` + `<台号>-*.csv` 一起读, 按列名对齐 + 按时间戳去重
|
|
|
+ (**主件优先**,重叠条数如实写在 `attrs['scada_overlap']` 里)。
|
|
|
"""
|
|
|
import pandas as pd
|
|
|
- csv_dir, _ = dirs_of(cls, cfg)
|
|
|
- fp = csv_dir / f'{turbine}.csv'
|
|
|
- if fp.is_file():
|
|
|
- if columns:
|
|
|
+ fps = files_of(cls, turbine, cfg)
|
|
|
+ if fps:
|
|
|
+ frames, per = [], {}
|
|
|
+ for i, fp in enumerate(fps):
|
|
|
have = set(pd.read_csv(fp, nrows=0).columns)
|
|
|
- cols = [c for c in columns if c in have]
|
|
|
- else:
|
|
|
- cols = None
|
|
|
- d = pd.read_csv(fp, usecols=cols, engine='pyarrow') if cols else pd.read_csv(fp, engine='pyarrow')
|
|
|
- if month:
|
|
|
- ts = d.columns[0]
|
|
|
+ take = None
|
|
|
+ if columns:
|
|
|
+ take = []
|
|
|
+ for c in columns:
|
|
|
+ if c in have:
|
|
|
+ take.append(c)
|
|
|
+ continue
|
|
|
+ for pre in PREFIXES: # 新导出把前缀去了 ⇒ 按去前缀名取列
|
|
|
+ if c.startswith(pre) and c[len(pre):] in have:
|
|
|
+ take.append(c[len(pre):])
|
|
|
+ break
|
|
|
+ d = pd.read_csv(fp, usecols=take, engine='pyarrow') if take else \
|
|
|
+ (pd.read_csv(fp, usecols=lambda c: c in have, engine='pyarrow') if columns else
|
|
|
+ pd.read_csv(fp, engine='pyarrow'))
|
|
|
+ d, n_back = _align_columns(d, columns or [])
|
|
|
+ first_cols = list(d.columns)
|
|
|
+ d['_scada_rank'] = i # i=0 是主件 ⇒ 去重时优先
|
|
|
+ frames.append(d)
|
|
|
+ per[fp.name] = len(d)
|
|
|
+ ts = _time_col(frames, first_cols)
|
|
|
+ d = pd.concat(frames, ignore_index=True, sort=False)
|
|
|
+ if month and ts:
|
|
|
d = d[d[ts].astype(str).str.startswith(month)]
|
|
|
+ overlap = 0
|
|
|
+ if len(frames) > 1 and ts:
|
|
|
+ before = len(d)
|
|
|
+ d = d.sort_values(['_scada_rank'], kind='stable').drop_duplicates(subset=[ts], keep='first')
|
|
|
+ overlap = before - len(d)
|
|
|
+ d = d.drop(columns=['_scada_rank'])
|
|
|
+ if ts in d.columns:
|
|
|
+ d = d.sort_values(ts, kind='stable').reset_index(drop=True)
|
|
|
+ else:
|
|
|
+ d = d.reset_index(drop=True)
|
|
|
+ if len(fps) > 1:
|
|
|
+ detail = '、'.join(f'{k} {v} 行' for k, v in per.items())
|
|
|
+ print(f' [i] {turbine}: 同台 {len(fps)} 件合并读入({detail})'
|
|
|
+ + (f';同时间戳重叠 {overlap} 行按"主件优先"去重' if overlap else ';无重叠'),
|
|
|
+ flush=True)
|
|
|
d.attrs['scada_source'] = 'csv'
|
|
|
+ d.attrs['scada_files'] = [f.name for f in fps]
|
|
|
+ d.attrs['scada_rows'] = per
|
|
|
+ d.attrs['scada_overlap'] = overlap
|
|
|
return d.reset_index(drop=True)
|
|
|
if source_of(cls, turbine, cfg) != 'mdb':
|
|
|
raise RuntimeError(
|
|
|
- f'{cls} 取数失败: {turbine} 既没有 CSV ({fp}),也没有含该台的 MDB '
|
|
|
+ f'{cls} 取数失败: {turbine} 既没有 CSV ({dirs_of(cls, cfg)[0] / f"{turbine}.csv"}),也没有含该台的 MDB '
|
|
|
f'({dirs_of(cls, cfg)[1]})。放原始件到 data/raw/<场站>/ 后重跑 scripts/rebuild_all.py。')
|
|
|
d = read_mdb(cls, turbine, columns, cfg, month)
|
|
|
d.attrs['scada_source'] = 'mdb'
|