| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348 |
- #!/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 '<dir>' '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]:
- """列名消毒 —— **单一实现**在 `src/windscada/mdb_names.py`(读侧 scada_source 要用同一套规则,
- 否则从 MDB 取数时按契约列名找不到列,静默少列)。这里只做转发,保持本脚本既有调用点不变。
- """
- from src.windscada.mdb_names import safe_names as _sn
- return _sn(cols)
- def build_chunks(files, month: str, tmp: pathlib.Path, limit_rows: int):
- r"""逐台 CSV → 分片 CSV(供 TransferText 装入)。
- ★2026-09-19 两处修正(都是为了"这份 MDB 能被**接回去当输入**用", 见 docs/现场收资接入_v0.1.md §8):
- ① **加 `turbine` 列**:原先每行只有 `rid` + 数据列, 逐台文件**堆在一起**, 而 1min 导出件里
- **没有任何机组标识列**(实测: `source_file` 全行同值 "如海风机1分钟数据/如海测点_2025-01.csv")
- ⇒ 那份 1min 库**按台取不到数**, 等于只能看不能算。现在首列写文件名的 stem(10min=`WTG01`,
- 1min=`01E`), 10min 另有的 `WTG` 列原样保留。
- ② **按列名对齐, 不再按位置堆**:原先用**第一个文件**的表头当 schema, 其余文件整行按位置写入。
- 10min 实测 38 件表头**完全一致**(1 种) ⇒ 无害; 但 1min 实测 **16 种表头**(同 79 列, 名与序都不同)
- ⇒ 位置堆法把 A 台的列写进 B 台的列位, **静默串列**(比缺列危险得多)。现在先扫一遍所有表头取并集,
- 逐行按列名落位, 缺列留空。
- """
- # ① 扫表头 → 并集(首见顺序) + 逐文件的列名→位置映射
- heads: dict[pathlib.Path, list[str]] = {}
- union: list[str] = []
- for f in files:
- with open(f, encoding='utf-8', errors='replace', newline='') as fh:
- cols = fh.readline().rstrip('\n').split(',')
- heads[f] = cols
- for c in cols:
- if c not in union:
- union.append(c)
- if len({tuple(v) for v in heads.values()}) > 1:
- print(f' [i] {len({tuple(v) for v in heads.values()})} 种表头 → 按**列名**对齐(并集 {len(union)} 列); '
- f'缺列留空 (原先按位置堆会静默串列)')
- names = safe_names(['rid', 'turbine'] + union)
- nt = max(1, -(-len(union) // CHUNK))
- total, handles = 0, {}
- try:
- for i in range(nt):
- lo, hi = i * CHUNK, min(len(union), (i + 1) * CHUNK)
- p = tmp / f't{i+1}.csv'
- p.write_text(','.join(names[0:2] + names[2 + lo:2 + hi]) + '\n', encoding='utf-8', newline='')
- handles[i] = open(p, 'a', encoding='utf-8', newline='')
- ini = []
- for i in range(nt):
- ini += [f'[t{i+1}.csv]', 'Format=CSVDelimited', 'ColNameHeader=True',
- 'CharacterSet=65001', '']
- (tmp / 'schema.ini').write_text('\n'.join(ini), encoding='utf-8')
- for f in files:
- cols = heads[f]
- idx = [union.index(c) for c in cols]
- with open(f, encoding='utf-8', errors='replace', newline='') as fh:
- fh.readline()
- n = 0
- for line in fh:
- if line[:7] != month:
- continue
- vals = line.rstrip('\n').split(',')
- row = [''] * len(union)
- for k, j in enumerate(idx):
- if k < len(vals):
- row[j] = vals[k]
- n += 1
- for i in range(nt):
- lo, hi = i * CHUNK, min(len(union), (i + 1) * CHUNK)
- handles[i].write(','.join([str(n), f.stem] + row[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)
- got = [x.strip() for x in (r.stdout or '').strip().splitlines() if x.strip()]
- counts = [x for x in got if re.fullmatch(r't\d+=\d+', x)]
- if r.returncode != 0:
- # ★2026-09-18 实逮: 全量转换里 1min 2025-07 / 2025-09 两个月退出码非 0, 而**数据其实已经装完** ——
- # stdout 里每张表的 `tN=<行数>` 都打出来了(且与切片行数逐月相符), 失败发生在装载之后的
- # COM 收尾 ($c.Close()/$acc.Quit() 那一步偶发 "未指定的错误")。原先这里直接判失败 ⇒
- # 整月被当成没跑成: 不打印 `+ …mdb`, 也不写 zip, 而 1.3 GB 的 .mdb 已经在盘上躺着
- # (每月的抽行+装入 ≈ 10 分钟, 白扔)。判据改成**看事实**: 所有分片表都报出了行数 ⇒ 算成功,
- # 但把退出码与 stderr 尾巴**响亮打出来**(不静默), 让人知道收尾那步出过岔子。
- if len(counts) < nt:
- print(' [X] ' + (r.stderr or r.stdout or '').strip()[:240])
- return []
- print(f' [!] PowerShell 收尾非零退出 (rc={r.returncode}), 但 {len(counts)}/{nt} 张分片表都报出了行数 '
- f'⇒ 按"已装成"处理; 尾巴: ' + (r.stderr or '').strip().replace('\n', ' ')[:160])
- return counts
- 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, short, over, minor = [], [], [], []
- BOUNDARY = 0.005 # 超日历 ≤0.5% 视为**月末边界行**(有些台次月 00:00 那一行落在本月)
- 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
- # 深核对: 必须与源逐行相等; 日历模式: 少行=该月数据本身不齐(正常), 多行超过 0.5% 才算转换错
- if deep:
- ok = bool(rows) and aligned and rows[0] == want
- else:
- ok = bool(rows) and aligned and d <= want * BOUNDARY
- tag = (('OK' if ok else '异常') if deep else
- ('一致' if d == 0 else (f'缺口{d:+d}' if d < 0 else
- (f'超日历{d:+d}' if d > want * BOUNDARY 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))
- if rows and d < 0:
- short.append((cls, ym, rows, want))
- elif rows and d > 0:
- (over if d > want * BOUNDARY else minor).append((cls, ym, rows, want))
- # ★2026-09-19 修: 原来汇总只报"一致 N / 差异 K", 而日历判据是 `rows <= want` ⇒ **任何缺口都算一致**,
- # 于是 1min 那张表有 8 个月短几百到一万多行, 汇总却是"一致 18 · 差异 0"(自相矛盾, 只有看表才发现)。
- # 现在四档分开报: 一致 / 缺口(少, 数据本身) / 边界(多但 ≤0.5%, 月末那一行) / 超日历(多且超 0.5%, 才是转换错)。
- print(f' 合计: {len(mdbs)} 个月库 · 一致 {len(mdbs) - len(bad)} · 缺口 {len(short)} · '
- f'边界 {len(minor)} · 超日历 {len(over)}'
- + ('' if deep else '(按日历推算; 少行=该月数据本身不齐, 多行>0.5% 才是转换错; 逐行对源加 --deep)'))
- for cls, ym, rows, want in sorted(short, key=lambda x: x[2][0] - x[3])[:3]:
- print(f' [i] 缺口最大: {cls} {ym}: {rows[0]:,} 行 (差 {rows[0] - want:+,} = '
- f'{100 * (rows[0] - want) / want:.2f}%)')
- for cls, ym, rows, want in bad[:8]:
- print(f' [i] {cls} {ym}: mdb={rows} vs 应得 {want}')
- print(' 结论: ' + ('全部一致(各分片表行数相同 ⇒ 逐行对齐)' if not bad else
- ('有差异, 见上' if over or deep else '分片表行数整齐; 少行属当月数据不齐(见上), 无转换错')))
- return 0 if not bad else 5
- def rezip(out: pathlib.Path, cls: str) -> int:
- """月库在位但 zip 缺失 → 只补 zip (不重算, 不碰 .mdb)。
- 为什么需要: 2026-09-18 全量转换实测 1min 2025-07 / 2025-09 两个月"装载后收尾失败" ⇒ 当月的
- zip 没写出来(见 to_mdb 的说明), 而 .mdb 是好的。重算一个月要 ~10 分钟, 补个 zip 只要几秒;
- 归档形态(现场是年度 zip)又不能缺, 故单列一个只读 .mdb 的补救口。
- """
- miss = [m for m in sorted(out.glob(f'*年/*月/*-{cls}.mdb')) if not m.with_suffix('.zip').exists()]
- if not miss:
- print(f' 没有缺 zip 的 {cls} 月库')
- return 0
- for m in miss:
- z = m.with_suffix('.zip')
- print(f' + 补 zip {m.name} ({m.stat().st_size / 1e6:.1f} MB) …', end='', flush=True)
- with zipfile.ZipFile(z, 'w', zipfile.ZIP_DEFLATED, compresslevel=6) as zf:
- zf.write(m, m.name)
- print(f' {z.stat().st_size / 1e6:.1f} MB')
- print(f' 共补 {len(miss)} 个 zip')
- return 0
- 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('--rezip', action='store_true', help='给已在位但缺 zip 的月库补 zip(不重算)')
- 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 a.rezip:
- return rezip(out, a.cls)
- 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())
|