windscada_watch_channels_build.py 6.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138
  1. #!/usr/bin/env python3
  2. # -*- coding: utf-8 -*-
  3. r"""油液化验索引摄入: data/raw/<场站>/油样报告/**/*.pdf → <store>/oil_samples_index.parquet (2026-09-11)
  4. ## 为什么有这个脚本
  5. 维护页「数据层 · 油液化验」写着 `python scripts/windscada_watch_channels_build.py (产 oil_samples_index)`,
  6. v0.2.0 包里没有这个文件; 而油样时效轴 (页面上的"油液时效"胶囊) 与 fusion 的"油样窗"都吃这份索引。
  7. ## 源件形态
  8. 化验报告按 `<台号>\<部件>(1台N份)\<报告>.pdf` 落位, 文件名自带全部关键字段:
  9. 17072025_10303681_中广核江苏如东海上项目_1#_主轴后.pdf
  10. └ 日期(DDMMYYYY) └sample_id └项目 └台号 └部件
  11. → turbine=WTG01 · comp=主轴后 · date=2025-07-17 · sample_id=10303681 · file=文件名
  12. ## 与随包件的关系 (2026-09-11 实测)
  13. 随包 oil_samples_index.parquet 506 行里:
  14. · 404 行的源件就在落位目录里, 文件名可解析 → 本脚本**逐字段复现** (校验见 scripts/rebuild_from_raw.py --verify);
  15. · 102 行来自一份**合并报告** `BG-2026-07-YP013 中广核新能源如东海上风电场.pdf` (华标 2026-07 批):
  16. 该 PDF 一份覆盖多台, 台号/部件只在 PDF 表格里而不在文件名; **且这个文件不在现场数据包里**
  17. (全盘搜过)。本脚本因此采取**合并而非重建**语义: 认不出的老行原样保留, 只报出来, 绝不删。
  18. round / date_note 两列是化验批次口径 (SGS2024-2025 / 华标2026-07), 文件名里没有:
  19. 已登记的行**沿用原值**, 新文件才给默认标签并在报告里点名要人工补。
  20. 用法:
  21. python scripts/windscada_watch_channels_build.py --dry-run
  22. python scripts/windscada_watch_channels_build.py
  23. """
  24. from __future__ import annotations
  25. import argparse
  26. import pathlib
  27. import re
  28. import sys
  29. ROOT = pathlib.Path(__file__).resolve().parents[1]
  30. sys.path.insert(0, str(ROOT))
  31. import pandas as pd # noqa: E402
  32. COLS = ['turbine', 'comp', 'date', 'sample_id', 'file', 'round', 'date_note']
  33. _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>.+)$')
  34. DEFAULT_ROUND = '' # 新文件的批次标签留空, 报告里点名要人工补 (不编造批次)
  35. def parse_name(stem: str):
  36. m = _RE.match(stem)
  37. if not m:
  38. return None
  39. g = m.groupdict()
  40. try:
  41. d = pd.Timestamp(year=int(g['y']), month=int(g['m']), day=int(g['d']))
  42. except ValueError:
  43. return None
  44. return dict(turbine=f"WTG{int(g['t']):02d}", comp=g['comp'], date=d.strftime('%Y-%m-%d'),
  45. sample_id=g['sid'])
  46. def main() -> int:
  47. ap = argparse.ArgumentParser()
  48. ap.add_argument('--farm', default=None)
  49. ap.add_argument('--raw', default=None)
  50. ap.add_argument('--dry-run', action='store_true')
  51. a = ap.parse_args()
  52. from src.windscada.config import farm
  53. cfg = farm(a.farm)
  54. src = pathlib.Path(a.raw) if a.raw else pathlib.Path(str(cfg['src_oil']))
  55. store = pathlib.Path(cfg['store'])
  56. out_path = store / 'oil_samples_index.parquet'
  57. print(f'源: {src}\n出: {out_path}\n')
  58. pdfs = sorted(src.rglob('*.pdf'))
  59. if not pdfs:
  60. raise SystemExit(f'源目录没有 pdf: {src}')
  61. rows, unparsed = [], []
  62. for p in pdfs:
  63. got = parse_name(p.stem)
  64. if not got:
  65. unparsed.append(p.name)
  66. continue
  67. got['file'] = p.name
  68. rows.append(got)
  69. new = pd.DataFrame(rows, columns=[c for c in COLS if c != 'round' and c != 'date_note'])
  70. print(f'== 源件 ==\n pdf {len(pdfs)} 个 → 文件名可解析 {len(new)} 个')
  71. if unparsed:
  72. print(f' 文件名不带关键字段的 {len(unparsed)} 个 (需人工登记, 已跳过):')
  73. for n in unparsed[:5]:
  74. print(f' {n}')
  75. old = pd.read_parquet(out_path) if out_path.exists() else pd.DataFrame(columns=COLS)
  76. for c in COLS:
  77. if c not in old:
  78. old[c] = pd.NA
  79. # round/date_note 是批次口径, 文件名给不出 → 已登记的行沿用原值
  80. meta = old.set_index('file')[['round', 'date_note']]
  81. new = new.join(meta, on='file')
  82. filled = int(new['round'].notna().sum())
  83. missing_meta = new[new['round'].isna()]['file'].tolist()
  84. if missing_meta:
  85. new.loc[new['round'].isna(), 'round'] = DEFAULT_ROUND
  86. new.loc[new['date_note'].isna(), 'date_note'] = '文件名解析'
  87. print(f'\n== round/date_note ==\n 沿用原值 {filled} 行; 新件需人工补批次 {len(missing_meta)} 行')
  88. touch = set(new['file'])
  89. # ★产物 = **当前源件的函数** (用户令 2026-09-20「重算要按 data/raw 的最新变化」): 原先 `keep_old` 保留
  90. # "源件已不在盘"的行 ⇒ 撤掉一份报告后旧记录仍留在产物里。现在只留仍在盘的那些件的行, 消失的报出来。
  91. present = {p.name for p in pdfs}
  92. gone = sorted(set(old['file'].dropna().unique()) - present) if len(old) else []
  93. if gone:
  94. _n = int(old['file'].isin(gone).sum())
  95. print(f'\n== [!] {len(gone)} 个油样报告已不在盘 ⇒ 它们带来的 {_n} 行不再产出 (口径: 产物=当前 data/raw 的函数): '
  96. + '、'.join(gone[:5]) + (' …' if len(gone) > 5 else ''))
  97. keep_old = old[~old['file'].isin(touch | set(gone))]
  98. merged = pd.concat([keep_old, new[COLS]], ignore_index=True)
  99. merged['date'] = merged['date'].astype('string')
  100. merged = merged.sort_values(['turbine', 'comp', 'date'], na_position='last').reset_index(drop=True)
  101. d = merged['date'].dropna()
  102. print(f'\n== 结果 ==\n 旧 {len(old)} 行 → 新 {len(merged)} 行 '
  103. f'(本次 {len(new)} 行; 保留未匹配到源件的 {len(keep_old)} 行)')
  104. if len(keep_old):
  105. print(f' 保留的老行来源文件 (源件不在目录里, 例如合并报告): {keep_old["file"].nunique()} 个文件')
  106. for f in keep_old['file'].unique()[:3]:
  107. print(f' {f} ({int((keep_old["file"] == f).sum())} 行)')
  108. print(f' 覆盖 {d.min()} ~ {d.max()}; 台 {merged.turbine.nunique()}; 部件 {merged.comp.nunique()} 种')
  109. if a.dry_run:
  110. print('\n(dry-run, 未写盘)')
  111. return 0
  112. store.mkdir(parents=True, exist_ok=True)
  113. merged.to_parquet(out_path, index=False)
  114. print(f'\n已写 {out_path}')
  115. return 0
  116. if __name__ == '__main__':
  117. sys.exit(main())