Explorar el Código

CSV→MDB: --verify 一条命令全量核对(日历模式区分 一致/缺口/超日历; --deep 按源 CSV 逐行核对, 带缓存)

zhouyang.xie hace 3 semanas
padre
commit
69c707dbee
Se han modificado 1 ficheros con 142 adiciones y 62 borrados
  1. 142 62
      scripts/csv_to_mdb.py

+ 142 - 62
scripts/csv_to_mdb.py

@@ -1,29 +1,28 @@
 #!/usr/bin/env python3
 # -*- coding: utf-8 -*-
-r"""CSV → MDB(仿照现场 `25年.zip`/`26年.zip` 的形态)。2026-09-17 用户令。
+r"""CSV → MDB(仿照现场 `25年.zip`/`26年.zip` 的形态) + 全量核对。2026-09-17 用户令。
 
-形态: `<年>年/<月>月/<年>-<月>-<类>.zip → <年>-<月>-<类>.mdb`;每类一库(库内含该月全部机组)。
+形态: `data/raw/<场站>/scada_mdb/<年>年/<月>月/<年>-<月>-<类>.mdb`(+ 同名 zip);每类一库、库内按 250 列拆表。
 
-**实测可用路线**(探路过程见 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 '<dir>' 'text;…'` 全部报
-      「这种对象类型不支持该操作」;DAO `CreateDatabase` 报「找不到可安装的 ISAM」⇒ 都不用)
-  ④ 读回校验行数 ✔
+实测可用路线(其余写法都失败, 见 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 '<dir>' 'text;…'` → 「这种对象类型不支持该操作」
+  ✗ DAO `CreateDatabase` → 「找不到可安装的 ISAM」
+  ⚠ `TransferText` 必须给 **CodePage=65001**,否则 1min 的中文列名/值按 cp936 解成乱码(实测)。
 
-硬限制(必须偏离现场形态): 单表 ≤255 列(scada_10min 每台 598 列 ⇒ 拆表,每片 250 列);
-单库 ≤2 GB(⇒ 按月分库;不按月分库时 1 分钟数据 13.7M 行会超)。
+硬限制: 单表 ≤255 列(10min 每台 598 列 ⇒ 拆 3 张表)· 单库 ≤2 GB(⇒ 按月分库)。
 
 用法:
-    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 失败
+    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
@@ -37,14 +36,15 @@ 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)
+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')      # ★ 源 CSV 表头带 UTF-8 BOM ⇒ 不清掉会被净化成 _TimeStamp(实测)
+        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
@@ -55,10 +55,8 @@ def safe_names(cols: list[str]) -> list[str]:
     return out
 
 
-def build_chunks(files: list[pathlib.Path], month: str, tmp: pathlib.Path, limit_rows: int):
-    """把各台该月行按列切片汇到 `<tmp>/t{i}.csv`(首列 rid)。→ (分片数, 总行数, 每台行数)"""
-    nt, total, per = 0, 0, {}
-    handles = {}
+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:
@@ -86,100 +84,182 @@ def build_chunks(files: list[pathlib.Path], month: str, tmp: pathlib.Path, limit
                         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
+    return nt, total
 
 
 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)']
+    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):
-        lines.append(f'$acc.DoCmd.TransferText(0, "", "t{i}", "{tmp / f"t{i}.csv"}", $true, "", 65001)')
-    lines += ['$acc.CloseCurrentDatabase()', '$acc.Quit()',
-              '$c = New-Object -ComObject ADODB.Connection',
-              '$c.Open("Provider=Microsoft.ACE.OLEDB.12.0;Data Source=$db;")']
+        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):
-        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')
+        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)
-    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
+    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', required=True, help='YYYY-MM')
+    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 = 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')
+    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.glob('*.csv')) if (not want or f.stem in want)]
+    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} 下没有匹配 CSV')
+        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, per = build_chunks(files, a.month, tmp, a.limit_rows)
+    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()
-    got = to_mdb(db, tmp, nt)
-    for ln in got:
+    res = to_mdb(db, tmp, nt)
+    for ln in res:
         print('  ' + ln)
-    if not got:
+    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'  + {P_rel(out, db)}  mdb {db.stat().st_size / 1e6:.1f} MB · 装入 {time.time() - ts:.0f}s'
-          f' · 合计 {time.time() - t0:.0f}s')
+    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
 
 
-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: