#!/usr/bin/env python3 # -*- coding: utf-8 -*- r"""油液化验索引摄入: data/raw/<场站>/油样报告/**/*.pdf → /oil_samples_index.parquet (2026-09-11) ## 为什么有这个脚本 维护页「数据层 · 油液化验」写着 `python scripts/windscada_watch_channels_build.py (产 oil_samples_index)`, v0.2.0 包里没有这个文件; 而油样时效轴 (页面上的"油液时效"胶囊) 与 fusion 的"油样窗"都吃这份索引。 ## 源件形态 化验报告按 `<台号>\<部件>(1台N份)\<报告>.pdf` 落位, 文件名自带全部关键字段: 17072025_10303681_中广核江苏如东海上项目_1#_主轴后.pdf └ 日期(DDMMYYYY) └sample_id └项目 └台号 └部件 → turbine=WTG01 · comp=主轴后 · date=2025-07-17 · sample_id=10303681 · file=文件名 ## 与随包件的关系 (2026-09-11 实测) 随包 oil_samples_index.parquet 506 行里: · 404 行的源件就在落位目录里, 文件名可解析 → 本脚本**逐字段复现** (校验见 scripts/rebuild_from_raw.py --verify); · 102 行来自一份**合并报告** `BG-2026-07-YP013 中广核新能源如东海上风电场.pdf` (华标 2026-07 批): 该 PDF 一份覆盖多台, 台号/部件只在 PDF 表格里而不在文件名; **且这个文件不在现场数据包里** (全盘搜过)。本脚本因此采取**合并而非重建**语义: 认不出的老行原样保留, 只报出来, 绝不删。 round / date_note 两列是化验批次口径 (SGS2024-2025 / 华标2026-07), 文件名里没有: 已登记的行**沿用原值**, 新文件才给默认标签并在报告里点名要人工补。 用法: python scripts/windscada_watch_channels_build.py --dry-run python scripts/windscada_watch_channels_build.py """ from __future__ import annotations import argparse import pathlib import re import sys ROOT = pathlib.Path(__file__).resolve().parents[1] sys.path.insert(0, str(ROOT)) import pandas as pd # noqa: E402 COLS = ['turbine', 'comp', 'date', 'sample_id', 'file', 'round', 'date_note'] _RE = re.compile(r'^(?P\d{2})(?P\d{2})(?P\d{4})_(?P\d+)_(?P[^_]+)_(?P\d{1,2})#_(?P.+)$') DEFAULT_ROUND = '' # 新文件的批次标签留空, 报告里点名要人工补 (不编造批次) def parse_name(stem: str): m = _RE.match(stem) if not m: return None g = m.groupdict() try: d = pd.Timestamp(year=int(g['y']), month=int(g['m']), day=int(g['d'])) except ValueError: return None return dict(turbine=f"WTG{int(g['t']):02d}", comp=g['comp'], date=d.strftime('%Y-%m-%d'), sample_id=g['sid']) def main() -> int: ap = argparse.ArgumentParser() ap.add_argument('--farm', default=None) ap.add_argument('--raw', default=None) ap.add_argument('--dry-run', action='store_true') a = ap.parse_args() from src.windscada.config import farm cfg = farm(a.farm) src = pathlib.Path(a.raw) if a.raw else pathlib.Path(str(cfg['src_oil'])) store = pathlib.Path(cfg['store']) out_path = store / 'oil_samples_index.parquet' print(f'源: {src}\n出: {out_path}\n') pdfs = sorted(src.rglob('*.pdf')) if not pdfs: raise SystemExit(f'源目录没有 pdf: {src}') rows, unparsed = [], [] for p in pdfs: got = parse_name(p.stem) if not got: unparsed.append(p.name) continue got['file'] = p.name rows.append(got) new = pd.DataFrame(rows, columns=[c for c in COLS if c != 'round' and c != 'date_note']) print(f'== 源件 ==\n pdf {len(pdfs)} 个 → 文件名可解析 {len(new)} 个') if unparsed: print(f' 文件名不带关键字段的 {len(unparsed)} 个 (需人工登记, 已跳过):') for n in unparsed[:5]: print(f' {n}') old = pd.read_parquet(out_path) if out_path.exists() else pd.DataFrame(columns=COLS) for c in COLS: if c not in old: old[c] = pd.NA # round/date_note 是批次口径, 文件名给不出 → 已登记的行沿用原值 meta = old.set_index('file')[['round', 'date_note']] new = new.join(meta, on='file') filled = int(new['round'].notna().sum()) missing_meta = new[new['round'].isna()]['file'].tolist() if missing_meta: new.loc[new['round'].isna(), 'round'] = DEFAULT_ROUND new.loc[new['date_note'].isna(), 'date_note'] = '文件名解析' print(f'\n== round/date_note ==\n 沿用原值 {filled} 行; 新件需人工补批次 {len(missing_meta)} 行') touch = set(new['file']) # ★产物 = **当前源件的函数** (用户令 2026-09-20「重算要按 data/raw 的最新变化」): 原先 `keep_old` 保留 # "源件已不在盘"的行 ⇒ 撤掉一份报告后旧记录仍留在产物里。现在只留仍在盘的那些件的行, 消失的报出来。 present = {p.name for p in pdfs} gone = sorted(set(old['file'].dropna().unique()) - present) if len(old) else [] if gone: _n = int(old['file'].isin(gone).sum()) print(f'\n== [!] {len(gone)} 个油样报告已不在盘 ⇒ 它们带来的 {_n} 行不再产出 (口径: 产物=当前 data/raw 的函数): ' + '、'.join(gone[:5]) + (' …' if len(gone) > 5 else '')) keep_old = old[~old['file'].isin(touch | set(gone))] merged = pd.concat([keep_old, new[COLS]], ignore_index=True) merged['date'] = merged['date'].astype('string') merged = merged.sort_values(['turbine', 'comp', 'date'], na_position='last').reset_index(drop=True) d = merged['date'].dropna() print(f'\n== 结果 ==\n 旧 {len(old)} 行 → 新 {len(merged)} 行 ' f'(本次 {len(new)} 行; 保留未匹配到源件的 {len(keep_old)} 行)') if len(keep_old): print(f' 保留的老行来源文件 (源件不在目录里, 例如合并报告): {keep_old["file"].nunique()} 个文件') for f in keep_old['file'].unique()[:3]: print(f' {f} ({int((keep_old["file"] == f).sum())} 行)') print(f' 覆盖 {d.min()} ~ {d.max()}; 台 {merged.turbine.nunique()}; 部件 {merged.comp.nunique()} 种') if a.dry_run: print('\n(dry-run, 未写盘)') return 0 store.mkdir(parents=True, exist_ok=True) merged.to_parquet(out_path, index=False) print(f'\n已写 {out_path}') return 0 if __name__ == '__main__': sys.exit(main())