Selaa lähdekoodia

CSV→MDB 落地并试点成功: scripts/csv_to_mdb.py(ADOX 建库 + Access.TransferText 批量装入)

实测可用路线(前四种写法全部失败, 结论记在 docs §7):
· SELECT...INTO [text;…]; 目标库 INSERT...SELECT [text;…]; IN '<dir>' 'text;…'; DAO CreateDatabase
  —— 分别报「这种对象类型不支持该操作」/语法错误/找不到可安装的 ISAM。
· 可用: ADOX.Catalog 建库 + Access.Application.DoCmd.TransferText(0,'',表,分片csv,true) 批量装入 + ADODB 读回校验。
试点(WGT01 · 2025-01 · 10min): 3 张表 × 4,464 行 · mdb 30.2 MB · 抽行 2s + 装入 11s = 13s;
逐值抽检 TimeStamp/StationId/cnt_wtc_TrdHrT 与源 CSV 一致。
形态: data/raw/如东/scada_mdb/<年>年/<月>月/<年>-<月>-<类>.mdb(+同名 zip), 每类一库、库内按 250 列拆表。
zhouyang.xie 3 viikkoa sitten
vanhempi
commit
b28ba4f97b
1 muutettua tiedostoa jossa 188 lisäystä ja 0 poistoa
  1. 188 0
      scripts/csv_to_mdb.py

+ 188 - 0
scripts/csv_to_mdb.py

@@ -0,0 +1,188 @@
+#!/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 '<dir>' '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):
+        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: list[pathlib.Path], month: str, tmp: pathlib.Path, limit_rows: int):
+    """把各台该月行按列切片汇到 `<tmp>/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())