csv_to_mdb.py 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269
  1. #!/usr/bin/env python3
  2. # -*- coding: utf-8 -*-
  3. r"""CSV → MDB(仿照现场 `25年.zip`/`26年.zip` 的形态) + 全量核对。2026-09-17 用户令。
  4. 形态: `data/raw/<场站>/scada_mdb/<年>年/<月>月/<年>-<月>-<类>.mdb`(+ 同名 zip);每类一库、库内按 250 列拆表。
  5. 实测可用路线(其余写法都失败, 见 docs/现场收资接入_v0.1.md §7):
  6. ADOX.Catalog 建库(ACE 12.0) → 分片 CSV(首列 rid, schema.ini CharacterSet=65001)
  7. → **Access.Application.DoCmd.TransferText(0,"",表,csv,$true,"",65001)** 批量装入 → ADODB 读回校验
  8. ✗ `SELECT … INTO [text;…]` / 目标库 `INSERT … SELECT [text;…]` / `IN '<dir>' 'text;…'` → 「这种对象类型不支持该操作」
  9. ✗ DAO `CreateDatabase` → 「找不到可安装的 ISAM」
  10. ⚠ `TransferText` 必须给 **CodePage=65001**,否则 1min 的中文列名/值按 cp936 解成乱码(实测)。
  11. 硬限制: 单表 ≤255 列(10min 每台 598 列 ⇒ 拆 3 张表)· 单库 ≤2 GB(⇒ 按月分库)。
  12. 用法:
  13. python scripts/csv_to_mdb.py --month 2025-01 --class 10min # 转换某月(全部机组)
  14. python scripts/csv_to_mdb.py --verify # ★ 一条命令: 全量核对报告(日历推算, 秒出)
  15. python scripts/csv_to_mdb.py --verify --deep # 深核对: 扫源 CSV 逐月计数(几分钟, 带缓存)
  16. 退出码: 0 成功/全对 · 2 参数或源不齐 · 4 无 Access/ACE · 5 有差异/失败
  17. """
  18. from __future__ import annotations
  19. import argparse
  20. import calendar
  21. import pathlib
  22. import re
  23. import shutil
  24. import subprocess
  25. import sys
  26. import tempfile
  27. import time
  28. import zipfile
  29. ROOT = pathlib.Path(__file__).resolve().parents[1]
  30. sys.path.insert(0, str(ROOT))
  31. from src.windscada.config import raw_station_dir # noqa: E402
  32. CHUNK = 250
  33. CLS_DIR = {'10min': 'scada_10min', '1min': 'scada_1min'}
  34. CAD = {'10min': 144, '1min': 1440} # 每天行数(10 分钟 / 1 分钟节拍)
  35. def safe_names(cols: list[str]) -> list[str]:
  36. out, seen = [], {}
  37. for i, c in enumerate(cols):
  38. c = c.lstrip('\ufeff') # 源表头带 UTF-8 BOM ⇒ 不清掉会变成 _TimeStamp(实测)
  39. n = re.sub(r'[^0-9A-Za-z_\u4e00-\u9fff]', '_', c)[:60] or f'c{i}'
  40. if n in seen:
  41. seen[n] += 1
  42. n = f'{n[:56]}_{seen[n]}'
  43. else:
  44. seen[n] = 0
  45. out.append(n)
  46. return out
  47. def build_chunks(files, month: str, tmp: pathlib.Path, limit_rows: int):
  48. nt, total, handles = 0, 0, {}
  49. try:
  50. for f in files:
  51. with open(f, encoding='utf-8', errors='replace', newline='') as fh:
  52. cols = fh.readline().rstrip('\n').split(',')
  53. names = safe_names(['rid'] + cols)
  54. if not nt:
  55. nt = max(1, -(-len(cols) // CHUNK))
  56. ini = []
  57. for i in range(nt):
  58. lo, hi = i * CHUNK, min(len(cols), (i + 1) * CHUNK)
  59. p = tmp / f't{i+1}.csv'
  60. p.write_text(','.join(names[lo:hi + 1]) + '\n', encoding='utf-8', newline='')
  61. handles[i] = open(p, 'a', encoding='utf-8', newline='')
  62. ini += [f'[{p.name}]', 'Format=CSVDelimited', 'ColNameHeader=True',
  63. 'CharacterSet=65001', '']
  64. (tmp / 'schema.ini').write_text('\n'.join(ini), encoding='utf-8')
  65. n = 0
  66. for line in fh:
  67. if line[:7] != month:
  68. continue
  69. vals = line.rstrip('\n').split(',')
  70. n += 1
  71. for i in range(nt):
  72. lo, hi = i * CHUNK, min(len(cols), (i + 1) * CHUNK)
  73. handles[i].write(','.join([str(n)] + vals[lo:hi]) + '\n')
  74. if limit_rows and n >= limit_rows:
  75. break
  76. total += n
  77. finally:
  78. for h in handles.values():
  79. h.close()
  80. return nt, total
  81. def to_mdb(db: pathlib.Path, tmp: pathlib.Path, nt: int) -> list[str]:
  82. db.parent.mkdir(parents=True, exist_ok=True)
  83. ps = pathlib.Path(tempfile.gettempdir()) / '_csv2mdb_run.ps1'
  84. L = ['$ErrorActionPreference="Stop"', f'$db = "{db}"',
  85. 'if (Test-Path $db) { Remove-Item $db -Force }',
  86. '$cat = New-Object -ComObject ADOX.Catalog',
  87. '[void]$cat.Create("Provider=Microsoft.ACE.OLEDB.12.0;Data Source=$db;")', '$cat = $null',
  88. '$acc = New-Object -ComObject Access.Application', '$acc.Visible = $false',
  89. '$acc.OpenCurrentDatabase($db)']
  90. for i in range(1, nt + 1):
  91. L.append(f'$acc.DoCmd.TransferText(0, "", "t{i}", "{tmp / f"t{i}.csv"}", $true, "", 65001)')
  92. L += ['$acc.CloseCurrentDatabase()', '$acc.Quit()',
  93. '$c = New-Object -ComObject ADODB.Connection',
  94. '$c.Open("Provider=Microsoft.ACE.OLEDB.12.0;Data Source=$db;")']
  95. for i in range(1, nt + 1):
  96. L.append(f'$rs = $c.Execute("SELECT COUNT(*) FROM [t{i}]"); '
  97. f'Write-Output ("t{i}=" + $rs.Fields.Item(0).Value); $rs.Close()')
  98. L.append('$c.Close()')
  99. ps.write_text('\n'.join(L), encoding='utf-8-sig')
  100. r = subprocess.run(['powershell', '-NoProfile', '-ExecutionPolicy', 'Bypass', '-File', str(ps)],
  101. capture_output=True, text=True, errors='replace', timeout=7200)
  102. if r.returncode != 0:
  103. print(' [X] ' + (r.stderr or r.stdout or '').strip()[:240])
  104. return []
  105. return [x.strip() for x in (r.stdout or '').strip().splitlines() if x.strip()]
  106. def mdb_rows(mdbs: list[pathlib.Path]) -> dict:
  107. """一次 PowerShell 批量读所有 MDB 的表行数 → {文件名: [t1,t2,…]}。"""
  108. L = ['$ErrorActionPreference="Continue"']
  109. for m in mdbs:
  110. L.append(f'$c=New-Object -ComObject ADODB.Connection;'
  111. f'$c.Open("Provider=Microsoft.ACE.OLEDB.12.0;Data Source={m};")')
  112. L.append('$o=@(); foreach($t in @("t1","t2","t3","t4")){ try {'
  113. ' $rs=$c.Execute("SELECT COUNT(*) FROM [$t]"); $o+=$rs.Fields.Item(0).Value;'
  114. ' $rs.Close() } catch {} }')
  115. L.append(f'Write-Output ("{m.name}|" + ($o -join ",")); $c.Close()')
  116. ps = pathlib.Path(tempfile.gettempdir()) / '_verify_mdb.ps1'
  117. ps.write_text('\n'.join(L), encoding='utf-8-sig')
  118. r = subprocess.run(['powershell', '-NoProfile', '-ExecutionPolicy', 'Bypass', '-File', str(ps)],
  119. capture_output=True, text=True, errors='replace', timeout=7200)
  120. out = {}
  121. for ln in (r.stdout or '').splitlines():
  122. if '|' in ln:
  123. k, v = ln.strip().split('|', 1)
  124. out[k] = [int(x) for x in v.split(',') if x.strip().isdigit()]
  125. return out
  126. def src_month_rows(cls_filter: str | None, cache: pathlib.Path, refresh: bool) -> dict:
  127. """扫源 CSV 逐月计数(二进制整文件 + 正则, 带缓存)。"""
  128. import json
  129. if cache.is_file() and not refresh:
  130. try:
  131. return json.loads(cache.read_text(encoding='utf-8'))
  132. except Exception:
  133. pass
  134. print(' [扫源 CSV 逐月计数…首次几分钟,之后走缓存]')
  135. t0, rx, out = time.time(), re.compile(rb'(?:^|\n)(\d{4}-\d{2})'), {}
  136. for cls, fn in CLS_DIR.items():
  137. if cls_filter and cls != cls_filter:
  138. continue
  139. for f in sorted((pathlib.Path(raw_station_dir()) / fn).glob('*.csv')):
  140. data = f.read_bytes()
  141. nl = data.find(b'\n')
  142. body = data[nl + 1:] if nl >= 0 else data
  143. for mm in rx.findall(body):
  144. k = f'{cls}|{mm.decode()}'
  145. out[k] = out.get(k, 0) + 1
  146. del data, body
  147. cache.write_text(json.dumps(out, indent=0), encoding='utf-8')
  148. print(f' 源计数完成 {time.time() - t0:.0f}s → {cache.name}')
  149. return out
  150. def verify(out: pathlib.Path, cls_filter: str | None, refresh: bool, deep: bool) -> int:
  151. mdbs = sorted(m for m in out.rglob('*-*.mdb') if not m.name.startswith('_'))
  152. if cls_filter:
  153. mdbs = [m for m in mdbs if m.stem.endswith(f'-{cls_filter}')]
  154. if not mdbs:
  155. print(f'[X] {out} 下还没有产物 mdb')
  156. return 2
  157. got = mdb_rows(mdbs)
  158. src = src_month_rows(cls_filter, out / '_csv_month_rows.json', refresh) if deep else {}
  159. print(f'== CSV → MDB 全量核对 · {out} ==')
  160. print(f' {"类":6s} {"月份":8s} {"mdb 各表行数":18s} {"应得行数":10s} 判定 zip')
  161. bad = []
  162. for m in mdbs:
  163. cls = m.stem.rsplit('-', 1)[1]
  164. ym = '-'.join(m.stem.split('-')[:2])
  165. rows = got.get(m.name, [])
  166. want = src.get(f'{cls}|{ym}') if deep else None
  167. if want is None:
  168. y, mo = int(ym[:4]), int(ym[5:7])
  169. want = CAD.get(cls, 0) * calendar.monthrange(y, mo)[1] * 38
  170. aligned = len(set(rows)) == 1 if rows else False
  171. d = (rows[0] - want) if rows else -want
  172. # 日历模式: 少于日历=数据缺口(正常), 多于日历才是错; 深核对模式: 必须与源逐行相等
  173. ok = bool(rows) and aligned and (rows[0] == want if deep else rows[0] <= want)
  174. tag = ('OK' if ok else '异常') if deep else ('一致' if d == 0 else (f'缺口{d:+d}' if d < 0 else f'超日历{d:+d}'))
  175. zp = m.parent / f'{ym}-{cls}.zip'
  176. ztxt = f'有 {zp.stat().st_size / 1e6:.1f} MB' if zp.is_file() else '无'
  177. print(f' {cls:6s} {ym:8s} {",".join(map(str, rows)):18s} {want:<10d} '
  178. f'{tag:12s} {ztxt}')
  179. if not ok:
  180. bad.append((cls, ym, rows, want))
  181. print(f' 合计: {len(mdbs)} 个月库 · 一致 {len(mdbs) - len(bad)} · 差异 {len(bad)}'
  182. + ('' if deep else '(按日历推算; 差异多为该月首尾不完整, 属正常 ⇒ 要按源逐行核对加 --deep)'))
  183. for cls, ym, rows, want in bad[:8]:
  184. print(f' [i] {cls} {ym}: mdb={rows} vs 应得 {want}')
  185. print(' 结论: ' + ('全部一致(各分片表行数相同 ⇒ 逐行对齐)' if not bad else '有差异, 见上'))
  186. return 0 if not bad else 5
  187. def main() -> int:
  188. ap = argparse.ArgumentParser(description='CSV → MDB(仿现场年度归档形态)+ 全量核对')
  189. ap.add_argument('--month', default=None, help='YYYY-MM(--verify 时可省)')
  190. ap.add_argument('--class', dest='cls', default='10min', choices=('10min', '1min'))
  191. ap.add_argument('--turbines', default=None, help='逗号分隔(WTG01 或 01E);不给=全部')
  192. ap.add_argument('--out', default=None)
  193. ap.add_argument('--limit-rows', type=int, default=0)
  194. ap.add_argument('--no-zip', action='store_true')
  195. ap.add_argument('--verify', action='store_true', help='全量核对报告(默认日历推算, 秒出)')
  196. ap.add_argument('--deep', action='store_true', help='--verify 深核对: 扫源 CSV 逐月计数')
  197. ap.add_argument('--refresh', action='store_true', help='--deep 时重扫源(不用缓存)')
  198. a = ap.parse_args()
  199. out = pathlib.Path(a.out) if a.out else (pathlib.Path(raw_station_dir()) / 'scada_mdb')
  200. if a.verify:
  201. return verify(out, a.cls, a.refresh, a.deep)
  202. if not a.month:
  203. print('[X] 需要 --month(或用 --verify)')
  204. return 2
  205. if not sys.platform.startswith('win'):
  206. print('[X] 需要 Windows(Access/ACE)—— Linux 侧无写库能力')
  207. return 4
  208. y, m = a.month.split('-')
  209. src_dir = pathlib.Path(raw_station_dir()) / CLS_DIR[a.cls]
  210. want = {t.strip() for t in a.turbines.split(',')} if a.turbines else None
  211. files = [f for f in sorted(src_dir.glob('*.csv')) if (not want or f.stem in want)]
  212. if not files:
  213. print(f'[X] {src_dir} 下没有匹配 CSV')
  214. return 2
  215. tmp = pathlib.Path(tempfile.gettempdir()) / f'csv2mdb_{y}{m}_{a.cls}'
  216. if tmp.exists():
  217. shutil.rmtree(tmp, ignore_errors=True)
  218. tmp.mkdir(parents=True, exist_ok=True)
  219. t0 = time.time()
  220. nt, total = build_chunks(files, a.month, tmp, a.limit_rows)
  221. print(f' 切片: {nt} 张表 · {total} 行 · {len(files)} 台 · 抽行 {time.time() - t0:.0f}s')
  222. if not total:
  223. print(f' [i] {a.month} 无数据')
  224. return 0
  225. db = out / f'{y}年' / f'{m}月' / f'{y}-{m}-{a.cls}.mdb'
  226. ts = time.time()
  227. res = to_mdb(db, tmp, nt)
  228. for ln in res:
  229. print(' ' + ln)
  230. if not res:
  231. return 5
  232. if not a.no_zip:
  233. z = db.parent / f'{y}-{m}-{a.cls}.zip'
  234. with zipfile.ZipFile(z, 'w', zipfile.ZIP_DEFLATED, compresslevel=6) as zf:
  235. zf.write(db, db.name)
  236. print(f' + {db.relative_to(out.parent).as_posix()} mdb {db.stat().st_size / 1e6:.1f} MB'
  237. f' · 装入 {time.time() - ts:.0f}s · 合计 {time.time() - t0:.0f}s')
  238. shutil.rmtree(tmp, ignore_errors=True)
  239. return 0
  240. if __name__ == '__main__':
  241. for _s in (sys.stdout, sys.stderr):
  242. try:
  243. _s.reconfigure(errors='replace')
  244. except Exception:
  245. pass
  246. sys.exit(main())