#!/usr/bin/env python3 # -*- coding: utf-8 -*- r"""CSV → MDB(仿照现场 `25年.zip`/`26年.zip` 的形态)。2026-09-17 用户令。 形态: `<年>年/<月>月/<年>-<月>-<类>.zip → <年>-<月>-<类>.mdb`;每类一库(库内含该月全部机组)。 **实测可用路线**(探路过程见 docs/现场收资接入_v0.1.md §7): ① `ADOX.Catalog` + ACE 12.0 建库 ✔ ② 按列切片写分片 CSV(首列 `rid`,配 `schema.ini` 声明 UTF-8)✔ ③ **`Access.Application.DoCmd.TransferText` 批量装入** ✔ ← 唯一可用的批量导入 (`SELECT … INTO [text;…]`、目标库 `INSERT … SELECT [text;…]`、`IN '' 'text;…'` 全部报 「这种对象类型不支持该操作」;DAO `CreateDatabase` 报「找不到可安装的 ISAM」⇒ 都不用) ④ 读回校验行数 ✔ 硬限制(必须偏离现场形态): 单表 ≤255 列(scada_10min 每台 598 列 ⇒ 拆表,每片 250 列); 单库 ≤2 GB(⇒ 按月分库;不按月分库时 1 分钟数据 13.7M 行会超)。 用法: python scripts/csv_to_mdb.py --month 2025-01 --turbines WTG01 # 试点(1 台 1 月) python scripts/csv_to_mdb.py --month 2025-01 # 该月全部机组 python scripts/csv_to_mdb.py --month 2025-01 --class 1min [--turbines 01E] 退出码: 0 成功 · 2 参数/源不齐 · 4 无 Access/ACE 能力 · 5 失败 """ from __future__ import annotations import argparse import pathlib import re import shutil import subprocess import sys import tempfile import time import zipfile ROOT = pathlib.Path(__file__).resolve().parents[1] sys.path.insert(0, str(ROOT)) from src.windscada.config import raw_station_dir # noqa: E402 CHUNK = 250 # 每张表列数(含 rid) CLS_DIR = {'10min': 'scada_10min', '1min': 'scada_1min'} def safe_names(cols: list[str]) -> list[str]: out, seen = [], {} for i, c in enumerate(cols): c = c.lstrip('\ufeff') # ★ 源 CSV 表头带 UTF-8 BOM ⇒ 不清掉会被净化成 _TimeStamp(实测) if n in seen: seen[n] += 1 n = f'{n[:56]}_{seen[n]}' else: seen[n] = 0 out.append(n) return out def build_chunks(files: list[pathlib.Path], month: str, tmp: pathlib.Path, limit_rows: int): """把各台该月行按列切片汇到 `/t{i}.csv`(首列 rid)。→ (分片数, 总行数, 每台行数)""" nt, total, per = 0, 0, {} handles = {} try: for f in files: with open(f, encoding='utf-8', errors='replace', newline='') as fh: cols = fh.readline().rstrip('\n').split(',') names = safe_names(['rid'] + cols) if not nt: nt = max(1, -(-len(cols) // CHUNK)) ini = [] for i in range(nt): lo, hi = i * CHUNK, min(len(cols), (i + 1) * CHUNK) p = tmp / f't{i+1}.csv' p.write_text(','.join(names[lo:hi + 1]) + '\n', encoding='utf-8', newline='') handles[i] = open(p, 'a', encoding='utf-8', newline='') ini += [f'[{p.name}]', 'Format=CSVDelimited', 'ColNameHeader=True', 'CharacterSet=65001', ''] (tmp / 'schema.ini').write_text('\n'.join(ini), encoding='utf-8') n = 0 for line in fh: if line[:7] != month: continue vals = line.rstrip('\n').split(',') n += 1 for i in range(nt): lo, hi = i * CHUNK, min(len(cols), (i + 1) * CHUNK) handles[i].write(','.join([str(n)] + vals[lo:hi]) + '\n') if limit_rows and n >= limit_rows: break per[f.stem] = n total += n finally: for h in handles.values(): h.close() return nt, total, per def to_mdb(db: pathlib.Path, tmp: pathlib.Path, nt: int) -> list[str]: """ADOX 建库 + Access.TransferText 逐分片装入 + 读回行数。→ 每表行数描述。""" db.parent.mkdir(parents=True, exist_ok=True) ps = pathlib.Path(tempfile.gettempdir()) / '_csv2mdb_run.ps1' lines = ['$ErrorActionPreference="Stop"', f'$db = "{db}"', 'if (Test-Path $db) { Remove-Item $db -Force }', '$cat = New-Object -ComObject ADOX.Catalog', '[void]$cat.Create("Provider=Microsoft.ACE.OLEDB.12.0;Data Source=$db;")', '$cat = $null', '$acc = New-Object -ComObject Access.Application', '$acc.Visible = $false', '$acc.OpenCurrentDatabase($db)'] for i in range(1, nt + 1): lines.append(f'$acc.DoCmd.TransferText(0, "", "t{i}", "{tmp / f"t{i}.csv"}", $true)') lines += ['$acc.CloseCurrentDatabase()', '$acc.Quit()', '$c = New-Object -ComObject ADODB.Connection', '$c.Open("Provider=Microsoft.ACE.OLEDB.12.0;Data Source=$db;")'] for i in range(1, nt + 1): lines.append(f'$rs = $c.Execute("SELECT COUNT(*) FROM [t{i}]"); ' f'Write-Output ("t{i}=" + $rs.Fields.Item(0).Value); $rs.Close()') lines.append('$c.Close()') ps.write_text('\n'.join(lines), encoding='utf-8-sig') r = subprocess.run(['powershell', '-NoProfile', '-ExecutionPolicy', 'Bypass', '-File', str(ps)], capture_output=True, text=True, errors='replace', timeout=7200) got = [ln.strip() for ln in (r.stdout or '').strip().splitlines() if ln.strip()] if r.returncode != 0: print(' [X] ' + (r.stderr or r.stdout or '').strip()[:240]) return [] return got def main() -> int: ap = argparse.ArgumentParser(description='CSV → MDB(仿现场年度归档形态)') ap.add_argument('--month', required=True, help='YYYY-MM') ap.add_argument('--class', dest='cls', default='10min', choices=('10min', '1min')) ap.add_argument('--turbines', default=None, help='逗号分隔(WTG01 或 01E);不给=全部') ap.add_argument('--out', default=None) ap.add_argument('--limit-rows', type=int, default=0) ap.add_argument('--no-zip', action='store_true') a = ap.parse_args() if not sys.platform.startswith('win'): print('[X] 需要 Windows(Access/ACE)—— Linux 侧无写库能力') return 4 y, m = a.month.split('-') src = pathlib.Path(raw_station_dir()) / CLS_DIR[a.cls] out = pathlib.Path(a.out) if a.out else (pathlib.Path(raw_station_dir()) / 'scada_mdb') want = {t.strip() for t in a.turbines.split(',')} if a.turbines else None files = [f for f in sorted(src.glob('*.csv')) if (not want or f.stem in want)] if not files: print(f'[X] {src} 下没有匹配 CSV') return 2 tmp = pathlib.Path(tempfile.gettempdir()) / f'csv2mdb_{y}{m}_{a.cls}' if tmp.exists(): shutil.rmtree(tmp, ignore_errors=True) tmp.mkdir(parents=True, exist_ok=True) t0 = time.time() nt, total, per = build_chunks(files, a.month, tmp, a.limit_rows) print(f' 切片: {nt} 张表 · {total} 行 · {len(files)} 台 · 抽行 {time.time() - t0:.0f}s') if not total: print(f' [i] {a.month} 无数据') return 0 db = out / f'{y}年' / f'{m}月' / f'{y}-{m}-{a.cls}.mdb' ts = time.time() got = to_mdb(db, tmp, nt) for ln in got: print(' ' + ln) if not got: return 5 if not a.no_zip: z = db.parent / f'{y}-{m}-{a.cls}.zip' with zipfile.ZipFile(z, 'w', zipfile.ZIP_DEFLATED, compresslevel=6) as zf: zf.write(db, db.name) print(f' + {P_rel(out, db)} mdb {db.stat().st_size / 1e6:.1f} MB · 装入 {time.time() - ts:.0f}s' f' · 合计 {time.time() - t0:.0f}s') shutil.rmtree(tmp, ignore_errors=True) return 0 def P_rel(base: pathlib.Path, p: pathlib.Path) -> str: try: return p.relative_to(base.parent).as_posix() except ValueError: return str(p) if __name__ == '__main__': for _s in (sys.stdout, sys.stderr): try: _s.reconfigure(errors='replace') except Exception: pass sys.exit(main())