scada_source.py 22 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444
  1. # -*- coding: utf-8 -*-
  2. r"""SCADA 原始数据接入层:**CSV 与 MDB 两种形态都认**(用户令 2026-09-19)。
  3. ## 为什么要有这一层
  4. 现场交来的 SCADA 收资有两种形态,此前只有第一种能被摄入:
  5. ```
  6. data/raw/<场站>/scada_10min/WTG01.csv … WTG38.csv ← 已导出 CSV(38 件 / 14.4 GB)
  7. data/raw/<场站>/scada_mdb/<年>年/<月>月/<年>-<月>-<类>.mdb ← 现场年度归档的 Access 库(25年/26年)
  8. ```
  9. `src/windscada/data.py::load_10min()` 原来把路径**写死**成 `<src_10min>/<台>.csv`。用户令之后取数统一走本层:
  10. **CSV 在位就用 CSV**(既有 10 个构建器一直在用的形态,行数已核过),**该台没有 CSV 才回落 MDB**,
  11. 并把"这次是从哪种形态取的"记在返回表的 `attrs['scada_source']` 里,供产物台账/页面说明来路 —— 不静默换源。
  12. ## MDB 那一侧的两个硬事实(决定了本层怎么写)
  13. 1. **单表 ≤255 列 ⇒ 库内按 250 列拆表**(`t1/t2/…`)。新库每行前两列是 `rid`(该台该月内的行号)与
  14. `turbine`(机组标识)⇒ 取一台的多列数据要**跨分片表按 (rid, turbine) 拼**。
  15. 2. **2026-09-19 之前建的老库缺 `turbine` 列**(那时按位置堆行)。10min 老库还能救:源件自带 `WTG` 列,
  16. 本层用它兜底;**1min 老库救不了**(源件没有任何机组标识列,且 38 件表头实测 16 种 ⇒ 位置堆法还会串列)。
  17. 处置:重跑 `scripts/csv_to_mdb.py`(现已写 `turbine` 列 + 按列名对齐),见 `docs/现场收资接入_v0.1.md` §8。
  18. 读 MDB 走 ACE(`Microsoft.ACE.OLEDB.12.0`)+ PowerShell —— 与写库同一套(本机无 pyodbc/mdbtools;
  19. 实测可用路线见 `docs/现场收资接入_v0.1.md` §7)。日期字段在 PowerShell 里显式格式化成
  20. `yyyy-MM-dd HH:mm:ss`(否则会按系统区域设置输出,下游时区/日序解析会踩雷)。
  21. ## 用法
  22. from src.windscada import scada_source as SS
  23. SS.sources('10min', cfg) # 有哪些形态、各多少件(如实,不猜)
  24. SS.turbines('10min', cfg) # 可用的机组名(CSV stem ∪ MDB 里的 turbine)
  25. SS.load('10min', 'WTG01', columns=[...], cfg=cfg) # → DataFrame(列名原样,ts 未解析)
  26. SS.source_of('10min', 'WTG01', cfg) # → 'csv' / 'mdb'
  27. 纪律: 两种形态都取不到时**响亮 raise**(不返回空表 —— 空表会被下游当成"这台没数据")。
  28. """
  29. from __future__ import annotations
  30. import functools
  31. import pathlib
  32. import re
  33. import subprocess
  34. import tempfile
  35. CLS = {'10min': ('src_10min', 'scada_10min'), '1min': ('src_1min', 'scada_1min')}
  36. TURBINE_COL = 'turbine' # 新库的机组列(csv_to_mdb.py 2026-09-19 起写)
  37. LEGACY_TURBINE_COL = 'WTG' # 老库兜底:10min 源件自带 WTG 列
  38. # ★同台多件 (用户令 2026-09-20「真正把新月份的 10min 件吃进来」): 现场后来补的导出件命名与主件不同
  39. # (实测 `WTG01-B2.csv` 856 列 / 2026-08 数据),而原实现只认 `<台号>.csv` ⇒ 这批数据**从未被读过**,
  40. # 页面时间窗自然停在旧月份。现在本台取数 = 主件 + 同台补充件,按时间戳合并去重(主件优先)。
  41. SUPPLEMENT_GLOBS = ('{t}-*.csv', '{t}_*.csv')
  42. # 新旧导出的列名差: 新件把"类型前缀族"去掉了(`din_wtc_HydLevel_timeon` → `wtc_HydLevel_timeon`)。
  43. # 实测 WTG01 老件 598 列里 597 列能这样对上(唯一对不上的是 `WTG`,新件用 `StationId`)。
  44. PREFIXES = ('din_', 'dot_', 'prs_', 'tmp_', 'grd_', 'tur_', 'cnt_', 'flg_', 'int_', 'evt_')
  45. TIME_CANDIDATES = ('TimeStamp', 'timestamp', 'Time', 'time', 'time_stamp', 'occur_time')
  46. _PS = r'''
  47. $ErrorActionPreference = "Continue"
  48. $c = New-Object -ComObject ADODB.Connection
  49. $c.Open("Provider=Microsoft.ACE.OLEDB.12.0;Data Source=@@MDB@@;")
  50. $turb = "@@TURB@@"
  51. $want = @(@@WANT@@)
  52. # 表名清单: 用 ADO 的 OpenSchema(20=adSchemaTables)。★`$c.GetSchema("Tables")` 会报
  53. # "Arguments are of the wrong type..."(那是 .NET 的方法名), 实测 2026-09-19 踩过。
  54. $tabs = @()
  55. $rs0 = $c.OpenSchema(20)
  56. while (-not $rs0.EOF) {
  57. if ($rs0.Fields.Item("TABLE_TYPE").Value -eq "TABLE") { $tabs += $rs0.Fields.Item("TABLE_NAME").Value }
  58. $rs0.MoveNext()
  59. }
  60. $rs0.Close()
  61. foreach ($t in $tabs) {
  62. $rs = $c.Execute("SELECT TOP 1 * FROM [$t]")
  63. $names = @()
  64. foreach ($f in $rs.Fields) { $names += $f.Name }
  65. $rs.Close()
  66. # 机组列: 先试 turbine(新库), 再试 WTG(老 10min 库兜底); 两个都没有 ⇒ 这库按台取不到数
  67. $tcol = $null
  68. foreach ($p in @("turbine", "WTG")) {
  69. try { $null = $c.Execute("SELECT TOP 1 [$p] FROM [$t]"); $tcol = $p; break } catch { }
  70. }
  71. $turbines = @()
  72. if ($tcol) {
  73. $rs1 = $c.Execute("SELECT DISTINCT [$tcol] FROM [$t] WHERE [$tcol] IS NOT NULL")
  74. while (-not $rs1.EOF) { $turbines += ("" + $rs1.Fields.Item(0).Value); $rs1.MoveNext() }
  75. $rs1.Close()
  76. }
  77. Write-Output ("#SCHEMA " + $t + "|" + ($names -join ",") + "|" + ($turbines -join ";"))
  78. if (-not $tcol -or $turb -eq "") { continue }
  79. $keep = @()
  80. foreach ($w in $want) { if (($names -contains $w) -and $w -ne "rid" -and $w -ne $tcol) { $keep += $w } }
  81. if ($want.Count -eq 0) { # 不给白名单 = 全列(与 CSV 形态 columns=None 的语义一致)
  82. foreach ($w in $names) { if ($w -ne "rid" -and $w -ne $tcol) { $keep += $w } }
  83. }
  84. $sel = @("rid", $tcol) + $keep
  85. $collist = ($sel | ForEach-Object { "[$_]" }) -join ","
  86. Write-Output ("#TABLE " + $t + "|" + ($sel -join ","))
  87. $rs2 = $c.Execute("SELECT $collist FROM [$t] WHERE [$tcol] = '$turb'")
  88. while (-not $rs2.EOF) {
  89. $vals = @()
  90. for ($i = 0; $i -lt $rs2.Fields.Count; $i++) {
  91. $v = $rs2.Fields.Item($i).Value
  92. if ($v -is [datetime]) { $vals += $v.ToString("yyyy-MM-dd HH:mm:ss") }
  93. elseif ($null -eq $v) { $vals += "" }
  94. else { $vals += ("" + $v) }
  95. }
  96. Write-Output ($vals -join "`t")
  97. $rs2.MoveNext()
  98. }
  99. $rs2.Close()
  100. }
  101. $c.Close()
  102. '''
  103. def dirs_of(cls: str, cfg=None) -> tuple[pathlib.Path, pathlib.Path]:
  104. """→ (CSV 目录, MDB 根)。两者都不必存在。"""
  105. cfg = _cfg(cfg)
  106. key, sub = CLS[cls]
  107. csv_dir = pathlib.Path(str(cfg.get(key) or (pathlib.Path(str(cfg['raw_station_dir'])) / sub)))
  108. mdb_root = pathlib.Path(str(cfg['raw_station_dir'])) / 'scada_mdb'
  109. return csv_dir, mdb_root
  110. def _cfg(cfg=None):
  111. if cfg is not None:
  112. return cfg
  113. from src.windscada.config import farm # config 属后端模块(P4),暂留原路径
  114. return farm()
  115. def sources(cls: str = '10min', cfg=None) -> dict:
  116. """盘上有哪些形态(件数)—— 供页面/自检如实说明。"""
  117. csv_dir, mdb_root = dirs_of(cls, cfg)
  118. csvs = sorted(csv_dir.glob('*.csv')) if csv_dir.is_dir() else []
  119. mdbs = sorted(mdb_root.rglob(f'*-{cls}.mdb')) if mdb_root.is_dir() else []
  120. zips = sorted(mdb_root.rglob(f'*-{cls}.zip')) if mdb_root.is_dir() else []
  121. return dict(csv_dir=csv_dir, mdb_root=mdb_root, n_csv=len(csvs), n_mdb=len(mdbs), n_zip=len(zips),
  122. csvs=csvs, mdbs=mdbs, zips=zips,
  123. kind=('csv' if csvs else ('mdb' if (mdbs or zips) else None)))
  124. @functools.lru_cache(maxsize=8)
  125. def schema_and_turbines(mdb: str) -> tuple:
  126. """一个 MDB 的 (分片表 → 列名, 分片表 → 该列里的机组值)。老库缺 turbine 列时退 WTG 列。"""
  127. lines = _ps(_PS.replace('@@MDB@@', str(mdb)).replace('@@TURB@@', '').replace('@@WANT@@', ''))
  128. cols, turbs = {}, {}
  129. for ln in lines:
  130. if ln.startswith('#SCHEMA '):
  131. body = ln[8:]
  132. t, rest = body.split('|', 1)
  133. names, tv = rest.split('|', 1)
  134. cols[t] = [x for x in names.split(',') if x]
  135. turbs[t] = tuple(sorted({x for x in tv.split(';') if x}))
  136. return cols, turbs
  137. def _ps(script: str) -> list[str]:
  138. """跑一段 PowerShell(ACE 只能经 COM 用)→ stdout 行。
  139. ★退出码不可尽信(2026-09-19 实测): ACE 在**进程退出时清理 COM** 偶发 `0xC0000005`(访问违例),
  140. 此时 stderr 为空、stdout 完整(脚本里每一行都打出来了)。原先 `rc != 0 ⇒ raise` 会把这种"数据都拿到了"
  141. 的情况报成读库失败。判据改成**看输出**: 有输出且退出码是那个崩溃码 ⇒ 收下并提示一次; 否则才 raise。
  142. """
  143. p = pathlib.Path(tempfile.gettempdir()) / '_scada_source.ps1'
  144. p.write_text(script + '\nexit 0\n', encoding='utf-8-sig')
  145. r = subprocess.run(['powershell', '-NoProfile', '-ExecutionPolicy', 'Bypass', '-File', str(p)],
  146. capture_output=True, text=True, errors='replace', timeout=3600)
  147. out = [x for x in (r.stdout or '').splitlines() if x.strip()]
  148. crash = r.returncode in (-1073741819, 3221225477)
  149. if r.returncode != 0 and not (crash and out):
  150. raise RuntimeError((r.stderr or r.stdout or '').strip()[:400] or 'PowerShell 读 MDB 失败')
  151. if crash and out and 'ace_exit_crash' not in _WARNED:
  152. _WARNED.add('ace_exit_crash')
  153. print(' [i] ACE 退出时 COM 清理崩溃 (0xC0000005) —— 输出完整, 已按成功处理', flush=True)
  154. return out
  155. _WARNED: set = set()
  156. def mdb_turbines(cls: str = '10min', cfg=None, limit: int = 3) -> set:
  157. """MDB 侧的可用机组名(取前几个库的并集 —— 逐库扫全量太慢,够用)。"""
  158. s = sources(cls, cfg)
  159. got = set()
  160. for m in s['mdbs'][:limit]:
  161. _cols, turbs = schema_and_turbines(str(m))
  162. for t in turbs.values():
  163. got |= set(t)
  164. return got
  165. def turbines(cls: str = '10min', cfg=None) -> list[str]:
  166. """可用机组名 = CSV 文件名 stem ∪ MDB 里的 turbine 值。"""
  167. s = sources(cls, cfg)
  168. got = {f.stem for f in s['csvs']} | mdb_turbines(cls, cfg)
  169. return sorted(got)
  170. def source_of(cls: str, turbine: str, cfg=None) -> str | None:
  171. """这一台这次会从哪种形态取(csv 优先)。"""
  172. csv_dir, _ = dirs_of(cls, cfg)
  173. if (csv_dir / f'{turbine}.csv').is_file():
  174. return 'csv'
  175. return 'mdb' if turbine in mdb_turbines(cls, cfg, limit=1) else None
  176. def files_of(cls: str, turbine: str, cfg=None) -> list:
  177. """本台这一次要读的 CSV 件: **主件 `<台号>.csv` 在前**, 后面是同台补充件(`<台号>-*.csv` 等)。
  178. 用户令 2026-09-20: 现场补的件命名与主件不同(实测 `WTG01-B2.csv`),只认 `<台号>.csv` 会让
  179. "新来的月份"静默漏读。这里如实列出这次会读的件(顺序=优先级),供取数/台账/页面说明来路。
  180. """
  181. csv_dir, _ = dirs_of(cls, cfg)
  182. out = []
  183. pri = csv_dir / f'{turbine}.csv'
  184. if pri.is_file():
  185. out.append(pri)
  186. for pat in SUPPLEMENT_GLOBS:
  187. out += sorted(csv_dir.glob(pat.format(t=turbine)))
  188. seen, res = set(), []
  189. for p in out:
  190. if p.name not in seen:
  191. seen.add(p.name)
  192. res.append(p)
  193. return res
  194. def _align_columns(df, want):
  195. """按列名对齐新旧导出: 想要的列不在、但"去掉类型前缀"后在同表里 ⇒ 认回来。
  196. 返回 (df, 认回的列数)。实测: 老 598 列 → 新 856 列, 597/598 靠这条规则对上。
  197. """
  198. if not want or df is None or not len(df.columns):
  199. return df, 0
  200. have = set(df.columns)
  201. back = {}
  202. for c in want:
  203. if c in have:
  204. continue
  205. for pre in PREFIXES:
  206. if c.startswith(pre) and c[len(pre):] in have:
  207. back[c[len(pre):]] = c
  208. break
  209. return (df.rename(columns=back) if back else df), len(back)
  210. def _time_col(frames, first_cols):
  211. """时间列: 先按已知名找(合并后的表列名可能来自任一形态),再退回主件的首列。"""
  212. cols = set(first_cols)
  213. for f in frames:
  214. cols |= set(f.columns)
  215. for c in TIME_CANDIDATES:
  216. if c in cols:
  217. return c
  218. return first_cols[0] if first_cols else None
  219. def load(cls: str, turbine: str, columns=None, cfg=None, month: str | None = None):
  220. """取一台的数据(**CSV 优先**;该台没有 CSV 时回落 MDB)。
  221. columns = 白名单(两形态都按"实际有的列"取交集,缺列不报错 —— 与 `data.load_10min` 的既有口径一致)。
  222. month = 'YYYY-MM',只取该月(MDB 形态天然是按月的;CSV 形态按首列前缀过滤)。
  223. ★同台多件 (用户令 2026-09-20): `<台号>.csv` + `<台号>-*.csv` 一起读, 按列名对齐 + 按时间戳去重
  224. (**主件优先**,重叠条数如实写在 `attrs['scada_overlap']` 里)。
  225. """
  226. import pandas as pd
  227. fps = files_of(cls, turbine, cfg)
  228. if fps:
  229. frames, per = [], {}
  230. for i, fp in enumerate(fps):
  231. have = set(pd.read_csv(fp, nrows=0).columns)
  232. take = None
  233. if columns:
  234. take = []
  235. for c in columns:
  236. if c in have:
  237. take.append(c)
  238. continue
  239. for pre in PREFIXES: # 新导出把前缀去了 ⇒ 按去前缀名取列
  240. if c.startswith(pre) and c[len(pre):] in have:
  241. take.append(c[len(pre):])
  242. break
  243. d = pd.read_csv(fp, usecols=take, engine='pyarrow') if take else \
  244. (pd.read_csv(fp, usecols=lambda c: c in have, engine='pyarrow') if columns else
  245. pd.read_csv(fp, engine='pyarrow'))
  246. d, n_back = _align_columns(d, columns or [])
  247. first_cols = list(d.columns)
  248. d['_scada_rank'] = i # i=0 是主件 ⇒ 去重时优先
  249. frames.append(d)
  250. per[fp.name] = len(d)
  251. ts = _time_col(frames, first_cols)
  252. d = pd.concat(frames, ignore_index=True, sort=False)
  253. if month and ts:
  254. d = d[d[ts].astype(str).str.startswith(month)]
  255. overlap = 0
  256. if len(frames) > 1 and ts:
  257. before = len(d)
  258. d = d.sort_values(['_scada_rank'], kind='stable').drop_duplicates(subset=[ts], keep='first')
  259. overlap = before - len(d)
  260. d = d.drop(columns=['_scada_rank'])
  261. if ts in d.columns:
  262. d = d.sort_values(ts, kind='stable').reset_index(drop=True)
  263. else:
  264. d = d.reset_index(drop=True)
  265. if len(fps) > 1:
  266. detail = '、'.join(f'{k} {v} 行' for k, v in per.items())
  267. print(f' [i] {turbine}: 同台 {len(fps)} 件合并读入({detail})'
  268. + (f';同时间戳重叠 {overlap} 行按"主件优先"去重' if overlap else ';无重叠'),
  269. flush=True)
  270. d.attrs['scada_source'] = 'csv'
  271. d.attrs['scada_files'] = [f.name for f in fps]
  272. d.attrs['scada_rows'] = per
  273. d.attrs['scada_overlap'] = overlap
  274. return d.reset_index(drop=True)
  275. if source_of(cls, turbine, cfg) != 'mdb':
  276. raise RuntimeError(
  277. f'{cls} 取数失败: {turbine} 既没有 CSV ({dirs_of(cls, cfg)[0] / f"{turbine}.csv"}),也没有含该台的 MDB '
  278. f'({dirs_of(cls, cfg)[1]})。放原始件到 data/raw/<场站>/ 后重跑 scripts/rebuild_all.py。')
  279. d = read_mdb(cls, turbine, columns, cfg, month)
  280. d.attrs['scada_source'] = 'mdb'
  281. return d
  282. def read_mdb(cls: str, turbine: str, columns=None, cfg=None, month: str | None = None):
  283. """从月库里取一台:逐分片表按 (rid, turbine) 拼列。"""
  284. import pandas as pd
  285. from app_ETL.app_ETL_guanlan.mdb_names import safe_name
  286. s = sources(cls, cfg)
  287. mdbs = [m for m in s['mdbs'] if not month or m.stem.startswith(month)]
  288. if not mdbs:
  289. raise RuntimeError(f'{cls} 取数失败: scada_mdb 下没有 {month or "任何"} 的 .mdb')
  290. # 契约列名 → 库内列名(Access 列名不能含 `.` `!` `[]` 等, 建库时消毒过; 两边同一套规则)
  291. cand = {c: safe_name(c) for c in (columns or [])}
  292. want_stored = sorted(set(cand.values()))
  293. frames, missing_col = [], None
  294. for m in mdbs:
  295. want = want_stored
  296. script = (_PS.replace('@@MDB@@', str(m)).replace('@@TURB@@', str(turbine))
  297. .replace('@@WANT@@', ','.join('"' + w + '"' for w in want)))
  298. cur_t, cur_cols, buf, seen_cols = None, [], [], set()
  299. def _flush():
  300. if cur_t and buf:
  301. frames.append(_frame(cur_cols, buf))
  302. for ln in _ps(script):
  303. if ln.startswith('#SCHEMA '):
  304. continue
  305. if ln.startswith('#TABLE '):
  306. _flush()
  307. body = ln[7:]
  308. cur_t, cols = body.split('|', 1)
  309. if not re.fullmatch(r't\d+', cur_t): # Access 自建的杂表(如"名称自动更正保存失败")
  310. cur_t, cur_cols, buf = None, [], []
  311. continue
  312. cur_cols, buf = cols.split(','), []
  313. seen_cols |= set(cur_cols)
  314. else:
  315. buf.append(ln)
  316. _flush()
  317. if columns and not any(c in seen_cols for c in columns):
  318. missing_col = [c for c in columns if c not in seen_cols][:3]
  319. # 列位错位检测: 用 **#SCHEMA**(每个分片表的列清单, 与本次取了哪些列无关)
  320. dup = set()
  321. for m2 in mdbs:
  322. sc = {}
  323. for ln in _ps(_PS.replace('@@MDB@@', str(m2)).replace('@@TURB@@', '').replace('@@WANT@@', '')):
  324. if ln.startswith('#SCHEMA '):
  325. body = ln[8:]
  326. t2, rest = body.split('|', 1)
  327. if re.fullmatch(r't\d+', t2):
  328. sc[t2] = set(x for x in rest.split('|', 1)[0].split(',') if x)
  329. ks = list(sc)
  330. for i in range(len(ks)):
  331. for j in range(i + 1, len(ks)):
  332. dup |= (sc[ks[i]] & sc[ks[j]]) - {'rid', TURBINE_COL}
  333. if dup:
  334. break
  335. if not frames:
  336. raise RuntimeError(
  337. f'{cls} 取数失败: {turbine} 在 {[m.name for m in mdbs]} 里取不到行 —— 老库可能没有 turbine 列'
  338. f'(源件也没有 WTG 列)⇒ 重跑 scripts/csv_to_mdb.py 建新库(现已写 turbine 列)')
  339. frames = [f.rename(columns={LEGACY_TURBINE_COL: TURBINE_COL}) if
  340. (LEGACY_TURBINE_COL in f.columns and TURBINE_COL not in f.columns) else f for f in frames]
  341. # 库内列名 → 契约列名(撞名时建库侧写过 `_1` 后缀, 这里按前缀认回来)
  342. if cand:
  343. for f in frames:
  344. back, used = {}, set()
  345. for orig, st in cand.items():
  346. if st in f.columns and st not in used:
  347. back[st], _ = orig, used.add(st)
  348. continue
  349. stem = st[:56]
  350. for c2 in f.columns:
  351. if c2 not in used and c2.startswith(stem) and c2[len(stem):].startswith('_'):
  352. back[c2], _ = orig, used.add(c2)
  353. break
  354. f.rename(columns=back, inplace=True)
  355. # ★2026-09-19 逮到的**老库列位错位**: 2026-09-19 之前建的 10min 库, 分片表表头与数据行差一列
  356. # (旧 build_chunks 表头写 names[lo:hi+1] 而行写 vals[lo:hi]) ⇒ 第 250 列之后的**列名整体错位一格**,
  357. # 且相邻分片表会共用同一个列名(实测 t1 末列 == t2 首列 == flg_wtc_ScReToOp_endvalue)。这种库读出来的
  358. # 是"名字对不上数"的数据 —— 比缺数据危险, 故**响亮拒绝**, 让人重建而不是将就着算。
  359. if dup:
  360. raise RuntimeError(
  361. f'{cls} 取数失败: {[m.name for m in mdbs]} 是**列位错位的老库**(分片表共用列名 {sorted(dup)[:3]}…) '
  362. f'—— 2026-09-19 前的建库脚本表头与数据行差一列, 第 250 列之后的列名整体错位。'
  363. f'处置: python scripts/csv_to_mdb.py --month <YYYY-MM> --class {cls} 重建该月库; '
  364. f'或放 CSV 到 data/raw/<场站>/{CLS[cls][1]}/ 直接用 CSV 形态。')
  365. d = frames[0]
  366. for f in frames[1:]:
  367. d = d.merge(f, on=['rid', TURBINE_COL], how='outer')
  368. # 行序回到**源件顺序**: rid 在库里是文本 ⇒ 直接合并得到的是 '1','10','100' 这种字典序
  369. # (2026-09-19 实逮: 与 CSV 逐值对拍对不上, 数据没错、是行序被按字符串排了)
  370. d['rid'] = pd.to_numeric(d['rid'], errors='coerce')
  371. d = d.sort_values([TURBINE_COL, 'rid'], kind='stable').reset_index(drop=True)
  372. if missing_col:
  373. d.attrs['mdb_missing'] = missing_col
  374. return d
  375. def _frame(cols: list[str], rows: list[str]):
  376. """PowerShell 落的制表符分隔文本 → DataFrame。
  377. 数值列转成数值: MDB 里存的是**数值类型**, ACE 读出来 `100.0` 会变成 `100` —— 若原样留成字符串,
  378. 下游按 CSV 口径写的老构建器会拿到 object 列 (`+`/`mean` 全崩)。这里逐列试探: 非空值全能解析成数字
  379. 就转数值, 否则留字符串 (台号/状态字这类)。
  380. """
  381. import pandas as pd
  382. ncol = len(cols)
  383. recs = []
  384. for ln in rows:
  385. vals = ln.split('\t')
  386. recs.append((vals + [''] * ncol)[:ncol])
  387. d = pd.DataFrame(recs, columns=cols)
  388. for c in d.columns:
  389. if c in ('rid', TURBINE_COL, LEGACY_TURBINE_COL):
  390. continue
  391. s = d[c].replace('', None)
  392. num = pd.to_numeric(s, errors='coerce')
  393. if num.notna().sum() == s.notna().sum() and s.notna().any():
  394. d[c] = num
  395. return d