extract_site_data.py 19 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393
  1. #!/usr/bin/env python3
  2. # -*- coding: utf-8 -*-
  3. r"""现场收资接入:从场站交来的年度 zip(内层是每月每类 `.mdb`)提取并转成观澜 `data/raw/<场站>/`。
  4. ## 场站交来的东西长什么样(2026-09-17 实测 `F:\temp\如东风场数据`)
  5. ```
  6. 25年.zip (1,022 MB) / 26年.zip (961 MB)
  7. └ 26年/1月/2026-01-cnt.zip → 2026-01-cnt.mdb (52.7 MB)
  8. 2026-01-std.zip → 2026-01-std.mdb ( 2.3 MB) ← tblGrid: 10 分钟电网/机组量测
  9. 2026-01-tur.zip → 2026-01-tur.mdb (47.6 MB)
  10. 2026-01-scd.zip → 2026-01-scd.mdb ( 5.6 MB)
  11. … 每月 12 类:cnt din dot flg grd int prs scd std sum tmp tur
  12. ```
  13. 即 **Access 数据库按月分类**(不是 CSV)。实测 `std.mdb` 里是单表 `tblGrid`:
  14. ```
  15. TimeStamp | Station | WPSStatus | TimestampStation | CurrentL1..L3 | VoltageL1..L3 |
  16. ActivePower | ReActivePower | ActivePowerExport | ReActivePowerExport | …
  17. ```
  18. `TimeStamp` 是 10 分钟节拍,`Station` 是机组号(如 91/92 —— 与 WTG 号的对应关系见 §映射)。
  19. ## 读 MDB 的路线(离线、零 pip 依赖)
  20. Windows 上用 **.NET `System.Data.Odbc` + 系统自带的 `Microsoft Access Driver (*.mdb, *.accdb)`**
  21. (本机实测可用;驱动来自 Office/Access Runtime,纯服务器上可能没有 ⇒ 本器会**先探针、缺了就响亮报并给补救**)。
  22. Linux 上用 `mdbtools` 的 `mdb-export`(本机没有)。两条路都不需要往包里塞新 pip 包。
  23. ## 用法
  24. python scripts/extract_site_data.py --src <年度zip所在目录> --list
  25. python scripts/extract_site_data.py --src <…> --extract [--year 26] [--month 1] [--out data/raw/如东/scada_mdb]
  26. python scripts/extract_site_data.py --src <…> --probe [--class std]
  27. python scripts/extract_site_data.py --src <…> --to-csv --class std [--year 26 --month 1] [--turbines WTG01,WTG02]
  28. 退出码: 0 成功 · 2 来源/参数不齐 · 4 平台无 MDB 读取能力(已给补救办法) · 5 转换/校验失败
  29. """
  30. from __future__ import annotations
  31. import argparse
  32. import hashlib
  33. import json
  34. import pathlib
  35. import shutil
  36. import subprocess
  37. import sys
  38. import tempfile
  39. import zipfile
  40. ROOT = pathlib.Path(__file__).resolve().parents[1]
  41. sys.path.insert(0, str(ROOT))
  42. from src import paths as P # noqa: E402
  43. WIN32 = sys.name == 'nt' if hasattr(sys, 'name') else False
  44. WIN32 = sys.platform.startswith('win')
  45. DRIVER = 'Microsoft Access Driver (*.mdb, *.accdb)'
  46. # 类目 → 观澜侧的落点(先按已知的上线:只有 std 与 sum 是 10min 量测;其余待现场口径确认)
  47. CLASS_TO_RAW = {
  48. 'std': 'scada_10min', # 实测: tblGrid = 10 分钟电网/机组量测(TimeStamp/Station/ActivePower/…)
  49. 'sum': 'scada_10min', # 待确认(体量最大, 疑为统计汇总)
  50. 'cnt': 'scada_10min', # 待确认
  51. 'tur': 'scada_10min', # 待确认(47 MB, 疑为机组状态)
  52. }
  53. # 现场交来的**已导出 CSV** zip → 观澜 raw 目录(2026-09-17 实测:与 data/raw 现有 38 件逐名对应)
  54. ZIP_TO_RAW = {
  55. '10分钟数据': 'scada_10min', # WTG01.csv … WTG38.csv(38 件 / 14.4 GB)
  56. '1分钟数据': 'scada_1min', # 01E.csv … 38B.csv(38 件 / 12.8 GB,机组代号命名)
  57. }
  58. def verify_raw(src: pathlib.Path, farm: str | None = None) -> int:
  59. """对账:现场 zip 里的 CSV ↔ `data/raw/<场>/<类>/` 现有文件(逐件比大小)。
  60. 这一步回答的是"**现有的输入数据到底是不是从这两个 zip 来的**" —— 不去解压十几 GB 就能验。
  61. """
  62. farm = farm or P.farm()
  63. zips = [z for z in find_year_zips(src) if z.stem in ZIP_TO_RAW]
  64. if not zips:
  65. print(f'[X] {src} 下没有 {"/".join(ZIP_TO_RAW)}.zip')
  66. return 2
  67. print(f'== 现场 zip ↔ data/raw/{farm}/ 对账 ==')
  68. bad = 0
  69. for z in zips:
  70. d = station_dir(farm) / ZIP_TO_RAW[z.stem]
  71. with zipfile.ZipFile(z) as zf:
  72. ents = [e for e in zf.infolist() if not e.is_dir()]
  73. same = miss = diff = 0
  74. miss_names, diff_names = [], []
  75. for e in ents:
  76. tgt = d / pathlib.Path(e.filename).name
  77. if not tgt.is_file():
  78. miss += 1
  79. miss_names.append(tgt.name)
  80. elif tgt.stat().st_size == e.file_size:
  81. same += 1
  82. else:
  83. diff += 1
  84. diff_names.append(f'{tgt.name}(zip {e.file_size:,} vs 盘 {tgt.stat().st_size:,})')
  85. bad += miss + diff
  86. print(f' {z.name} → {P.rel(d)}: zip {len(ents)} 件 · 逐字节同大小 {same} · 缺 {miss} · 大小不同 {diff}')
  87. for n in miss_names[:4]:
  88. print(f' 缺: {n}')
  89. for n in diff_names[:4]:
  90. print(f' 异: {n}')
  91. print(f' 结论: {"现有 raw 输入与现场 zip 逐件一致(可追溯)" if bad == 0 else str(bad) + " 件需要重新提取"} '
  92. f'(缺件提取: --extract-csv --only-missing)')
  93. return 0 if bad == 0 else 5
  94. def extract_csv(src: pathlib.Path, farm: str | None = None, only_missing: bool = False) -> int:
  95. """把"已导出 CSV"的现场 zip 解到 `data/raw/<场>/<类>/`(幂等:同名同大小即跳过)。"""
  96. farm = farm or P.farm()
  97. zips = [z for z in find_year_zips(src) if z.stem in ZIP_TO_RAW]
  98. if not zips:
  99. print(f'[X] {src} 下没有 {"/".join(ZIP_TO_RAW)}.zip')
  100. return 2
  101. done = skipped = 0
  102. for z in zips:
  103. d = station_dir(farm) / ZIP_TO_RAW[z.stem]
  104. d.mkdir(parents=True, exist_ok=True)
  105. print(f' {z.name} → {P.rel(d)}')
  106. with zipfile.ZipFile(z) as zf:
  107. for e in zf.infolist():
  108. if e.is_dir():
  109. continue
  110. tgt = d / pathlib.Path(e.filename).name
  111. if tgt.is_file() and tgt.stat().st_size == e.file_size:
  112. skipped += 1
  113. continue
  114. if only_missing and tgt.is_file():
  115. skipped += 1
  116. continue
  117. print(f' + {tgt.name} {e.file_size / 1e6:.1f} MB')
  118. with zf.open(e) as fh, open(tgt, 'wb') as out:
  119. shutil.copyfileobj(fh, out, 1024 * 1024 * 8)
  120. done += 1
  121. print(f' 完成: 新解 {done} 件 · 跳过(已同大小){skipped} 件')
  122. print(' 接下去: python scripts/raw_data_check.py(体检)→ python scripts/rebuild_all.py(重算)')
  123. return 0
  124. def station_dir(farm: str | None = None) -> pathlib.Path:
  125. """`data/raw/<场站>/` —— ★场站**目录名**可与 farm 键不同(键 = rudong,目录 = 如东),
  126. 所以必须走唯一取用口 `raw_station_dir()`;第一版按 `P.RAW_ROOT / farm` 拼,于是"38 件全缺"
  127. (对着 data/raw/rudong 找, 而东西在 data/raw/如东)。"""
  128. from src.windscada.config import raw_station_dir
  129. return pathlib.Path(raw_station_dir(farm or P.farm()))
  130. def find_year_zips(src: pathlib.Path) -> list[pathlib.Path]:
  131. return sorted([p for p in src.glob('*.zip') if p.is_file()])
  132. def list_structure(src: pathlib.Path) -> int:
  133. zips = find_year_zips(src)
  134. if not zips:
  135. print(f'[X] {src} 下没有年度 zip')
  136. return 2
  137. print(f'== 现场收资结构 · {src} ==')
  138. for z in zips:
  139. with zipfile.ZipFile(z) as zf:
  140. names = zf.namelist()
  141. months = sorted({n.split('/')[1] for n in names if n.count('/') >= 1 and n.split('/')[1]})
  142. classes = sorted({n.split('/')[-1].split('-')[-1].replace('.zip', '')
  143. for n in names if n.lower().endswith('.zip')})
  144. print(f' {z.name} {z.stat().st_size / 1e6:.0f} MB · {len(names)} 条目')
  145. print(f' 月份 {len(months)} 个: {", ".join(months)}')
  146. print(f' 类目 {len(classes)} 个: {", ".join(classes)}')
  147. return 0
  148. def extract(src: pathlib.Path, out: pathlib.Path, year: str | None, month: str | None) -> int:
  149. zips = [z for z in find_year_zips(src) if (not year or z.stem.startswith(year))]
  150. if not zips:
  151. print('[X] 没有匹配的年度 zip')
  152. return 2
  153. out.mkdir(parents=True, exist_ok=True)
  154. man_p = out / '_manifest.json'
  155. man = json.loads(man_p.read_text(encoding='utf-8')) if man_p.is_file() else dict(zips=[], files=[])
  156. done = 0
  157. for z in zips:
  158. with zipfile.ZipFile(z) as zf:
  159. inner = [n for n in zf.namelist() if n.lower().endswith('.zip')
  160. and (not month or f'/{month}月/' in n)]
  161. print(f' {z.name}: 内层 {len(inner)} 个分类 zip' + (f'(限 {month} 月)' if month else ''))
  162. for n in inner:
  163. parts = n.split('/') # 26年/1月/2026-01-std.zip
  164. rel = pathlib.Path(*parts[1:]) if len(parts) > 1 else pathlib.Path(parts[0])
  165. dst_dir = out / rel.parent
  166. dst_dir.mkdir(parents=True, exist_ok=True)
  167. with zf.open(n) as fh, zipfile.ZipFile(fh) as inner_zf:
  168. for m in inner_zf.namelist():
  169. tgt = dst_dir / pathlib.Path(m).name
  170. if tgt.is_file() and tgt.stat().st_size == inner_zf.getinfo(m).file_size:
  171. continue
  172. tgt.write_bytes(inner_zf.read(m))
  173. h = hashlib.sha256(tgt.read_bytes()).hexdigest()[:16]
  174. man['files'] = [x for x in man['files'] if x['path'] != str(tgt.relative_to(out))]
  175. man['files'].append(dict(path=str(tgt.relative_to(out)), bytes=tgt.stat().st_size,
  176. sha256_16=h, from_zip=f'{z.name}:{n}'))
  177. done += 1
  178. print(f' + {tgt.relative_to(out)} {tgt.stat().st_size / 1e6:.1f} MB sha256:{h}')
  179. man['zips'] = sorted({*man.get('zips', []), *[z.name for z in zips]})
  180. man_p.write_text(json.dumps(man, ensure_ascii=False, indent=1) + '\n', encoding='utf-8')
  181. print(f' 完成: 新解 {done} 个 MDB → {P.rel(out)}(清单 {P.rel(man_p)})')
  182. return 0
  183. def have_mdb_route() -> tuple[bool, str]:
  184. if WIN32:
  185. ps = ('$d=(Get-OdbcDriver | Where-Object { $_.Name -like "*Access Driver (*.mdb*" }).Count;'
  186. 'if($d -gt 0){"ok"}else{"no"}')
  187. try:
  188. r = subprocess.run(['powershell', '-NoProfile', '-Command', ps], capture_output=True, text=True,
  189. errors='replace', timeout=120)
  190. ok = 'ok' in (r.stdout or '')
  191. return ok, 'Windows + .NET System.Data.Odbc + 系统 Access ODBC 驱动' if ok else '缺 Access ODBC 驱动'
  192. except Exception as e:
  193. return False, f'探测失败: {type(e).__name__}'
  194. ok = bool(shutil.which('mdb-export') or shutil.which('mdb-tables'))
  195. return ok, 'mdbtools (mdb-export)' if ok else '缺 mdbtools'
  196. def mdb_tables(mdb: pathlib.Path) -> list[str]:
  197. if WIN32:
  198. ps = (f'$cs="Driver={{{DRIVER}}};Dbq={mdb};";'
  199. '$cn=New-Object System.Data.Odbc.OdbcConnection($cs);$cn.Open();'
  200. '$cn.GetSchema("Tables") | Where-Object { $_.TABLE_TYPE -eq "TABLE" } | '
  201. 'ForEach-Object { $_.TABLE_NAME };$cn.Close()')
  202. r = subprocess.run(['powershell', '-NoProfile', '-Command', ps], capture_output=True, text=True,
  203. errors='replace', timeout=300)
  204. if r.returncode != 0:
  205. raise RuntimeError((r.stderr or r.stdout or '').strip()[:300])
  206. return [x.strip() for x in (r.stdout or '').splitlines() if x.strip()]
  207. r = subprocess.run(['mdb-tables', '-1', str(mdb)], capture_output=True, text=True, errors='replace', timeout=300)
  208. return [x.strip() for x in (r.stdout or '').splitlines() if x.strip()]
  209. def mdb_to_csv(mdb: pathlib.Path, table: str, out_csv: pathlib.Path, where: str = '') -> int:
  210. """MDB 单表 → CSV(UTF-8)。Windows 走 PowerShell/.NET Odbc;Linux 走 mdb-export。"""
  211. out_csv.parent.mkdir(parents=True, exist_ok=True)
  212. if WIN32:
  213. q = f'SELECT * FROM [{table}]' + (f' WHERE {where}' if where else '')
  214. ps = f'''$cs="Driver={{{DRIVER}}};Dbq={mdb};"
  215. $cn=New-Object System.Data.Odbc.OdbcConnection($cs);$cn.Open()
  216. $cmd=$cn.CreateCommand();$cmd.CommandText="{q}"
  217. $da=New-Object System.Data.Odbc.OdbcDataAdapter($cmd)
  218. $dt=New-Object System.Data.DataTable;[void]$da.Fill($dt)
  219. $dt | Export-Csv -Path "{out_csv}" -NoTypeInformation -Encoding UTF8
  220. "rows=$($dt.Rows.Count)"
  221. $cn.Close()'''
  222. tf = pathlib.Path(tempfile.gettempdir()) / '_mdb_dump.ps1'
  223. tf.write_text(ps, encoding='utf-8-sig')
  224. r = subprocess.run(['powershell', '-NoProfile', '-ExecutionPolicy', 'Bypass', '-File', str(tf)],
  225. capture_output=True, text=True, errors='replace', timeout=3600)
  226. print(' ' + (r.stdout or '').strip().replace('\n', ' ')[:80])
  227. if r.returncode != 0:
  228. print(' [X] ' + (r.stderr or r.stdout or '').strip()[:200])
  229. return 5
  230. return 0
  231. cmd = ['mdb-export', str(mdb), table] + (['-q', where] if where else [])
  232. with open(out_csv, 'w', encoding='utf-8', newline='') as fh:
  233. r = subprocess.run(cmd, stdout=fh, stderr=subprocess.PIPE, text=True, errors='replace', timeout=3600)
  234. if r.returncode != 0:
  235. print(' [X] ' + (r.stderr or '')[:200])
  236. return 5
  237. return 0
  238. def probe(out: pathlib.Path, cls: str | None) -> int:
  239. ok, route = have_mdb_route()
  240. print(f'== MDB 读取能力 ==\n {route}: {"可用" if ok else "不可用"}')
  241. if not ok:
  242. print(' 补救: Windows 装 Access Runtime(或 Office)后自带 Access ODBC 驱动; '
  243. 'Linux 装 mdbtools(apt/yum install mdbtools); 或让现场把 MDB 导出成 CSV。')
  244. return 4
  245. mdbs = sorted(out.rglob('*.mdb'))
  246. mdbs = [m for m in mdbs if (not cls or f'-{cls}.' in m.name)]
  247. if not mdbs:
  248. print(f' [i] {P.rel(out)} 下还没有 MDB —— 先跑 --extract')
  249. return 0
  250. print(f' 已解出 MDB {len(mdbs)} 个(限类目 {cls or "全部"})')
  251. for m in mdbs[:8]:
  252. try:
  253. t = mdb_tables(m)
  254. print(f' {m.relative_to(out)}: 表 {t}')
  255. except Exception as e:
  256. print(f' {m.relative_to(out)}: [X] {type(e).__name__}: {str(e)[:120]}')
  257. return 0
  258. def main() -> int:
  259. ap = argparse.ArgumentParser(description='现场收资(年度 zip → 每月每类 MDB)提取与接入')
  260. ap.add_argument('--src', default=str(pathlib.Path('F:/temp/如东风场数据')),
  261. help='年度 zip 所在目录(默认按现场交来的目录名)')
  262. ap.add_argument('--out', default=str(station_dir() / 'scada_mdb'),
  263. help='MDB 落点(默认 data/raw/<场站>/scada_mdb)')
  264. ap.add_argument('--list', action='store_true', help='只列结构(年度 zip / 月份 / 类目)')
  265. ap.add_argument('--extract', action='store_true', help='解压出 MDB(幂等, 带清单)')
  266. ap.add_argument('--probe', action='store_true', help='探针: 读取能力 + 已解出 MDB 的表')
  267. ap.add_argument('--verify-raw', action='store_true', help='对账: 现场 CSV zip ↔ data/raw 现有文件(逐件比大小)')
  268. ap.add_argument('--extract-csv', action='store_true', help='把现场 CSV zip 解到 data/raw(幂等)')
  269. ap.add_argument('--only-missing', action='store_true', help='与 --extract-csv 合用: 只补缺件, 不覆盖已存在的')
  270. ap.add_argument('--to-csv', action='store_true', help='MDB → 观澜 data/raw 的 CSV(按类目映射)')
  271. ap.add_argument('--class', dest='cls', default=None, help='类目: cnt/sum/std/tur/…')
  272. ap.add_argument('--year', default=None, help='年度前缀(如 26)')
  273. ap.add_argument('--month', default=None, help='月份数字(如 1)')
  274. ap.add_argument('--table', default=None, help='指定表名(默认取该 MDB 的第一张表)')
  275. ap.add_argument('--turbines', default=None, help='只导这些机组(逗号分隔, 如 WTG01,WTG02)')
  276. a = ap.parse_args()
  277. src, out = pathlib.Path(a.src), pathlib.Path(a.out)
  278. if a.list:
  279. return list_structure(src)
  280. if a.extract:
  281. return extract(src, out, a.year, a.month)
  282. if a.verify_raw:
  283. return verify_raw(src)
  284. if a.extract_csv:
  285. return extract_csv(src, only_missing=a.only_missing)
  286. if a.probe:
  287. return probe(out, a.cls)
  288. if a.to_csv:
  289. return to_csv(out, a.cls, a.year, a.month, a.table, a.turbines)
  290. ap.print_help()
  291. return 2
  292. def to_csv(out: pathlib.Path, cls: str | None, year: str | None, month: str | None,
  293. table: str | None, turbines: str | None) -> int:
  294. ok, route = have_mdb_route()
  295. if not ok:
  296. print(f'[X] 本机没有 MDB 读取能力({route})—— 见 --probe 的补救办法')
  297. return 4
  298. if not cls:
  299. print('[X] --to-csv 需要 --class(如 std)')
  300. return 2
  301. mdbs = [m for m in sorted(out.rglob(f'*-{cls}.mdb')) if (not year or f'{year}年' in str(m))
  302. and (not month or f'{month}月' in str(m))]
  303. if not mdbs:
  304. print(f'[X] 没找到 *-{cls}.mdb(先 --extract)')
  305. return 2
  306. raw_dir = station_dir() / CLASS_TO_RAW.get(cls, f'scada_{cls}')
  307. print(f'== MDB → CSV · 类目 {cls} · {len(mdbs)} 个月 · 落点 {P.rel(raw_dir)} ==')
  308. want = {t.strip().upper() for t in turbines.split(',')} if turbines else None
  309. for m in mdbs:
  310. try:
  311. tbls = mdb_tables(m)
  312. except Exception as e:
  313. print(f' {m.name}: [X] {type(e).__name__}: {str(e)[:120]}')
  314. continue
  315. tbl = table or (tbls[0] if tbls else None)
  316. if not tbl:
  317. print(f' {m.name}: 没有表')
  318. continue
  319. tmp_csv = pathlib.Path(tempfile.gettempdir()) / f'_{m.stem}_{tbl}.csv'
  320. rc = mdb_to_csv(m, tbl, tmp_csv)
  321. if rc:
  322. continue
  323. import pandas as pd
  324. df = pd.read_csv(tmp_csv, dtype=str)
  325. print(f' {m.name} / {tbl}: {len(df)} 行 × {len(df.columns)} 列 → {list(df.columns)[:8]}')
  326. # 现场口径: Station = 机组号(91→WTG01), 每个机组一个文件、时间列名统一成 Time
  327. st = next((c for c in df.columns if c.lower() in ('station', 'stationid', 'wtg')), None)
  328. ts = next((c for c in df.columns if 'time' in c.lower() and 'station' not in c.lower()), None)
  329. if not st or not ts:
  330. print(f' [i] 没识别出 Station/TimeStamp 列 —— 该月份先原样留存 {P.rel(tmp_csv)}')
  331. continue
  332. base = df[st].astype(str).str.extract(r'(\d+)')[0].astype('Int64')
  333. df['_wtg'] = 'WTG' + (base - (90 if base.min() and base.min() >= 90 else 0)).astype('Int64').astype(str).str.zfill(2)
  334. for w, g in df.groupby('_wtg'):
  335. if want and w.upper() not in want:
  336. continue
  337. g2 = g.drop(columns=['_wtg']).rename(columns={ts: 'Time'})
  338. dst = raw_dir / f'{w}.csv'
  339. head = not dst.is_file()
  340. g2.to_csv(dst, mode='a' if not head else 'w', header=head, index=False, encoding='utf-8')
  341. if head:
  342. print(f' → {P.rel(dst)} {len(g2)} 行(Time 列 = {ts})')
  343. print(' 完成。接下去: python scripts/raw_data_check.py(体检)→ python scripts/rebuild_all.py')
  344. return 0
  345. if __name__ == '__main__':
  346. for _s in (sys.stdout, sys.stderr):
  347. try:
  348. _s.reconfigure(errors='replace')
  349. except Exception:
  350. pass
  351. sys.exit(main())