#!/usr/bin/env python3 # -*- coding: utf-8 -*- r"""现场收资接入:从场站交来的年度 zip(内层是每月每类 `.mdb`)提取并转成观澜 `data/raw/<场站>/`。 ## 场站交来的东西长什么样(2026-09-17 实测 `F:\temp\如东风场数据`) ``` 25年.zip (1,022 MB) / 26年.zip (961 MB) └ 26年/1月/2026-01-cnt.zip → 2026-01-cnt.mdb (52.7 MB) 2026-01-std.zip → 2026-01-std.mdb ( 2.3 MB) ← tblGrid: 10 分钟电网/机组量测 2026-01-tur.zip → 2026-01-tur.mdb (47.6 MB) 2026-01-scd.zip → 2026-01-scd.mdb ( 5.6 MB) … 每月 12 类:cnt din dot flg grd int prs scd std sum tmp tur ``` 即 **Access 数据库按月分类**(不是 CSV)。实测 `std.mdb` 里是单表 `tblGrid`: ``` TimeStamp | Station | WPSStatus | TimestampStation | CurrentL1..L3 | VoltageL1..L3 | ActivePower | ReActivePower | ActivePowerExport | ReActivePowerExport | … ``` `TimeStamp` 是 10 分钟节拍,`Station` 是机组号(如 91/92 —— 与 WTG 号的对应关系见 §映射)。 ## 读 MDB 的路线(离线、零 pip 依赖) Windows 上用 **.NET `System.Data.Odbc` + 系统自带的 `Microsoft Access Driver (*.mdb, *.accdb)`** (本机实测可用;驱动来自 Office/Access Runtime,纯服务器上可能没有 ⇒ 本器会**先探针、缺了就响亮报并给补救**)。 Linux 上用 `mdbtools` 的 `mdb-export`(本机没有)。两条路都不需要往包里塞新 pip 包。 ## 用法 python scripts/extract_site_data.py --src <年度zip所在目录> --list python scripts/extract_site_data.py --src <…> --extract [--year 26] [--month 1] [--out data/raw/如东/scada_mdb] python scripts/extract_site_data.py --src <…> --probe [--class std] python scripts/extract_site_data.py --src <…> --to-csv --class std [--year 26 --month 1] [--turbines WTG01,WTG02] 退出码: 0 成功 · 2 来源/参数不齐 · 4 平台无 MDB 读取能力(已给补救办法) · 5 转换/校验失败 """ from __future__ import annotations import argparse import hashlib import json import os import pathlib import shutil import subprocess import sys import tempfile import zipfile ROOT = pathlib.Path(__file__).resolve().parents[1] sys.path.insert(0, str(ROOT)) from src import paths as P # noqa: E402 WIN32 = sys.name == 'nt' if hasattr(sys, 'name') else False WIN32 = sys.platform.startswith('win') DRIVER = 'Microsoft Access Driver (*.mdb, *.accdb)' # 类目 → 观澜侧的落点(先按已知的上线:只有 std 与 sum 是 10min 量测;其余待现场口径确认) CLASS_TO_RAW = { 'std': 'scada_10min', # 实测: tblGrid = 10 分钟电网/机组量测(TimeStamp/Station/ActivePower/…) 'sum': 'scada_10min', # 待确认(体量最大, 疑为统计汇总) 'cnt': 'scada_10min', # 待确认 'tur': 'scada_10min', # 待确认(47 MB, 疑为机组状态) } # 现场交来的**已导出 CSV** zip → 观澜 raw 目录(2026-09-17 实测:与 data/raw 现有 38 件逐名对应) ZIP_TO_RAW = { '10分钟数据': 'scada_10min', # WTG01.csv … WTG38.csv(38 件 / 14.4 GB) '1分钟数据': 'scada_1min', # 01E.csv … 38B.csv(38 件 / 12.8 GB,机组代号命名) } def verify_raw(src: pathlib.Path, farm: str | None = None) -> int: """对账:现场 zip 里的 CSV ↔ `data/raw/<场>/<类>/` 现有文件(逐件比大小)。 这一步回答的是"**现有的输入数据到底是不是从这两个 zip 来的**" —— 不去解压十几 GB 就能验。 """ farm = farm or P.farm() zips = [z for z in find_year_zips(src) if z.stem in ZIP_TO_RAW] if not zips: print(f'[X] {src} 下没有 {"/".join(ZIP_TO_RAW)}.zip') return 2 print(f'== 现场 zip ↔ data/raw/{farm}/ 对账 ==') bad = 0 for z in zips: d = station_dir(farm) / ZIP_TO_RAW[z.stem] with zipfile.ZipFile(z) as zf: ents = [e for e in zf.infolist() if not e.is_dir()] same = miss = diff = 0 miss_names, diff_names = [], [] for e in ents: tgt = d / pathlib.Path(e.filename).name if not tgt.is_file(): miss += 1 miss_names.append(tgt.name) elif tgt.stat().st_size == e.file_size: same += 1 else: diff += 1 diff_names.append(f'{tgt.name}(zip {e.file_size:,} vs 盘 {tgt.stat().st_size:,})') bad += miss + diff print(f' {z.name} → {P.rel(d)}: zip {len(ents)} 件 · 逐字节同大小 {same} · 缺 {miss} · 大小不同 {diff}') for n in miss_names[:4]: print(f' 缺: {n}') for n in diff_names[:4]: print(f' 异: {n}') print(f' 结论: {"现有 raw 输入与现场 zip 逐件一致(可追溯)" if bad == 0 else str(bad) + " 件需要重新提取"} ' f'(缺件提取: --extract-csv --only-missing)') return 0 if bad == 0 else 5 def extract_csv(src: pathlib.Path, farm: str | None = None, only_missing: bool = False) -> int: """把"已导出 CSV"的现场 zip 解到 `data/raw/<场>/<类>/`(幂等:同名同大小即跳过)。""" farm = farm or P.farm() zips = [z for z in find_year_zips(src) if z.stem in ZIP_TO_RAW] if not zips: print(f'[X] {src} 下没有 {"/".join(ZIP_TO_RAW)}.zip') return 2 done = skipped = 0 for z in zips: d = station_dir(farm) / ZIP_TO_RAW[z.stem] d.mkdir(parents=True, exist_ok=True) print(f' {z.name} → {P.rel(d)}') with zipfile.ZipFile(z) as zf: for e in zf.infolist(): if e.is_dir(): continue tgt = d / pathlib.Path(e.filename).name if tgt.is_file() and tgt.stat().st_size == e.file_size: skipped += 1 continue if only_missing and tgt.is_file(): skipped += 1 continue print(f' + {tgt.name} {e.file_size / 1e6:.1f} MB') with zf.open(e) as fh, open(tgt, 'wb') as out: shutil.copyfileobj(fh, out, 1024 * 1024 * 8) done += 1 print(f' 完成: 新解 {done} 件 · 跳过(已同大小){skipped} 件') print(' 接下去: python scripts/raw_data_check.py(体检)→ python scripts/rebuild_all.py(重算)') return 0 def station_dir(farm: str | None = None) -> pathlib.Path: """`data/raw/<场站>/` —— ★场站**目录名**可与 farm 键不同(键 = rudong,目录 = 如东), 所以必须走唯一取用口 `raw_station_dir()`;第一版按 `P.RAW_ROOT / farm` 拼,于是"38 件全缺" (对着 data/raw/rudong 找, 而东西在 data/raw/如东)。""" from src.windscada.config import raw_station_dir return pathlib.Path(raw_station_dir(farm or P.farm())) def find_year_zips(src: pathlib.Path) -> list[pathlib.Path]: return sorted([p for p in src.glob('*.zip') if p.is_file()]) def list_structure(src: pathlib.Path) -> int: zips = find_year_zips(src) if not zips: print(f'[X] {src} 下没有年度 zip') return 2 print(f'== 现场收资结构 · {src} ==') for z in zips: with zipfile.ZipFile(z) as zf: names = zf.namelist() months = sorted({n.split('/')[1] for n in names if n.count('/') >= 1 and n.split('/')[1]}) classes = sorted({n.split('/')[-1].split('-')[-1].replace('.zip', '') for n in names if n.lower().endswith('.zip')}) print(f' {z.name} {z.stat().st_size / 1e6:.0f} MB · {len(names)} 条目') print(f' 月份 {len(months)} 个: {", ".join(months)}') print(f' 类目 {len(classes)} 个: {", ".join(classes)}') return 0 def extract(src: pathlib.Path, out: pathlib.Path, year: str | None, month: str | None) -> int: zips = [z for z in find_year_zips(src) if (not year or z.stem.startswith(year))] if not zips: print('[X] 没有匹配的年度 zip') return 2 out.mkdir(parents=True, exist_ok=True) man_p = out / '_manifest.json' man = json.loads(man_p.read_text(encoding='utf-8')) if man_p.is_file() else dict(zips=[], files=[]) done = 0 for z in zips: with zipfile.ZipFile(z) as zf: inner = [n for n in zf.namelist() if n.lower().endswith('.zip') and (not month or f'/{month}月/' in n)] print(f' {z.name}: 内层 {len(inner)} 个分类 zip' + (f'(限 {month} 月)' if month else '')) for n in inner: parts = n.split('/') # 26年/1月/2026-01-std.zip rel = pathlib.Path(*parts[1:]) if len(parts) > 1 else pathlib.Path(parts[0]) dst_dir = out / rel.parent dst_dir.mkdir(parents=True, exist_ok=True) with zf.open(n) as fh, zipfile.ZipFile(fh) as inner_zf: for m in inner_zf.namelist(): tgt = dst_dir / pathlib.Path(m).name if tgt.is_file() and tgt.stat().st_size == inner_zf.getinfo(m).file_size: continue tgt.write_bytes(inner_zf.read(m)) h = hashlib.sha256(tgt.read_bytes()).hexdigest()[:16] man['files'] = [x for x in man['files'] if x['path'] != str(tgt.relative_to(out))] man['files'].append(dict(path=str(tgt.relative_to(out)), bytes=tgt.stat().st_size, sha256_16=h, from_zip=f'{z.name}:{n}')) done += 1 print(f' + {tgt.relative_to(out)} {tgt.stat().st_size / 1e6:.1f} MB sha256:{h}') man['zips'] = sorted({*man.get('zips', []), *[z.name for z in zips]}) man_p.write_text(json.dumps(man, ensure_ascii=False, indent=1) + '\n', encoding='utf-8') print(f' 完成: 新解 {done} 个 MDB → {P.rel(out)}(清单 {P.rel(man_p)})') return 0 def have_mdb_route() -> tuple[bool, str]: if WIN32: ps = ('$d=(Get-OdbcDriver | Where-Object { $_.Name -like "*Access Driver (*.mdb*" }).Count;' 'if($d -gt 0){"ok"}else{"no"}') try: r = subprocess.run(['powershell', '-NoProfile', '-Command', ps], capture_output=True, text=True, errors='replace', timeout=120) ok = 'ok' in (r.stdout or '') return ok, 'Windows + .NET System.Data.Odbc + 系统 Access ODBC 驱动' if ok else '缺 Access ODBC 驱动' except Exception as e: return False, f'探测失败: {type(e).__name__}' ok = bool(shutil.which('mdb-export') or shutil.which('mdb-tables')) return ok, 'mdbtools (mdb-export)' if ok else '缺 mdbtools' def mdb_tables(mdb: pathlib.Path) -> list[str]: if WIN32: ps = (f'$cs="Driver={{{DRIVER}}};Dbq={mdb};";' '$cn=New-Object System.Data.Odbc.OdbcConnection($cs);$cn.Open();' '$cn.GetSchema("Tables") | Where-Object { $_.TABLE_TYPE -eq "TABLE" } | ' 'ForEach-Object { $_.TABLE_NAME };$cn.Close()') r = subprocess.run(['powershell', '-NoProfile', '-Command', ps], capture_output=True, text=True, errors='replace', timeout=300) if r.returncode != 0: raise RuntimeError((r.stderr or r.stdout or '').strip()[:300]) return [x.strip() for x in (r.stdout or '').splitlines() if x.strip()] r = subprocess.run(['mdb-tables', '-1', str(mdb)], capture_output=True, text=True, errors='replace', timeout=300) return [x.strip() for x in (r.stdout or '').splitlines() if x.strip()] def mdb_to_csv(mdb: pathlib.Path, table: str, out_csv: pathlib.Path, where: str = '') -> int: """MDB 单表 → CSV(UTF-8)。Windows 走 PowerShell/.NET Odbc;Linux 走 mdb-export。""" out_csv.parent.mkdir(parents=True, exist_ok=True) if WIN32: q = f'SELECT * FROM [{table}]' + (f' WHERE {where}' if where else '') ps = f'''$cs="Driver={{{DRIVER}}};Dbq={mdb};" $cn=New-Object System.Data.Odbc.OdbcConnection($cs);$cn.Open() $cmd=$cn.CreateCommand();$cmd.CommandText="{q}" $da=New-Object System.Data.Odbc.OdbcDataAdapter($cmd) $dt=New-Object System.Data.DataTable;[void]$da.Fill($dt) $dt | Export-Csv -Path "{out_csv}" -NoTypeInformation -Encoding UTF8 "rows=$($dt.Rows.Count)" $cn.Close()''' tf = pathlib.Path(tempfile.gettempdir()) / '_mdb_dump.ps1' tf.write_text(ps, encoding='utf-8-sig') r = subprocess.run(['powershell', '-NoProfile', '-ExecutionPolicy', 'Bypass', '-File', str(tf)], capture_output=True, text=True, errors='replace', timeout=3600) print(' ' + (r.stdout or '').strip().replace('\n', ' ')[:80]) if r.returncode != 0: print(' [X] ' + (r.stderr or r.stdout or '').strip()[:200]) return 5 return 0 cmd = ['mdb-export', str(mdb), table] + (['-q', where] if where else []) with open(out_csv, 'w', encoding='utf-8', newline='') as fh: r = subprocess.run(cmd, stdout=fh, stderr=subprocess.PIPE, text=True, errors='replace', timeout=3600) if r.returncode != 0: print(' [X] ' + (r.stderr or '')[:200]) return 5 return 0 def probe(out: pathlib.Path, cls: str | None) -> int: ok, route = have_mdb_route() print(f'== MDB 读取能力 ==\n {route}: {"可用" if ok else "不可用"}') if not ok: print(' 补救: Windows 装 Access Runtime(或 Office)后自带 Access ODBC 驱动; ' 'Linux 装 mdbtools(apt/yum install mdbtools); 或让现场把 MDB 导出成 CSV。') return 4 mdbs = sorted(out.rglob('*.mdb')) mdbs = [m for m in mdbs if (not cls or f'-{cls}.' in m.name)] if not mdbs: print(f' [i] {P.rel(out)} 下还没有 MDB —— 先跑 --extract') return 0 print(f' 已解出 MDB {len(mdbs)} 个(限类目 {cls or "全部"})') for m in mdbs[:8]: try: t = mdb_tables(m) print(f' {m.relative_to(out)}: 表 {t}') except Exception as e: print(f' {m.relative_to(out)}: [X] {type(e).__name__}: {str(e)[:120]}') return 0 def main() -> int: ap = argparse.ArgumentParser(description='现场收资(年度 zip → 每月每类 MDB)提取与接入') # ★2026-09-19 移植性扫描逮到: 原默认值是**开发机的绝对路径** `F:/temp/如东风场数据` # —— 换台机器上这个默认值毫无意义, 而且会让 `--src` 忘写时去翻一个不存在的目录。 # 改成: env `GUANLAN_SITE_SRC` 可配, 否则 None 并在缺参时响亮提示(不猜路径)。 ap.add_argument('--src', default=os.environ.get('GUANLAN_SITE_SRC') or None, help='年度 zip 所在目录(默认取 env GUANLAN_SITE_SRC; 没有则必填)') ap.add_argument('--out', default=str(station_dir() / 'scada_mdb'), help='MDB 落点(默认 data/raw/<场站>/scada_mdb)') ap.add_argument('--list', action='store_true', help='只列结构(年度 zip / 月份 / 类目)') ap.add_argument('--extract', action='store_true', help='解压出 MDB(幂等, 带清单)') ap.add_argument('--probe', action='store_true', help='探针: 读取能力 + 已解出 MDB 的表') ap.add_argument('--verify-raw', action='store_true', help='对账: 现场 CSV zip ↔ data/raw 现有文件(逐件比大小)') ap.add_argument('--extract-csv', action='store_true', help='把现场 CSV zip 解到 data/raw(幂等)') ap.add_argument('--only-missing', action='store_true', help='与 --extract-csv 合用: 只补缺件, 不覆盖已存在的') ap.add_argument('--to-csv', action='store_true', help='MDB → 观澜 data/raw 的 CSV(按类目映射)') ap.add_argument('--class', dest='cls', default=None, help='类目: cnt/sum/std/tur/…') ap.add_argument('--year', default=None, help='年度前缀(如 26)') ap.add_argument('--month', default=None, help='月份数字(如 1)') ap.add_argument('--table', default=None, help='指定表名(默认取该 MDB 的第一张表)') ap.add_argument('--turbines', default=None, help='只导这些机组(逗号分隔, 如 WTG01,WTG02)') a = ap.parse_args() if not a.src: print('[X] 需要 --src <年度 zip 所在目录>(或设 env GUANLAN_SITE_SRC)—— 本器不猜路径:' '开发机上的默认目录在别的机器上不存在') return 2 src, out = pathlib.Path(a.src), pathlib.Path(a.out) if a.list: return list_structure(src) if a.extract: return extract(src, out, a.year, a.month) if a.verify_raw: return verify_raw(src) if a.extract_csv: return extract_csv(src, only_missing=a.only_missing) if a.probe: return probe(out, a.cls) if a.to_csv: return to_csv(out, a.cls, a.year, a.month, a.table, a.turbines) ap.print_help() return 2 def to_csv(out: pathlib.Path, cls: str | None, year: str | None, month: str | None, table: str | None, turbines: str | None) -> int: ok, route = have_mdb_route() if not ok: print(f'[X] 本机没有 MDB 读取能力({route})—— 见 --probe 的补救办法') return 4 if not cls: print('[X] --to-csv 需要 --class(如 std)') return 2 mdbs = [m for m in sorted(out.rglob(f'*-{cls}.mdb')) if (not year or f'{year}年' in str(m)) and (not month or f'{month}月' in str(m))] if not mdbs: print(f'[X] 没找到 *-{cls}.mdb(先 --extract)') return 2 raw_dir = station_dir() / CLASS_TO_RAW.get(cls, f'scada_{cls}') print(f'== MDB → CSV · 类目 {cls} · {len(mdbs)} 个月 · 落点 {P.rel(raw_dir)} ==') want = {t.strip().upper() for t in turbines.split(',')} if turbines else None for m in mdbs: try: tbls = mdb_tables(m) except Exception as e: print(f' {m.name}: [X] {type(e).__name__}: {str(e)[:120]}') continue tbl = table or (tbls[0] if tbls else None) if not tbl: print(f' {m.name}: 没有表') continue tmp_csv = pathlib.Path(tempfile.gettempdir()) / f'_{m.stem}_{tbl}.csv' rc = mdb_to_csv(m, tbl, tmp_csv) if rc: continue import pandas as pd df = pd.read_csv(tmp_csv, dtype=str) print(f' {m.name} / {tbl}: {len(df)} 行 × {len(df.columns)} 列 → {list(df.columns)[:8]}') # 现场口径: Station = 机组号(91→WTG01), 每个机组一个文件、时间列名统一成 Time st = next((c for c in df.columns if c.lower() in ('station', 'stationid', 'wtg')), None) ts = next((c for c in df.columns if 'time' in c.lower() and 'station' not in c.lower()), None) if not st or not ts: print(f' [i] 没识别出 Station/TimeStamp 列 —— 该月份先原样留存 {P.rel(tmp_csv)}') continue base = df[st].astype(str).str.extract(r'(\d+)')[0].astype('Int64') df['_wtg'] = 'WTG' + (base - (90 if base.min() and base.min() >= 90 else 0)).astype('Int64').astype(str).str.zfill(2) for w, g in df.groupby('_wtg'): if want and w.upper() not in want: continue g2 = g.drop(columns=['_wtg']).rename(columns={ts: 'Time'}) dst = raw_dir / f'{w}.csv' head = not dst.is_file() g2.to_csv(dst, mode='a' if not head else 'w', header=head, index=False, encoding='utf-8') if head: print(f' → {P.rel(dst)} {len(g2)} 行(Time 列 = {ts})') print(' 完成。接下去: python scripts/raw_data_check.py(体检)→ python scripts/rebuild_all.py') return 0 if __name__ == '__main__': for _s in (sys.stdout, sys.stderr): try: _s.reconfigure(errors='replace') except Exception: pass sys.exit(main())