#!/usr/bin/env python3 # -*- coding: utf-8 -*- r"""CSV → MDB(仿照现场 `25年.zip`/`26年.zip` 的形态) + 全量核对。2026-09-17 用户令。 形态: `data/raw/<场站>/scada_mdb/<年>年/<月>月/<年>-<月>-<类>.mdb`(+ 同名 zip);每类一库、库内按 250 列拆表。 实测可用路线(其余写法都失败, 见 docs/现场收资接入_v0.1.md §7): ADOX.Catalog 建库(ACE 12.0) → 分片 CSV(首列 rid, schema.ini CharacterSet=65001) → **Access.Application.DoCmd.TransferText(0,"",表,csv,$true,"",65001)** 批量装入 → ADODB 读回校验 ✗ `SELECT … INTO [text;…]` / 目标库 `INSERT … SELECT [text;…]` / `IN '' 'text;…'` → 「这种对象类型不支持该操作」 ✗ DAO `CreateDatabase` → 「找不到可安装的 ISAM」 ⚠ `TransferText` 必须给 **CodePage=65001**,否则 1min 的中文列名/值按 cp936 解成乱码(实测)。 硬限制: 单表 ≤255 列(10min 每台 598 列 ⇒ 拆 3 张表)· 单库 ≤2 GB(⇒ 按月分库)。 用法: python scripts/csv_to_mdb.py --month 2025-01 --class 10min # 转换某月(全部机组) python scripts/csv_to_mdb.py --verify # ★ 一条命令: 全量核对报告(日历推算, 秒出) python scripts/csv_to_mdb.py --verify --deep # 深核对: 扫源 CSV 逐月计数(几分钟, 带缓存) 退出码: 0 成功/全对 · 2 参数或源不齐 · 4 无 Access/ACE · 5 有差异/失败 """ from __future__ import annotations import argparse import calendar 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 CLS_DIR = {'10min': 'scada_10min', '1min': 'scada_1min'} CAD = {'10min': 144, '1min': 1440} # 每天行数(10 分钟 / 1 分钟节拍) def safe_names(cols: list[str]) -> list[str]: out, seen = [], {} for i, c in enumerate(cols): c = c.lstrip('\ufeff') # 源表头带 UTF-8 BOM ⇒ 不清掉会变成 _TimeStamp(实测) n = re.sub(r'[^0-9A-Za-z_\u4e00-\u9fff]', '_', c)[:60] or f'c{i}' 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, month: str, tmp: pathlib.Path, limit_rows: int): nt, total, handles = 0, 0, {} 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 total += n finally: for h in handles.values(): h.close() return nt, total def to_mdb(db: pathlib.Path, tmp: pathlib.Path, nt: int) -> list[str]: db.parent.mkdir(parents=True, exist_ok=True) ps = pathlib.Path(tempfile.gettempdir()) / '_csv2mdb_run.ps1' L = ['$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): L.append(f'$acc.DoCmd.TransferText(0, "", "t{i}", "{tmp / f"t{i}.csv"}", $true, "", 65001)') L += ['$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): L.append(f'$rs = $c.Execute("SELECT COUNT(*) FROM [t{i}]"); ' f'Write-Output ("t{i}=" + $rs.Fields.Item(0).Value); $rs.Close()') L.append('$c.Close()') ps.write_text('\n'.join(L), encoding='utf-8-sig') r = subprocess.run(['powershell', '-NoProfile', '-ExecutionPolicy', 'Bypass', '-File', str(ps)], capture_output=True, text=True, errors='replace', timeout=7200) if r.returncode != 0: print(' [X] ' + (r.stderr or r.stdout or '').strip()[:240]) return [] return [x.strip() for x in (r.stdout or '').strip().splitlines() if x.strip()] def mdb_rows(mdbs: list[pathlib.Path]) -> dict: """一次 PowerShell 批量读所有 MDB 的表行数 → {文件名: [t1,t2,…]}。""" L = ['$ErrorActionPreference="Continue"'] for m in mdbs: L.append(f'$c=New-Object -ComObject ADODB.Connection;' f'$c.Open("Provider=Microsoft.ACE.OLEDB.12.0;Data Source={m};")') L.append('$o=@(); foreach($t in @("t1","t2","t3","t4")){ try {' ' $rs=$c.Execute("SELECT COUNT(*) FROM [$t]"); $o+=$rs.Fields.Item(0).Value;' ' $rs.Close() } catch {} }') L.append(f'Write-Output ("{m.name}|" + ($o -join ",")); $c.Close()') ps = pathlib.Path(tempfile.gettempdir()) / '_verify_mdb.ps1' ps.write_text('\n'.join(L), encoding='utf-8-sig') r = subprocess.run(['powershell', '-NoProfile', '-ExecutionPolicy', 'Bypass', '-File', str(ps)], capture_output=True, text=True, errors='replace', timeout=7200) out = {} for ln in (r.stdout or '').splitlines(): if '|' in ln: k, v = ln.strip().split('|', 1) out[k] = [int(x) for x in v.split(',') if x.strip().isdigit()] return out def src_month_rows(cls_filter: str | None, cache: pathlib.Path, refresh: bool) -> dict: """扫源 CSV 逐月计数(二进制整文件 + 正则, 带缓存)。""" import json if cache.is_file() and not refresh: try: return json.loads(cache.read_text(encoding='utf-8')) except Exception: pass print(' [扫源 CSV 逐月计数…首次几分钟,之后走缓存]') t0, rx, out = time.time(), re.compile(rb'(?:^|\n)(\d{4}-\d{2})'), {} for cls, fn in CLS_DIR.items(): if cls_filter and cls != cls_filter: continue for f in sorted((pathlib.Path(raw_station_dir()) / fn).glob('*.csv')): data = f.read_bytes() nl = data.find(b'\n') body = data[nl + 1:] if nl >= 0 else data for mm in rx.findall(body): k = f'{cls}|{mm.decode()}' out[k] = out.get(k, 0) + 1 del data, body cache.write_text(json.dumps(out, indent=0), encoding='utf-8') print(f' 源计数完成 {time.time() - t0:.0f}s → {cache.name}') return out def verify(out: pathlib.Path, cls_filter: str | None, refresh: bool, deep: bool) -> int: mdbs = sorted(m for m in out.rglob('*-*.mdb') if not m.name.startswith('_')) if cls_filter: mdbs = [m for m in mdbs if m.stem.endswith(f'-{cls_filter}')] if not mdbs: print(f'[X] {out} 下还没有产物 mdb') return 2 got = mdb_rows(mdbs) src = src_month_rows(cls_filter, out / '_csv_month_rows.json', refresh) if deep else {} print(f'== CSV → MDB 全量核对 · {out} ==') print(f' {"类":6s} {"月份":8s} {"mdb 各表行数":18s} {"应得行数":10s} 判定 zip') bad = [] for m in mdbs: cls = m.stem.rsplit('-', 1)[1] ym = '-'.join(m.stem.split('-')[:2]) rows = got.get(m.name, []) want = src.get(f'{cls}|{ym}') if deep else None if want is None: y, mo = int(ym[:4]), int(ym[5:7]) want = CAD.get(cls, 0) * calendar.monthrange(y, mo)[1] * 38 aligned = len(set(rows)) == 1 if rows else False d = (rows[0] - want) if rows else -want # 日历模式: 少于日历=数据缺口(正常), 多于日历才是错; 深核对模式: 必须与源逐行相等 ok = bool(rows) and aligned and (rows[0] == want if deep else rows[0] <= want) tag = ('OK' if ok else '异常') if deep else ('一致' if d == 0 else (f'缺口{d:+d}' if d < 0 else f'超日历{d:+d}')) zp = m.parent / f'{ym}-{cls}.zip' ztxt = f'有 {zp.stat().st_size / 1e6:.1f} MB' if zp.is_file() else '无' print(f' {cls:6s} {ym:8s} {",".join(map(str, rows)):18s} {want:<10d} ' f'{tag:12s} {ztxt}') if not ok: bad.append((cls, ym, rows, want)) print(f' 合计: {len(mdbs)} 个月库 · 一致 {len(mdbs) - len(bad)} · 差异 {len(bad)}' + ('' if deep else '(按日历推算; 差异多为该月首尾不完整, 属正常 ⇒ 要按源逐行核对加 --deep)')) for cls, ym, rows, want in bad[:8]: print(f' [i] {cls} {ym}: mdb={rows} vs 应得 {want}') print(' 结论: ' + ('全部一致(各分片表行数相同 ⇒ 逐行对齐)' if not bad else '有差异, 见上')) return 0 if not bad else 5 def main() -> int: ap = argparse.ArgumentParser(description='CSV → MDB(仿现场年度归档形态)+ 全量核对') ap.add_argument('--month', default=None, help='YYYY-MM(--verify 时可省)') 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') ap.add_argument('--verify', action='store_true', help='全量核对报告(默认日历推算, 秒出)') ap.add_argument('--deep', action='store_true', help='--verify 深核对: 扫源 CSV 逐月计数') ap.add_argument('--refresh', action='store_true', help='--deep 时重扫源(不用缓存)') a = ap.parse_args() out = pathlib.Path(a.out) if a.out else (pathlib.Path(raw_station_dir()) / 'scada_mdb') if a.verify: return verify(out, a.cls, a.refresh, a.deep) if not a.month: print('[X] 需要 --month(或用 --verify)') return 2 if not sys.platform.startswith('win'): print('[X] 需要 Windows(Access/ACE)—— Linux 侧无写库能力') return 4 y, m = a.month.split('-') src_dir = pathlib.Path(raw_station_dir()) / CLS_DIR[a.cls] want = {t.strip() for t in a.turbines.split(',')} if a.turbines else None files = [f for f in sorted(src_dir.glob('*.csv')) if (not want or f.stem in want)] if not files: print(f'[X] {src_dir} 下没有匹配 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 = 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() res = to_mdb(db, tmp, nt) for ln in res: print(' ' + ln) if not res: 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' + {db.relative_to(out.parent).as_posix()} mdb {db.stat().st_size / 1e6:.1f} MB' f' · 装入 {time.time() - ts:.0f}s · 合计 {time.time() - t0:.0f}s') shutil.rmtree(tmp, ignore_errors=True) return 0 if __name__ == '__main__': for _s in (sys.stdout, sys.stderr): try: _s.reconfigure(errors='replace') except Exception: pass sys.exit(main())