| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130 |
- #!/usr/bin/env python3
- # -*- coding: utf-8 -*-
- r"""油液化验索引摄入: data/raw/<场站>/油样报告/**/*.pdf → <store>/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>\d{2})(?P<m>\d{2})(?P<y>\d{4})_(?P<sid>\d+)_(?P<proj>[^_]+)_(?P<t>\d{1,2})#_(?P<comp>.+)$')
- 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'])
- keep_old = old[~old['file'].isin(touch)]
- 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())
|