csv_to_mdb.py 19 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354
  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. try:
  20. from app_common.app_common_guanlan.api import install_root as _install_root
  21. except ImportError: # 理论不可达;包结构异常时回退到按位置上跳
  22. from pathlib import Path as _P
  23. def _install_root(_f): return _P(_f).resolve().parents[3]
  24. import argparse
  25. import calendar
  26. import pathlib
  27. import re
  28. import shutil
  29. import subprocess
  30. import sys
  31. import tempfile
  32. import time
  33. import zipfile
  34. ROOT = _install_root(__file__) # 模块化后按标记找安装根(原 parents[1] 已不成立)
  35. sys.path.insert(0, str(ROOT))
  36. from src.windscada.config import raw_station_dir # noqa: E402
  37. CHUNK = 250
  38. CLS_DIR = {'10min': 'scada_10min', '1min': 'scada_1min'}
  39. CAD = {'10min': 144, '1min': 1440} # 每天行数(10 分钟 / 1 分钟节拍)
  40. def safe_names(cols: list[str]) -> list[str]:
  41. """列名消毒 —— **单一实现**在 `src/windscada/mdb_names.py`(读侧 scada_source 要用同一套规则,
  42. 否则从 MDB 取数时按契约列名找不到列,静默少列)。这里只做转发,保持本脚本既有调用点不变。
  43. """
  44. from src.windscada.mdb_names import safe_names as _sn
  45. return _sn(cols)
  46. def build_chunks(files, month: str, tmp: pathlib.Path, limit_rows: int):
  47. r"""逐台 CSV → 分片 CSV(供 TransferText 装入)。
  48. ★2026-09-19 两处修正(都是为了"这份 MDB 能被**接回去当输入**用", 见 docs/现场收资接入_v0.1.md §8):
  49. ① **加 `turbine` 列**:原先每行只有 `rid` + 数据列, 逐台文件**堆在一起**, 而 1min 导出件里
  50. **没有任何机组标识列**(实测: `source_file` 全行同值 "如海风机1分钟数据/如海测点_2025-01.csv")
  51. ⇒ 那份 1min 库**按台取不到数**, 等于只能看不能算。现在首列写文件名的 stem(10min=`WTG01`,
  52. 1min=`01E`), 10min 另有的 `WTG` 列原样保留。
  53. ② **按列名对齐, 不再按位置堆**:原先用**第一个文件**的表头当 schema, 其余文件整行按位置写入。
  54. 10min 实测 38 件表头**完全一致**(1 种) ⇒ 无害; 但 1min 实测 **16 种表头**(同 79 列, 名与序都不同)
  55. ⇒ 位置堆法把 A 台的列写进 B 台的列位, **静默串列**(比缺列危险得多)。现在先扫一遍所有表头取并集,
  56. 逐行按列名落位, 缺列留空。
  57. """
  58. # ① 扫表头 → 并集(首见顺序) + 逐文件的列名→位置映射
  59. heads: dict[pathlib.Path, list[str]] = {}
  60. union: list[str] = []
  61. for f in files:
  62. with open(f, encoding='utf-8', errors='replace', newline='') as fh:
  63. cols = fh.readline().rstrip('\n').split(',')
  64. heads[f] = cols
  65. for c in cols:
  66. if c not in union:
  67. union.append(c)
  68. if len({tuple(v) for v in heads.values()}) > 1:
  69. print(f' [i] {len({tuple(v) for v in heads.values()})} 种表头 → 按**列名**对齐(并集 {len(union)} 列); '
  70. f'缺列留空 (原先按位置堆会静默串列)')
  71. names = safe_names(['rid', 'turbine'] + union)
  72. nt = max(1, -(-len(union) // CHUNK))
  73. total, handles = 0, {}
  74. try:
  75. for i in range(nt):
  76. lo, hi = i * CHUNK, min(len(union), (i + 1) * CHUNK)
  77. p = tmp / f't{i+1}.csv'
  78. p.write_text(','.join(names[0:2] + names[2 + lo:2 + hi]) + '\n', encoding='utf-8', newline='')
  79. handles[i] = open(p, 'a', encoding='utf-8', newline='')
  80. ini = []
  81. for i in range(nt):
  82. ini += [f'[t{i+1}.csv]', 'Format=CSVDelimited', 'ColNameHeader=True',
  83. 'CharacterSet=65001', '']
  84. (tmp / 'schema.ini').write_text('\n'.join(ini), encoding='utf-8')
  85. for f in files:
  86. cols = heads[f]
  87. idx = [union.index(c) for c in cols]
  88. with open(f, encoding='utf-8', errors='replace', newline='') as fh:
  89. fh.readline()
  90. n = 0
  91. for line in fh:
  92. if line[:7] != month:
  93. continue
  94. vals = line.rstrip('\n').split(',')
  95. row = [''] * len(union)
  96. for k, j in enumerate(idx):
  97. if k < len(vals):
  98. row[j] = vals[k]
  99. n += 1
  100. for i in range(nt):
  101. lo, hi = i * CHUNK, min(len(union), (i + 1) * CHUNK)
  102. handles[i].write(','.join([str(n), f.stem] + row[lo:hi]) + '\n')
  103. if limit_rows and n >= limit_rows:
  104. break
  105. total += n
  106. finally:
  107. for h in handles.values():
  108. h.close()
  109. return nt, total
  110. def to_mdb(db: pathlib.Path, tmp: pathlib.Path, nt: int) -> list[str]:
  111. db.parent.mkdir(parents=True, exist_ok=True)
  112. ps = pathlib.Path(tempfile.gettempdir()) / '_csv2mdb_run.ps1'
  113. L = ['$ErrorActionPreference="Stop"', f'$db = "{db}"',
  114. 'if (Test-Path $db) { Remove-Item $db -Force }',
  115. '$cat = New-Object -ComObject ADOX.Catalog',
  116. '[void]$cat.Create("Provider=Microsoft.ACE.OLEDB.12.0;Data Source=$db;")', '$cat = $null',
  117. '$acc = New-Object -ComObject Access.Application', '$acc.Visible = $false',
  118. '$acc.OpenCurrentDatabase($db)']
  119. for i in range(1, nt + 1):
  120. L.append(f'$acc.DoCmd.TransferText(0, "", "t{i}", "{tmp / f"t{i}.csv"}", $true, "", 65001)')
  121. L += ['$acc.CloseCurrentDatabase()', '$acc.Quit()',
  122. '$c = New-Object -ComObject ADODB.Connection',
  123. '$c.Open("Provider=Microsoft.ACE.OLEDB.12.0;Data Source=$db;")']
  124. for i in range(1, nt + 1):
  125. L.append(f'$rs = $c.Execute("SELECT COUNT(*) FROM [t{i}]"); '
  126. f'Write-Output ("t{i}=" + $rs.Fields.Item(0).Value); $rs.Close()')
  127. L.append('$c.Close()')
  128. ps.write_text('\n'.join(L), encoding='utf-8-sig')
  129. r = subprocess.run(['powershell', '-NoProfile', '-ExecutionPolicy', 'Bypass', '-File', str(ps)],
  130. capture_output=True, text=True, errors='replace', timeout=7200)
  131. got = [x.strip() for x in (r.stdout or '').strip().splitlines() if x.strip()]
  132. counts = [x for x in got if re.fullmatch(r't\d+=\d+', x)]
  133. if r.returncode != 0:
  134. # ★2026-09-18 实逮: 全量转换里 1min 2025-07 / 2025-09 两个月退出码非 0, 而**数据其实已经装完** ——
  135. # stdout 里每张表的 `tN=<行数>` 都打出来了(且与切片行数逐月相符), 失败发生在装载之后的
  136. # COM 收尾 ($c.Close()/$acc.Quit() 那一步偶发 "未指定的错误")。原先这里直接判失败 ⇒
  137. # 整月被当成没跑成: 不打印 `+ …mdb`, 也不写 zip, 而 1.3 GB 的 .mdb 已经在盘上躺着
  138. # (每月的抽行+装入 ≈ 10 分钟, 白扔)。判据改成**看事实**: 所有分片表都报出了行数 ⇒ 算成功,
  139. # 但把退出码与 stderr 尾巴**响亮打出来**(不静默), 让人知道收尾那步出过岔子。
  140. if len(counts) < nt:
  141. print(' [X] ' + (r.stderr or r.stdout or '').strip()[:240])
  142. return []
  143. print(f' [!] PowerShell 收尾非零退出 (rc={r.returncode}), 但 {len(counts)}/{nt} 张分片表都报出了行数 '
  144. f'⇒ 按"已装成"处理; 尾巴: ' + (r.stderr or '').strip().replace('\n', ' ')[:160])
  145. return counts
  146. def mdb_rows(mdbs: list[pathlib.Path]) -> dict:
  147. """一次 PowerShell 批量读所有 MDB 的表行数 → {文件名: [t1,t2,…]}。"""
  148. L = ['$ErrorActionPreference="Continue"']
  149. for m in mdbs:
  150. L.append(f'$c=New-Object -ComObject ADODB.Connection;'
  151. f'$c.Open("Provider=Microsoft.ACE.OLEDB.12.0;Data Source={m};")')
  152. L.append('$o=@(); foreach($t in @("t1","t2","t3","t4")){ try {'
  153. ' $rs=$c.Execute("SELECT COUNT(*) FROM [$t]"); $o+=$rs.Fields.Item(0).Value;'
  154. ' $rs.Close() } catch {} }')
  155. L.append(f'Write-Output ("{m.name}|" + ($o -join ",")); $c.Close()')
  156. ps = pathlib.Path(tempfile.gettempdir()) / '_verify_mdb.ps1'
  157. ps.write_text('\n'.join(L), encoding='utf-8-sig')
  158. r = subprocess.run(['powershell', '-NoProfile', '-ExecutionPolicy', 'Bypass', '-File', str(ps)],
  159. capture_output=True, text=True, errors='replace', timeout=7200)
  160. out = {}
  161. for ln in (r.stdout or '').splitlines():
  162. if '|' in ln:
  163. k, v = ln.strip().split('|', 1)
  164. out[k] = [int(x) for x in v.split(',') if x.strip().isdigit()]
  165. return out
  166. def src_month_rows(cls_filter: str | None, cache: pathlib.Path, refresh: bool) -> dict:
  167. """扫源 CSV 逐月计数(二进制整文件 + 正则, 带缓存)。"""
  168. import json
  169. if cache.is_file() and not refresh:
  170. try:
  171. return json.loads(cache.read_text(encoding='utf-8'))
  172. except Exception:
  173. pass
  174. print(' [扫源 CSV 逐月计数…首次几分钟,之后走缓存]')
  175. t0, rx, out = time.time(), re.compile(rb'(?:^|\n)(\d{4}-\d{2})'), {}
  176. for cls, fn in CLS_DIR.items():
  177. if cls_filter and cls != cls_filter:
  178. continue
  179. for f in sorted((pathlib.Path(raw_station_dir()) / fn).glob('*.csv')):
  180. data = f.read_bytes()
  181. nl = data.find(b'\n')
  182. body = data[nl + 1:] if nl >= 0 else data
  183. for mm in rx.findall(body):
  184. k = f'{cls}|{mm.decode()}'
  185. out[k] = out.get(k, 0) + 1
  186. del data, body
  187. cache.write_text(json.dumps(out, indent=0), encoding='utf-8')
  188. print(f' 源计数完成 {time.time() - t0:.0f}s → {cache.name}')
  189. return out
  190. def verify(out: pathlib.Path, cls_filter: str | None, refresh: bool, deep: bool) -> int:
  191. mdbs = sorted(m for m in out.rglob('*-*.mdb') if not m.name.startswith('_'))
  192. if cls_filter:
  193. mdbs = [m for m in mdbs if m.stem.endswith(f'-{cls_filter}')]
  194. if not mdbs:
  195. print(f'[X] {out} 下还没有产物 mdb')
  196. return 2
  197. got = mdb_rows(mdbs)
  198. src = src_month_rows(cls_filter, out / '_csv_month_rows.json', refresh) if deep else {}
  199. print(f'== CSV → MDB 全量核对 · {out} ==')
  200. print(f' {"类":6s} {"月份":8s} {"mdb 各表行数":18s} {"应得行数":10s} 判定 zip')
  201. bad, short, over, minor = [], [], [], []
  202. BOUNDARY = 0.005 # 超日历 ≤0.5% 视为**月末边界行**(有些台次月 00:00 那一行落在本月)
  203. for m in mdbs:
  204. cls = m.stem.rsplit('-', 1)[1]
  205. ym = '-'.join(m.stem.split('-')[:2])
  206. rows = got.get(m.name, [])
  207. want = src.get(f'{cls}|{ym}') if deep else None
  208. if want is None:
  209. y, mo = int(ym[:4]), int(ym[5:7])
  210. want = CAD.get(cls, 0) * calendar.monthrange(y, mo)[1] * 38
  211. aligned = len(set(rows)) == 1 if rows else False
  212. d = (rows[0] - want) if rows else -want
  213. # 深核对: 必须与源逐行相等; 日历模式: 少行=该月数据本身不齐(正常), 多行超过 0.5% 才算转换错
  214. if deep:
  215. ok = bool(rows) and aligned and rows[0] == want
  216. else:
  217. ok = bool(rows) and aligned and d <= want * BOUNDARY
  218. tag = (('OK' if ok else '异常') if deep else
  219. ('一致' if d == 0 else (f'缺口{d:+d}' if d < 0 else
  220. (f'超日历{d:+d}' if d > want * BOUNDARY else f'边界{d:+d}'))))
  221. zp = m.parent / f'{ym}-{cls}.zip'
  222. ztxt = f'有 {zp.stat().st_size / 1e6:.1f} MB' if zp.is_file() else '无'
  223. print(f' {cls:6s} {ym:8s} {",".join(map(str, rows)):18s} {want:<10d} '
  224. f'{tag:12s} {ztxt}')
  225. if not ok:
  226. bad.append((cls, ym, rows, want))
  227. if rows and d < 0:
  228. short.append((cls, ym, rows, want))
  229. elif rows and d > 0:
  230. (over if d > want * BOUNDARY else minor).append((cls, ym, rows, want))
  231. # ★2026-09-19 修: 原来汇总只报"一致 N / 差异 K", 而日历判据是 `rows <= want` ⇒ **任何缺口都算一致**,
  232. # 于是 1min 那张表有 8 个月短几百到一万多行, 汇总却是"一致 18 · 差异 0"(自相矛盾, 只有看表才发现)。
  233. # 现在四档分开报: 一致 / 缺口(少, 数据本身) / 边界(多但 ≤0.5%, 月末那一行) / 超日历(多且超 0.5%, 才是转换错)。
  234. print(f' 合计: {len(mdbs)} 个月库 · 一致 {len(mdbs) - len(bad)} · 缺口 {len(short)} · '
  235. f'边界 {len(minor)} · 超日历 {len(over)}'
  236. + ('' if deep else '(按日历推算; 少行=该月数据本身不齐, 多行>0.5% 才是转换错; 逐行对源加 --deep)'))
  237. for cls, ym, rows, want in sorted(short, key=lambda x: x[2][0] - x[3])[:3]:
  238. print(f' [i] 缺口最大: {cls} {ym}: {rows[0]:,} 行 (差 {rows[0] - want:+,} = '
  239. f'{100 * (rows[0] - want) / want:.2f}%)')
  240. for cls, ym, rows, want in bad[:8]:
  241. print(f' [i] {cls} {ym}: mdb={rows} vs 应得 {want}')
  242. print(' 结论: ' + ('全部一致(各分片表行数相同 ⇒ 逐行对齐)' if not bad else
  243. ('有差异, 见上' if over or deep else '分片表行数整齐; 少行属当月数据不齐(见上), 无转换错')))
  244. return 0 if not bad else 5
  245. def rezip(out: pathlib.Path, cls: str) -> int:
  246. """月库在位但 zip 缺失 → 只补 zip (不重算, 不碰 .mdb)。
  247. 为什么需要: 2026-09-18 全量转换实测 1min 2025-07 / 2025-09 两个月"装载后收尾失败" ⇒ 当月的
  248. zip 没写出来(见 to_mdb 的说明), 而 .mdb 是好的。重算一个月要 ~10 分钟, 补个 zip 只要几秒;
  249. 归档形态(现场是年度 zip)又不能缺, 故单列一个只读 .mdb 的补救口。
  250. """
  251. miss = [m for m in sorted(out.glob(f'*年/*月/*-{cls}.mdb')) if not m.with_suffix('.zip').exists()]
  252. if not miss:
  253. print(f' 没有缺 zip 的 {cls} 月库')
  254. return 0
  255. for m in miss:
  256. z = m.with_suffix('.zip')
  257. print(f' + 补 zip {m.name} ({m.stat().st_size / 1e6:.1f} MB) …', end='', flush=True)
  258. with zipfile.ZipFile(z, 'w', zipfile.ZIP_DEFLATED, compresslevel=6) as zf:
  259. zf.write(m, m.name)
  260. print(f' {z.stat().st_size / 1e6:.1f} MB')
  261. print(f' 共补 {len(miss)} 个 zip')
  262. return 0
  263. def main() -> int:
  264. ap = argparse.ArgumentParser(description='CSV → MDB(仿现场年度归档形态)+ 全量核对')
  265. ap.add_argument('--month', default=None, help='YYYY-MM(--verify 时可省)')
  266. ap.add_argument('--class', dest='cls', default='10min', choices=('10min', '1min'))
  267. ap.add_argument('--turbines', default=None, help='逗号分隔(WTG01 或 01E);不给=全部')
  268. ap.add_argument('--out', default=None)
  269. ap.add_argument('--limit-rows', type=int, default=0)
  270. ap.add_argument('--no-zip', action='store_true')
  271. ap.add_argument('--rezip', action='store_true', help='给已在位但缺 zip 的月库补 zip(不重算)')
  272. ap.add_argument('--verify', action='store_true', help='全量核对报告(默认日历推算, 秒出)')
  273. ap.add_argument('--deep', action='store_true', help='--verify 深核对: 扫源 CSV 逐月计数')
  274. ap.add_argument('--refresh', action='store_true', help='--deep 时重扫源(不用缓存)')
  275. a = ap.parse_args()
  276. out = pathlib.Path(a.out) if a.out else (pathlib.Path(raw_station_dir()) / 'scada_mdb')
  277. if a.verify:
  278. return verify(out, a.cls, a.refresh, a.deep)
  279. if a.rezip:
  280. return rezip(out, a.cls)
  281. if not a.month:
  282. print('[X] 需要 --month(或用 --verify)')
  283. return 2
  284. if not sys.platform.startswith('win'):
  285. print('[X] 需要 Windows(Access/ACE)—— Linux 侧无写库能力')
  286. return 4
  287. y, m = a.month.split('-')
  288. src_dir = pathlib.Path(raw_station_dir()) / CLS_DIR[a.cls]
  289. want = {t.strip() for t in a.turbines.split(',')} if a.turbines else None
  290. files = [f for f in sorted(src_dir.glob('*.csv')) if (not want or f.stem in want)]
  291. if not files:
  292. print(f'[X] {src_dir} 下没有匹配 CSV')
  293. return 2
  294. tmp = pathlib.Path(tempfile.gettempdir()) / f'csv2mdb_{y}{m}_{a.cls}'
  295. if tmp.exists():
  296. shutil.rmtree(tmp, ignore_errors=True)
  297. tmp.mkdir(parents=True, exist_ok=True)
  298. t0 = time.time()
  299. nt, total = build_chunks(files, a.month, tmp, a.limit_rows)
  300. print(f' 切片: {nt} 张表 · {total} 行 · {len(files)} 台 · 抽行 {time.time() - t0:.0f}s')
  301. if not total:
  302. print(f' [i] {a.month} 无数据')
  303. return 0
  304. db = out / f'{y}年' / f'{m}月' / f'{y}-{m}-{a.cls}.mdb'
  305. ts = time.time()
  306. res = to_mdb(db, tmp, nt)
  307. for ln in res:
  308. print(' ' + ln)
  309. if not res:
  310. return 5
  311. if not a.no_zip:
  312. z = db.parent / f'{y}-{m}-{a.cls}.zip'
  313. with zipfile.ZipFile(z, 'w', zipfile.ZIP_DEFLATED, compresslevel=6) as zf:
  314. zf.write(db, db.name)
  315. print(f' + {db.relative_to(out.parent).as_posix()} mdb {db.stat().st_size / 1e6:.1f} MB'
  316. f' · 装入 {time.time() - ts:.0f}s · 合计 {time.time() - t0:.0f}s')
  317. shutil.rmtree(tmp, ignore_errors=True)
  318. return 0
  319. if __name__ == '__main__':
  320. for _s in (sys.stdout, sys.stderr):
  321. try:
  322. _s.reconfigure(errors='replace')
  323. except Exception:
  324. pass
  325. sys.exit(main())