| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401 |
- #!/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())
|