csv_to_mdb.py 8.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188
  1. #!/usr/bin/env python3
  2. # -*- coding: utf-8 -*-
  3. r"""CSV → MDB(仿照现场 `25年.zip`/`26年.zip` 的形态)。2026-09-17 用户令。
  4. 形态: `<年>年/<月>月/<年>-<月>-<类>.zip → <年>-<月>-<类>.mdb`;每类一库(库内含该月全部机组)。
  5. **实测可用路线**(探路过程见 docs/现场收资接入_v0.1.md §7):
  6. ① `ADOX.Catalog` + ACE 12.0 建库 ✔
  7. ② 按列切片写分片 CSV(首列 `rid`,配 `schema.ini` 声明 UTF-8)✔
  8. ③ **`Access.Application.DoCmd.TransferText` 批量装入** ✔ ← 唯一可用的批量导入
  9. (`SELECT … INTO [text;…]`、目标库 `INSERT … SELECT [text;…]`、`IN '<dir>' 'text;…'` 全部报
  10. 「这种对象类型不支持该操作」;DAO `CreateDatabase` 报「找不到可安装的 ISAM」⇒ 都不用)
  11. ④ 读回校验行数 ✔
  12. 硬限制(必须偏离现场形态): 单表 ≤255 列(scada_10min 每台 598 列 ⇒ 拆表,每片 250 列);
  13. 单库 ≤2 GB(⇒ 按月分库;不按月分库时 1 分钟数据 13.7M 行会超)。
  14. 用法:
  15. python scripts/csv_to_mdb.py --month 2025-01 --turbines WTG01 # 试点(1 台 1 月)
  16. python scripts/csv_to_mdb.py --month 2025-01 # 该月全部机组
  17. python scripts/csv_to_mdb.py --month 2025-01 --class 1min [--turbines 01E]
  18. 退出码: 0 成功 · 2 参数/源不齐 · 4 无 Access/ACE 能力 · 5 失败
  19. """
  20. from __future__ import annotations
  21. import argparse
  22. import pathlib
  23. import re
  24. import shutil
  25. import subprocess
  26. import sys
  27. import tempfile
  28. import time
  29. import zipfile
  30. ROOT = pathlib.Path(__file__).resolve().parents[1]
  31. sys.path.insert(0, str(ROOT))
  32. from src.windscada.config import raw_station_dir # noqa: E402
  33. CHUNK = 250 # 每张表列数(含 rid)
  34. CLS_DIR = {'10min': 'scada_10min', '1min': 'scada_1min'}
  35. def safe_names(cols: list[str]) -> list[str]:
  36. out, seen = [], {}
  37. for i, c in enumerate(cols):
  38. n = re.sub(r'[^0-9A-Za-z_\u4e00-\u9fff]', '_', c)[:60] or f'c{i}'
  39. if n in seen:
  40. seen[n] += 1
  41. n = f'{n[:56]}_{seen[n]}'
  42. else:
  43. seen[n] = 0
  44. out.append(n)
  45. return out
  46. def build_chunks(files: list[pathlib.Path], month: str, tmp: pathlib.Path, limit_rows: int):
  47. """把各台该月行按列切片汇到 `<tmp>/t{i}.csv`(首列 rid)。→ (分片数, 总行数, 每台行数)"""
  48. nt, total, per = 0, 0, {}
  49. handles = {}
  50. try:
  51. for f in files:
  52. with open(f, encoding='utf-8', errors='replace', newline='') as fh:
  53. cols = fh.readline().rstrip('\n').split(',')
  54. names = safe_names(['rid'] + cols)
  55. if not nt:
  56. nt = max(1, -(-len(cols) // CHUNK))
  57. ini = []
  58. for i in range(nt):
  59. lo, hi = i * CHUNK, min(len(cols), (i + 1) * CHUNK)
  60. p = tmp / f't{i+1}.csv'
  61. p.write_text(','.join(names[lo:hi + 1]) + '\n', encoding='utf-8', newline='')
  62. handles[i] = open(p, 'a', encoding='utf-8', newline='')
  63. ini += [f'[{p.name}]', 'Format=CSVDelimited', 'ColNameHeader=True',
  64. 'CharacterSet=65001', '']
  65. (tmp / 'schema.ini').write_text('\n'.join(ini), encoding='utf-8')
  66. n = 0
  67. for line in fh:
  68. if line[:7] != month:
  69. continue
  70. vals = line.rstrip('\n').split(',')
  71. n += 1
  72. for i in range(nt):
  73. lo, hi = i * CHUNK, min(len(cols), (i + 1) * CHUNK)
  74. handles[i].write(','.join([str(n)] + vals[lo:hi]) + '\n')
  75. if limit_rows and n >= limit_rows:
  76. break
  77. per[f.stem] = n
  78. total += n
  79. finally:
  80. for h in handles.values():
  81. h.close()
  82. return nt, total, per
  83. def to_mdb(db: pathlib.Path, tmp: pathlib.Path, nt: int) -> list[str]:
  84. """ADOX 建库 + Access.TransferText 逐分片装入 + 读回行数。→ 每表行数描述。"""
  85. db.parent.mkdir(parents=True, exist_ok=True)
  86. ps = pathlib.Path(tempfile.gettempdir()) / '_csv2mdb_run.ps1'
  87. lines = ['$ErrorActionPreference="Stop"',
  88. f'$db = "{db}"',
  89. 'if (Test-Path $db) { Remove-Item $db -Force }',
  90. '$cat = New-Object -ComObject ADOX.Catalog',
  91. '[void]$cat.Create("Provider=Microsoft.ACE.OLEDB.12.0;Data Source=$db;")',
  92. '$cat = $null',
  93. '$acc = New-Object -ComObject Access.Application',
  94. '$acc.Visible = $false',
  95. '$acc.OpenCurrentDatabase($db)']
  96. for i in range(1, nt + 1):
  97. lines.append(f'$acc.DoCmd.TransferText(0, "", "t{i}", "{tmp / f"t{i}.csv"}", $true)')
  98. lines += ['$acc.CloseCurrentDatabase()', '$acc.Quit()',
  99. '$c = New-Object -ComObject ADODB.Connection',
  100. '$c.Open("Provider=Microsoft.ACE.OLEDB.12.0;Data Source=$db;")']
  101. for i in range(1, nt + 1):
  102. lines.append(f'$rs = $c.Execute("SELECT COUNT(*) FROM [t{i}]"); '
  103. f'Write-Output ("t{i}=" + $rs.Fields.Item(0).Value); $rs.Close()')
  104. lines.append('$c.Close()')
  105. ps.write_text('\n'.join(lines), encoding='utf-8-sig')
  106. r = subprocess.run(['powershell', '-NoProfile', '-ExecutionPolicy', 'Bypass', '-File', str(ps)],
  107. capture_output=True, text=True, errors='replace', timeout=7200)
  108. got = [ln.strip() for ln in (r.stdout or '').strip().splitlines() if ln.strip()]
  109. if r.returncode != 0:
  110. print(' [X] ' + (r.stderr or r.stdout or '').strip()[:240])
  111. return []
  112. return got
  113. def main() -> int:
  114. ap = argparse.ArgumentParser(description='CSV → MDB(仿现场年度归档形态)')
  115. ap.add_argument('--month', required=True, help='YYYY-MM')
  116. ap.add_argument('--class', dest='cls', default='10min', choices=('10min', '1min'))
  117. ap.add_argument('--turbines', default=None, help='逗号分隔(WTG01 或 01E);不给=全部')
  118. ap.add_argument('--out', default=None)
  119. ap.add_argument('--limit-rows', type=int, default=0)
  120. ap.add_argument('--no-zip', action='store_true')
  121. a = ap.parse_args()
  122. if not sys.platform.startswith('win'):
  123. print('[X] 需要 Windows(Access/ACE)—— Linux 侧无写库能力')
  124. return 4
  125. y, m = a.month.split('-')
  126. src = pathlib.Path(raw_station_dir()) / CLS_DIR[a.cls]
  127. out = pathlib.Path(a.out) if a.out else (pathlib.Path(raw_station_dir()) / 'scada_mdb')
  128. want = {t.strip() for t in a.turbines.split(',')} if a.turbines else None
  129. files = [f for f in sorted(src.glob('*.csv')) if (not want or f.stem in want)]
  130. if not files:
  131. print(f'[X] {src} 下没有匹配 CSV')
  132. return 2
  133. tmp = pathlib.Path(tempfile.gettempdir()) / f'csv2mdb_{y}{m}_{a.cls}'
  134. if tmp.exists():
  135. shutil.rmtree(tmp, ignore_errors=True)
  136. tmp.mkdir(parents=True, exist_ok=True)
  137. t0 = time.time()
  138. nt, total, per = build_chunks(files, a.month, tmp, a.limit_rows)
  139. print(f' 切片: {nt} 张表 · {total} 行 · {len(files)} 台 · 抽行 {time.time() - t0:.0f}s')
  140. if not total:
  141. print(f' [i] {a.month} 无数据')
  142. return 0
  143. db = out / f'{y}年' / f'{m}月' / f'{y}-{m}-{a.cls}.mdb'
  144. ts = time.time()
  145. got = to_mdb(db, tmp, nt)
  146. for ln in got:
  147. print(' ' + ln)
  148. if not got:
  149. return 5
  150. if not a.no_zip:
  151. z = db.parent / f'{y}-{m}-{a.cls}.zip'
  152. with zipfile.ZipFile(z, 'w', zipfile.ZIP_DEFLATED, compresslevel=6) as zf:
  153. zf.write(db, db.name)
  154. print(f' + {P_rel(out, db)} mdb {db.stat().st_size / 1e6:.1f} MB · 装入 {time.time() - ts:.0f}s'
  155. f' · 合计 {time.time() - t0:.0f}s')
  156. shutil.rmtree(tmp, ignore_errors=True)
  157. return 0
  158. def P_rel(base: pathlib.Path, p: pathlib.Path) -> str:
  159. try:
  160. return p.relative_to(base.parent).as_posix()
  161. except ValueError:
  162. return str(p)
  163. if __name__ == '__main__':
  164. for _s in (sys.stdout, sys.stderr):
  165. try:
  166. _s.reconfigure(errors='replace')
  167. except Exception:
  168. pass
  169. sys.exit(main())