# -*- coding: utf-8 -*- r"""SCADA 原始数据接入层:**CSV 与 MDB 两种形态都认**(用户令 2026-09-19)。 ## 为什么要有这一层 现场交来的 SCADA 收资有两种形态,此前只有第一种能被摄入: ``` data/raw/<场站>/scada_10min/WTG01.csv … WTG38.csv ← 已导出 CSV(38 件 / 14.4 GB) data/raw/<场站>/scada_mdb/<年>年/<月>月/<年>-<月>-<类>.mdb ← 现场年度归档的 Access 库(25年/26年) ``` `src/windscada/data.py::load_10min()` 原来把路径**写死**成 `/<台>.csv`。用户令之后取数统一走本层: **CSV 在位就用 CSV**(既有 10 个构建器一直在用的形态,行数已核过),**该台没有 CSV 才回落 MDB**, 并把"这次是从哪种形态取的"记在返回表的 `attrs['scada_source']` 里,供产物台账/页面说明来路 —— 不静默换源。 ## MDB 那一侧的两个硬事实(决定了本层怎么写) 1. **单表 ≤255 列 ⇒ 库内按 250 列拆表**(`t1/t2/…`)。新库每行前两列是 `rid`(该台该月内的行号)与 `turbine`(机组标识)⇒ 取一台的多列数据要**跨分片表按 (rid, turbine) 拼**。 2. **2026-09-19 之前建的老库缺 `turbine` 列**(那时按位置堆行)。10min 老库还能救:源件自带 `WTG` 列, 本层用它兜底;**1min 老库救不了**(源件没有任何机组标识列,且 38 件表头实测 16 种 ⇒ 位置堆法还会串列)。 处置:重跑 `scripts/csv_to_mdb.py`(现已写 `turbine` 列 + 按列名对齐),见 `docs/现场收资接入_v0.1.md` §8。 读 MDB 走 ACE(`Microsoft.ACE.OLEDB.12.0`)+ PowerShell —— 与写库同一套(本机无 pyodbc/mdbtools; 实测可用路线见 `docs/现场收资接入_v0.1.md` §7)。日期字段在 PowerShell 里显式格式化成 `yyyy-MM-dd HH:mm:ss`(否则会按系统区域设置输出,下游时区/日序解析会踩雷)。 ## 用法 from src.windscada import scada_source as SS SS.sources('10min', cfg) # 有哪些形态、各多少件(如实,不猜) SS.turbines('10min', cfg) # 可用的机组名(CSV stem ∪ MDB 里的 turbine) SS.load('10min', 'WTG01', columns=[...], cfg=cfg) # → DataFrame(列名原样,ts 未解析) SS.source_of('10min', 'WTG01', cfg) # → 'csv' / 'mdb' 纪律: 两种形态都取不到时**响亮 raise**(不返回空表 —— 空表会被下游当成"这台没数据")。 """ from __future__ import annotations import functools import pathlib import re import subprocess 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" $c = New-Object -ComObject ADODB.Connection $c.Open("Provider=Microsoft.ACE.OLEDB.12.0;Data Source=@@MDB@@;") $turb = "@@TURB@@" $want = @(@@WANT@@) # 表名清单: 用 ADO 的 OpenSchema(20=adSchemaTables)。★`$c.GetSchema("Tables")` 会报 # "Arguments are of the wrong type..."(那是 .NET 的方法名), 实测 2026-09-19 踩过。 $tabs = @() $rs0 = $c.OpenSchema(20) while (-not $rs0.EOF) { if ($rs0.Fields.Item("TABLE_TYPE").Value -eq "TABLE") { $tabs += $rs0.Fields.Item("TABLE_NAME").Value } $rs0.MoveNext() } $rs0.Close() foreach ($t in $tabs) { $rs = $c.Execute("SELECT TOP 1 * FROM [$t]") $names = @() foreach ($f in $rs.Fields) { $names += $f.Name } $rs.Close() # 机组列: 先试 turbine(新库), 再试 WTG(老 10min 库兜底); 两个都没有 ⇒ 这库按台取不到数 $tcol = $null foreach ($p in @("turbine", "WTG")) { try { $null = $c.Execute("SELECT TOP 1 [$p] FROM [$t]"); $tcol = $p; break } catch { } } $turbines = @() if ($tcol) { $rs1 = $c.Execute("SELECT DISTINCT [$tcol] FROM [$t] WHERE [$tcol] IS NOT NULL") while (-not $rs1.EOF) { $turbines += ("" + $rs1.Fields.Item(0).Value); $rs1.MoveNext() } $rs1.Close() } Write-Output ("#SCHEMA " + $t + "|" + ($names -join ",") + "|" + ($turbines -join ";")) if (-not $tcol -or $turb -eq "") { continue } $keep = @() foreach ($w in $want) { if (($names -contains $w) -and $w -ne "rid" -and $w -ne $tcol) { $keep += $w } } if ($want.Count -eq 0) { # 不给白名单 = 全列(与 CSV 形态 columns=None 的语义一致) foreach ($w in $names) { if ($w -ne "rid" -and $w -ne $tcol) { $keep += $w } } } $sel = @("rid", $tcol) + $keep $collist = ($sel | ForEach-Object { "[$_]" }) -join "," Write-Output ("#TABLE " + $t + "|" + ($sel -join ",")) $rs2 = $c.Execute("SELECT $collist FROM [$t] WHERE [$tcol] = '$turb'") while (-not $rs2.EOF) { $vals = @() for ($i = 0; $i -lt $rs2.Fields.Count; $i++) { $v = $rs2.Fields.Item($i).Value if ($v -is [datetime]) { $vals += $v.ToString("yyyy-MM-dd HH:mm:ss") } elseif ($null -eq $v) { $vals += "" } else { $vals += ("" + $v) } } Write-Output ($vals -join "`t") $rs2.MoveNext() } $rs2.Close() } $c.Close() ''' def dirs_of(cls: str, cfg=None) -> tuple[pathlib.Path, pathlib.Path]: """→ (CSV 目录, MDB 根)。两者都不必存在。""" cfg = _cfg(cfg) key, sub = CLS[cls] csv_dir = pathlib.Path(str(cfg.get(key) or (pathlib.Path(str(cfg['raw_station_dir'])) / sub))) mdb_root = pathlib.Path(str(cfg['raw_station_dir'])) / 'scada_mdb' return csv_dir, mdb_root def _cfg(cfg=None): if cfg is not None: return cfg from src.windscada.config import farm # config 属后端模块(P4),暂留原路径 return farm() def sources(cls: str = '10min', cfg=None) -> dict: """盘上有哪些形态(件数)—— 供页面/自检如实说明。""" csv_dir, mdb_root = dirs_of(cls, cfg) csvs = sorted(csv_dir.glob('*.csv')) if csv_dir.is_dir() else [] mdbs = sorted(mdb_root.rglob(f'*-{cls}.mdb')) if mdb_root.is_dir() else [] zips = sorted(mdb_root.rglob(f'*-{cls}.zip')) if mdb_root.is_dir() else [] return dict(csv_dir=csv_dir, mdb_root=mdb_root, n_csv=len(csvs), n_mdb=len(mdbs), n_zip=len(zips), csvs=csvs, mdbs=mdbs, zips=zips, kind=('csv' if csvs else ('mdb' if (mdbs or zips) else None))) @functools.lru_cache(maxsize=8) def schema_and_turbines(mdb: str) -> tuple: """一个 MDB 的 (分片表 → 列名, 分片表 → 该列里的机组值)。老库缺 turbine 列时退 WTG 列。""" lines = _ps(_PS.replace('@@MDB@@', str(mdb)).replace('@@TURB@@', '').replace('@@WANT@@', '')) cols, turbs = {}, {} for ln in lines: if ln.startswith('#SCHEMA '): body = ln[8:] t, rest = body.split('|', 1) names, tv = rest.split('|', 1) cols[t] = [x for x in names.split(',') if x] turbs[t] = tuple(sorted({x for x in tv.split(';') if x})) return cols, turbs def _ps(script: str) -> list[str]: """跑一段 PowerShell(ACE 只能经 COM 用)→ stdout 行。 ★退出码不可尽信(2026-09-19 实测): ACE 在**进程退出时清理 COM** 偶发 `0xC0000005`(访问违例), 此时 stderr 为空、stdout 完整(脚本里每一行都打出来了)。原先 `rc != 0 ⇒ raise` 会把这种"数据都拿到了" 的情况报成读库失败。判据改成**看输出**: 有输出且退出码是那个崩溃码 ⇒ 收下并提示一次; 否则才 raise。 """ p = pathlib.Path(tempfile.gettempdir()) / '_scada_source.ps1' p.write_text(script + '\nexit 0\n', encoding='utf-8-sig') r = subprocess.run(['powershell', '-NoProfile', '-ExecutionPolicy', 'Bypass', '-File', str(p)], capture_output=True, text=True, errors='replace', timeout=3600) out = [x for x in (r.stdout or '').splitlines() if x.strip()] crash = r.returncode in (-1073741819, 3221225477) if r.returncode != 0 and not (crash and out): raise RuntimeError((r.stderr or r.stdout or '').strip()[:400] or 'PowerShell 读 MDB 失败') if crash and out and 'ace_exit_crash' not in _WARNED: _WARNED.add('ace_exit_crash') print(' [i] ACE 退出时 COM 清理崩溃 (0xC0000005) —— 输出完整, 已按成功处理', flush=True) return out _WARNED: set = set() def mdb_turbines(cls: str = '10min', cfg=None, limit: int = 3) -> set: """MDB 侧的可用机组名(取前几个库的并集 —— 逐库扫全量太慢,够用)。""" s = sources(cls, cfg) got = set() for m in s['mdbs'][:limit]: _cols, turbs = schema_and_turbines(str(m)) for t in turbs.values(): got |= set(t) return got def turbines(cls: str = '10min', cfg=None) -> list[str]: """可用机组名 = CSV 文件名 stem ∪ MDB 里的 turbine 值。""" s = sources(cls, cfg) got = {f.stem for f in s['csvs']} | mdb_turbines(cls, cfg) return sorted(got) def source_of(cls: str, turbine: str, cfg=None) -> str | None: """这一台这次会从哪种形态取(csv 优先)。""" csv_dir, _ = dirs_of(cls, cfg) if (csv_dir / f'{turbine}.csv').is_file(): return 'csv' 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 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) 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 ({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' return d def read_mdb(cls: str, turbine: str, columns=None, cfg=None, month: str | None = None): """从月库里取一台:逐分片表按 (rid, turbine) 拼列。""" import pandas as pd from app_ETL.app_ETL_guanlan.mdb_names import safe_name s = sources(cls, cfg) mdbs = [m for m in s['mdbs'] if not month or m.stem.startswith(month)] if not mdbs: raise RuntimeError(f'{cls} 取数失败: scada_mdb 下没有 {month or "任何"} 的 .mdb') # 契约列名 → 库内列名(Access 列名不能含 `.` `!` `[]` 等, 建库时消毒过; 两边同一套规则) cand = {c: safe_name(c) for c in (columns or [])} want_stored = sorted(set(cand.values())) frames, missing_col = [], None for m in mdbs: want = want_stored script = (_PS.replace('@@MDB@@', str(m)).replace('@@TURB@@', str(turbine)) .replace('@@WANT@@', ','.join('"' + w + '"' for w in want))) cur_t, cur_cols, buf, seen_cols = None, [], [], set() def _flush(): if cur_t and buf: frames.append(_frame(cur_cols, buf)) for ln in _ps(script): if ln.startswith('#SCHEMA '): continue if ln.startswith('#TABLE '): _flush() body = ln[7:] cur_t, cols = body.split('|', 1) if not re.fullmatch(r't\d+', cur_t): # Access 自建的杂表(如"名称自动更正保存失败") cur_t, cur_cols, buf = None, [], [] continue cur_cols, buf = cols.split(','), [] seen_cols |= set(cur_cols) else: buf.append(ln) _flush() if columns and not any(c in seen_cols for c in columns): missing_col = [c for c in columns if c not in seen_cols][:3] # 列位错位检测: 用 **#SCHEMA**(每个分片表的列清单, 与本次取了哪些列无关) dup = set() for m2 in mdbs: sc = {} for ln in _ps(_PS.replace('@@MDB@@', str(m2)).replace('@@TURB@@', '').replace('@@WANT@@', '')): if ln.startswith('#SCHEMA '): body = ln[8:] t2, rest = body.split('|', 1) if re.fullmatch(r't\d+', t2): sc[t2] = set(x for x in rest.split('|', 1)[0].split(',') if x) ks = list(sc) for i in range(len(ks)): for j in range(i + 1, len(ks)): dup |= (sc[ks[i]] & sc[ks[j]]) - {'rid', TURBINE_COL} if dup: break if not frames: raise RuntimeError( f'{cls} 取数失败: {turbine} 在 {[m.name for m in mdbs]} 里取不到行 —— 老库可能没有 turbine 列' f'(源件也没有 WTG 列)⇒ 重跑 scripts/csv_to_mdb.py 建新库(现已写 turbine 列)') frames = [f.rename(columns={LEGACY_TURBINE_COL: TURBINE_COL}) if (LEGACY_TURBINE_COL in f.columns and TURBINE_COL not in f.columns) else f for f in frames] # 库内列名 → 契约列名(撞名时建库侧写过 `_1` 后缀, 这里按前缀认回来) if cand: for f in frames: back, used = {}, set() for orig, st in cand.items(): if st in f.columns and st not in used: back[st], _ = orig, used.add(st) continue stem = st[:56] for c2 in f.columns: if c2 not in used and c2.startswith(stem) and c2[len(stem):].startswith('_'): back[c2], _ = orig, used.add(c2) break f.rename(columns=back, inplace=True) # ★2026-09-19 逮到的**老库列位错位**: 2026-09-19 之前建的 10min 库, 分片表表头与数据行差一列 # (旧 build_chunks 表头写 names[lo:hi+1] 而行写 vals[lo:hi]) ⇒ 第 250 列之后的**列名整体错位一格**, # 且相邻分片表会共用同一个列名(实测 t1 末列 == t2 首列 == flg_wtc_ScReToOp_endvalue)。这种库读出来的 # 是"名字对不上数"的数据 —— 比缺数据危险, 故**响亮拒绝**, 让人重建而不是将就着算。 if dup: raise RuntimeError( f'{cls} 取数失败: {[m.name for m in mdbs]} 是**列位错位的老库**(分片表共用列名 {sorted(dup)[:3]}…) ' f'—— 2026-09-19 前的建库脚本表头与数据行差一列, 第 250 列之后的列名整体错位。' f'处置: python scripts/csv_to_mdb.py --month --class {cls} 重建该月库; ' f'或放 CSV 到 data/raw/<场站>/{CLS[cls][1]}/ 直接用 CSV 形态。') d = frames[0] for f in frames[1:]: d = d.merge(f, on=['rid', TURBINE_COL], how='outer') # 行序回到**源件顺序**: rid 在库里是文本 ⇒ 直接合并得到的是 '1','10','100' 这种字典序 # (2026-09-19 实逮: 与 CSV 逐值对拍对不上, 数据没错、是行序被按字符串排了) d['rid'] = pd.to_numeric(d['rid'], errors='coerce') d = d.sort_values([TURBINE_COL, 'rid'], kind='stable').reset_index(drop=True) if missing_col: d.attrs['mdb_missing'] = missing_col return d def _frame(cols: list[str], rows: list[str]): """PowerShell 落的制表符分隔文本 → DataFrame。 数值列转成数值: MDB 里存的是**数值类型**, ACE 读出来 `100.0` 会变成 `100` —— 若原样留成字符串, 下游按 CSV 口径写的老构建器会拿到 object 列 (`+`/`mean` 全崩)。这里逐列试探: 非空值全能解析成数字 就转数值, 否则留字符串 (台号/状态字这类)。 """ import pandas as pd ncol = len(cols) recs = [] for ln in rows: vals = ln.split('\t') recs.append((vals + [''] * ncol)[:ncol]) d = pd.DataFrame(recs, columns=cols) for c in d.columns: if c in ('rid', TURBINE_COL, LEGACY_TURBINE_COL): continue s = d[c].replace('', None) num = pd.to_numeric(s, errors='coerce') if num.notna().sum() == s.notna().sum() and s.notna().any(): d[c] = num return d