#!/usr/bin/env python3 # -*- coding: utf-8 -*- r"""产物清点与「输入 ↔ 产物」呼应校验 (2026-09-16 用户令 1、2)。 用户令: 1. 基于观澜源码, 认真检查、梳理**所有产物**(产物类型/功用/输出路径)确保无遗漏; 输出路径不统一的改源码统一; 完善到 `<安装目录>\docs\系统设计说明.md`。 2. 确保 `data/raw` 下的输入数据**可用于重算**, 且**产物与输入数据呼应**。 本脚本是三件事的**单一实现** (别在文档里手抄第二份, 会飘): · 清点: 每个产物仓的件数/大小/类型, 以及**逐件来源**(raw-derived / shipped, 读 `_provenance.json` 与 `_derived_manifest.json`); · 呼应: 对每类输入算它的**数据跨度/条数**, 与它喂出来的产物**逐项对拍** (产物跨度必须落在输入跨度内, 且输入更新时产物不该明显落后); · 写档: `--write-doc` 把清单块写进 `docs\系统设计说明.md` 的 ` … ` 之间 (文档其余部分手工维护)。 用法: python scripts/inventory_products.py # 打清单 + 呼应结论 (人看) python scripts/inventory_products.py --check # 只做呼应校验 (有问题 → 退出码 5) python scripts/inventory_products.py --write-doc # 把清单块写进 docs/系统设计说明.md python scripts/inventory_products.py --json # 机器可读 """ from __future__ import annotations import argparse import datetime as dt import json import pathlib import re import sys from collections import OrderedDict ROOT = pathlib.Path(__file__).resolve().parents[1] sys.path.insert(0, str(ROOT)) from src.console import soft # noqa: E402 from src import paths as P # noqa: E402 soft() DOC = P.DOCS / '系统设计说明.md' MARK_B, MARK_E = '' # ── 产物仓: 路径助手 → (仓名, 功用, 谁生成, 谁消费) ─────────────────────────────── # ★路径一律用 src/paths.py 的助手 (唯一真源); 这里写的是**相对安装根**的形式, 便于跨平台显示。 STORES = OrderedDict([ ('windscada', dict(helper='P.store()', purpose='L0 标准仓: SCADA/台账/派生分析的全部 parquet (页面主取数处)', producer='rebuild_from_raw.py (三门台账) + --scada (10 个构建器) + windscada_monthly_build.py', consumer='scripts/windscada_serve.py 各视图 · src/windscada/taxonomy · subsys/fusion')), ('ontology', dict(helper='P.ont()', purpose='本体对象库: 码表/手册/工单展开/失效树 + 检索索引 + 实机参数', producer='python -m src.ontology.kb_ingest → populate → chain_ingest → trend_ingest → retrieval.build', consumer='脚本 windscada_serve.py 本体页/问答 · scripts/guanlan_facts_contract.py')), ('windcms', dict(helper='P.cms()', purpose='CMS 振动诊断产物: 状态评估报告/逐台页/工作台页/知识库 + 厂家报告转录', producer='scripts/windcms.py report/kb · scripts/vib_reports_build.py', consumer='自服务 :18020 (每请求现读) · taxonomy.system_matrix (转录设备状态)')), ('m5_cms_tcm', dict(helper='P.m5()', purpose='振动线出件与窗级分析: handoff 接口 + 窗索引/谱库 + TCM 兼容件', producer='scripts/vib_raw_build.py (窗索引/谱) · 振动线出件 (handoff, 随包快照)', consumer='src/windscada/subsys/fusion.py · src/windcms/data.py · scripts/windcms.py report')), ('tcm_compatible_replay', dict(helper='P.tcm_replay()', purpose='TCM 兼容链回放资产 (模型表/掩码阈值/裁决记录)', producer='随包快照 (无生成端)', consumer='src/windcms/report*.py · config.mask_thresholds')), ('sop', dict(helper='P.sop()', purpose='SOP 中间件/评审/台账与事实契约底稿', producer='随包快照 (无生成端)', consumer='scripts/guanlan_facts_contract.py · 门户结论段')), ('guanlan', dict(helper='P.guanlan()', purpose='事实契约与对外派生 (可上云面孔)', producer='scripts/guanlan_facts_contract.py', consumer='门户 #findings · /api/facts')), ('pitch', dict(helper='P.pitch()', purpose='变桨侧派生件 (零位/日粒度)', producer='随包快照 + rebuild_from_raw --scada', consumer='脚本 windscada_serve.py 变桨面')), ('paradigm_r1', dict(helper='(无助手: P.out_root()/paradigm_r1)', purpose='范式实验件 (E3/E5/E8 底稿, 事实契约输入)', producer='随包快照 (无生成端)', consumer='scripts/guanlan_facts_contract.py')), ]) # ── 输入类 → (放什么, 喂哪些产物, 跨度判据) ──────────────────────────────────────── # span: 'csv_col' = 逐台 CSV 的首/末数据行; 'fname' = 从文件名年份/日期推; 'dir' = 目录名 (年/月) INPUTS = OrderedDict([ ('scada_10min', dict(what='SCADA 10min 导出 (逐台 WTG01.csv…WTG38.csv)', span='csv_col', time_col='Time', feeds=['windscada/powercurve*.parquet', 'windscada/loss_monthly.parquet', 'windscada/temp_monthly.parquet', 'windscada/temp_bins.parquet', 'windscada/stop_events.parquet', 'windscada/yaw_daily.parquet', 'windscada/hydraulic_accum.parquet', 'windscada/thermal_chain.parquet', 'windscada/curve_*.parquet', 'windscada/control_*.parquet', 'windscada/system_aux.parquet'])), ('scada_1min', dict(what='1min 导出 (逐台)', span='csv_col', time_col='occur_time', feeds=['windscada/watch_channels_monthly.parquet (watch 月度)'])), ('故障报警', dict(what='报警事件导出 (SpreadsheetML *.xls)', span='fname_year', feeds=['windscada/alarms.parquet'])), ('风机故障记录', dict(what='检修工单台账 (*.xls/xlsx)', span='fname_year', feeds=['windscada/workorders.parquet', 'windscada/mblub_monthly.parquet'])), ('油样报告', dict(what='油液化验报告 (*.pdf)', span='fname_date', feeds=['windscada/oil_samples_index.parquet'])), ('windcms', dict(what='CMS 原始测量导出 (Brande TCM *_decode.json) + 厂家评估报告', span='dir', feeds=['m5_cms_tcm/windows/<窗>/{index.parquet,spectra/*}', 'windcms/厂家报告提取_*.json'])), ('m5_cms_tcm', dict(what='振动线出件与 TCM 侧报告', span='none', feeds=['m5_cms_tcm/handoff_vibration_v2.json (现场正本优先)', 'windcms/报告_TCM传动链振动分析_*.md'])), ]) # 机理层 (不在场站目录下) TECH_INPUT = ('西门子4.0技术资料', '厂商技术资料 (故障处理手册/维护 WI/图纸/对译表)', ['ontology/objects.json', 'ontology/turbine_params.parquet', 'ontology/retrieval_index.json']) def _rows_of(path: pathlib.Path, probe: int = 1): """只读首/末若干行, 不整表载入 (scada csv 单文件几百 MB)。""" try: with open(path, 'rb') as f: head = f.readline() f.seek(max(0, path.stat().st_size - 65536)) tail = f.read().splitlines() return head, (tail[-1] if tail else b'') except Exception: return b'', b'' def input_span(station: pathlib.Path, sub: str, spec: dict): """→ (起, 止, 件数, 说明, 粒度)。粒度 ∈ {'日','月','年'}: 只有'日/月'才参与"产物落后"判定 —— 按文件名年份推出来的跨度天然是**年粒度**, 拿它当"输入到 2026-12"会造出假缺口 (本脚本第一版就这么 误报过 2 条: 报警/工单"落后 5 个月"), 故显式带粒度、按粒度决定能不能比。""" d = station / sub if not d.is_dir(): return None, None, 0, '目录不存在', '年' files = [p for p in d.rglob('*') if p.is_file()] kind = spec.get('span') if kind == 'csv_col': col = spec.get('time_col') starts, ends = [], [] for p in files[:60]: try: with open(p, 'rb') as f: first = f.readline() # 表头 first = f.readline() # 第一条数据行 f.seek(max(0, p.stat().st_size - 65536)) last = f.read().splitlines()[-1] hdr = first.decode('utf-8', 'replace').strip().split(',') i = hdr.index(col) if col in hdr else 0 starts.append(hdr[i]) ends.append(last.decode('utf-8', 'replace').split(',')[0]) except Exception: continue if starts and ends: return min(starts), max(ends), len(files), f'{len(files)} 件 (抽样 {len(starts)} 件首/末数据行)', '日' return None, None, len(files), 'CSV 首/末行解析失败', '年' if kind == 'fname_year': ys = sorted({int(m.group()) for p in files for m in [re.search(r'(20\d\d)', p.name)] if m}) return (f'{ys[0]}-01' if ys else None), (f'{ys[-1]}-12' if ys else None), len(files), \ f'{len(files)} 件 (文件名年份, 年粒度)', '年' if kind == 'fname_date': ds = [] for p in files: m = re.search(r'(\d{2})(\d{2})(20\d\d)', p.name) if m: ds.append(f'{m.group(3)}-{m.group(2)}-{m.group(1)}') return (min(ds) if ds else None), (max(ds) if ds else None), len(files), \ f'{len(files)} 件 (文件名日期)', '日' if kind == 'dir': months = sorted({f'{m.group(1)}-{m.group(2)}' for p in d.rglob('*') for m in [re.search(r'[/\\](20\d\d)[/\\](\d{2})[/\\]', str(p))] if m}) return (months[0] if months else None), (months[-1] if months else None), len(files), \ f'{len(files)} 件 (目录年月 {", ".join(months) or "未识别"})', '月' return None, None, len(files), '按件数清点 (无跨度判据)', '年' def product_stats(farm: str): """→ {仓: dict(n, bytes, files: [(名, 类型, 大小)], kinds)}。""" out = {} for name in STORES: d = P.out_root(farm) / name if not d.is_dir(): out[name] = dict(n=0, bytes=0, kinds={}, newest=None) continue files = [p for p in d.rglob('*') if p.is_file()] kinds = {} for p in files: k = p.suffix.lower() or '(无扩展名)' kinds[k] = kinds.get(k, 0) + 1 newest = max((p.stat().st_mtime for p in files), default=None) out[name] = dict(n=len(files), bytes=sum(p.stat().st_size for p in files), kinds=kinds, newest=(dt.datetime.fromtimestamp(newest).strftime('%Y-%m-%d %H:%M') if newest else None)) return out def provenance(farm: str): root = P.out_root(farm) prov = {} for fn in ('_provenance.json', '_derived_manifest.json'): f = root / fn if not f.exists(): continue try: d = json.loads(f.read_text(encoding='utf-8')) except Exception: continue for rel, v in (d.get('files') or {}).items(): src = v.get('source') if isinstance(v, dict) else None prov[rel] = dict(source=src or 'raw-derived', builder=(v.get('builder') if isinstance(v, dict) else str(v)) or '') return prov def span_check(farm: str, station: pathlib.Path, verbose=True): """逐类输入: 算输入跨度 + 对应产物的跨度/条数, 报联动结论。→ (rows, problems)""" import pandas as pd rows, problems = [], [] ST = P.store(farm) def pspan(p: pathlib.Path, col=None, path_cols=()): if p is None or not p.exists(): return None, None, 0 try: if p.suffix == '.parquet': d = pd.read_parquet(p) n = len(d) c = col if col in d.columns else next((c for c in path_cols if c in d.columns), None) if c is None: return None, None, n s = pd.to_datetime(d[c], errors='coerce') return ((str(s.min())[:10] if s.notna().any() else None), (str(s.max())[:10] if s.notna().any() else None), n) except Exception: return None, None, 0 return None, None, 0 # 输入类 → 它喂出来的产物 (路径, 时间列) PAIRS = { 'scada_10min': [('windscada/temp_monthly.parquet', ST / 'temp_monthly.parquet', 'month'), ('windscada/loss_monthly.parquet', ST / 'loss_monthly.parquet', 'month'), ('windscada/powercurve_bins.parquet', ST / 'powercurve_bins.parquet', None)], '故障报警': [('windscada/alarms.parquet', ST / 'alarms.parquet', 't_on')], '风机故障记录': [('windscada/workorders.parquet', ST / 'workorders.parquet', 't_report')], '油样报告': [('windscada/oil_samples_index.parquet', ST / 'oil_samples_index.parquet', 'date')], } for sub, spec in INPUTS.items(): if sub == 'm5_cms_tcm': continue # 出件不是"测量输入", 无跨度判据 i_s, i_e, i_n, note, gran = input_span(station, sub, spec) comparable = gran in ('日', '月') # 年粒度不参与"落后"判定 (见 input_span 注释) if sub == 'windcms': # ★glob 必须写 `w[0-9][0-9][0-9][0-9]`: `w[????]` 是"字符类里含 ? 之一", 匹配不到任何窗名 # (本脚本第一版就这么写, windcms 那行因此整行不出现 —— 静默漏项, 不是报错) wins = [w for w in sorted((P.m5(farm) / 'windows').glob('w[0-9][0-9][0-9][0-9]')) if (w / 'index.parquet').exists()] if wins: w = wins[-1] p_s, p_e, p_n = pspan(w / 'index.parquet', 'trigger_time') rows.append(dict(input=sub, note=note, in_span=f'{i_s} ~ {i_e}', files=i_n, product=f'm5_cms_tcm/windows/{w.name}/index.parquet', prod_span=f'{p_s} ~ {p_e}', rows_n=p_n)) if p_s and i_s and p_s[:7] < i_s[:7]: problems.append(f'{sub}: 窗 {w.name} 索引起点 {p_s} 早于输入最早月 {i_s} (跨窗混入?)') continue for label, path, col in PAIRS.get(sub, []): p_s, p_e, p_n = pspan(path, col) rows.append(dict(input=sub, note=note, in_span=f'{i_s} ~ {i_e}', files=i_n, gran=gran, product=label, prod_span=f'{p_s} ~ {p_e}', rows_n=p_n)) if comparable and i_e and p_e and _months_between(p_e[:7], i_e[:7]) >= 2: problems.append(f'{sub}: 产物 {label} 只到 {p_e}, 输入到 {i_e} ' f'(落后 {_months_between(p_e[:7], i_e[:7])} 个月) → 未重算 或 该产物无生成端') elif not comparable and p_e: pass # 年粒度: 跨度只作参考, 不判落后 (避免假缺口) if verbose: for r in rows: print(f" {r['input']:12s} 输入 {r['in_span']:26s} {r['files']:6d} 件 | " f"{r['product']:52s} {r['rows_n']:8d} 行 {r['prod_span']}") return rows, problems def _months_between(a: str, b: str) -> int: try: ya, ma = int(a[:4]), int(a[5:7]) yb, mb = int(b[:4]), int(b[5:7]) return (yb - ya) * 12 + (mb - ma) except Exception: return 0 def doc_block(farm: str) -> str: """生成写进 docs/系统设计说明.md 的清单块 (Markdown)。""" ps = product_stats(farm) prov = provenance(farm) from collections import Counter by_store = Counter() for rel, v in prov.items(): by_store[(rel.split('/')[0], v['source'])] += 1 L = [f'', f'*自动生成于 {dt.datetime.now():%Y-%m-%d %H:%M};数据源: `outputs/{farm}/_provenance.json` + `_derived_manifest.json` + 实际文件*', '', '| 产物仓 | 输出路径 (相对安装根) | 功用 | 生成端 | 消费端 | 件数 | 大小 | 来源(raw 重算/随包) |', '|---|---|---|---|---|---|---|---|'] for name, meta in STORES.items(): st = ps[name] rd = by_store.get((name, 'raw-derived'), 0) sh = by_store.get((name, 'shipped'), 0) path = f'outputs/{farm}/{name}/' L.append(f"| `{name}` | `{path}` | {meta['purpose']} | {meta['producer']} | {meta['consumer']} | " f"{st['n']} | {st['bytes'] / 1048576:.1f} MB | {rd} / {sh} |") L += ['', f'合计 {sum(s["n"] for s in ps.values())} 件 / {sum(s["bytes"] for s in ps.values()) / 1048576:.1f} MB;' f'其中 raw 重算 {sum(v for (_, k), v in by_store.items() if k == "raw-derived")} 件、' f'随包补齐 {sum(v for (_, k), v in by_store.items() if k == "shipped")} 件。', MARK_E] return '\n'.join(L) def write_doc(farm: str) -> int: block = doc_block(farm) txt = DOC.read_text(encoding='utf-8') if DOC.exists() else '' if MARK_B in txt and MARK_E in txt: head = txt.split(MARK_B)[0] tail = txt.split(MARK_E)[1] new = head + block + tail else: new = (txt + '\n\n## 附: 产物清单 (自动生成)\n\n' + block + '\n') if txt else block + '\n' DOC.parent.mkdir(parents=True, exist_ok=True) DOC.write_text(new, encoding='utf-8') print(f'已写产物清单 → {P.rel(DOC)}') return 0 def main() -> int: ap = argparse.ArgumentParser() ap.add_argument('--farm', default=os_farm()) ap.add_argument('--check', action='store_true', help='只做呼应校验 (有问题退出码 5)') ap.add_argument('--write-doc', action='store_true', help='把清单块写进 docs/系统设计说明.md') ap.add_argument('--json', action='store_true') a = ap.parse_args() from src.windscada.config import farm as _farm, raw_station_dir station = pathlib.Path(raw_station_dir(a.farm)) ps = product_stats(a.farm) if a.write_doc: return write_doc(a.farm) print(f'=== 产物仓 ({len(STORES)} 个) @ outputs/{a.farm}/ ===') for name, meta in STORES.items(): st = ps[name] kinds = ' '.join(f'{k}×{v}' for k, v in sorted(st['kinds'].items(), key=lambda kv: -kv[1])[:5]) print(f" {name:24s} {st['n']:6d} 件 {st['bytes'] / 1048576:9.1f} MB {meta['helper']:28s} {kinds}") print(f" 功用: {meta['purpose']}") print(f'\n=== 输入 ↔ 产物 呼应 (输入根 {P.rel(station)}) ===') rows, problems = span_check(a.farm, station) if not rows: print(' (无可对拍项)') print(f'\n=== 结论: {"呼应正常" if not problems else str(len(problems)) + " 项要处理"} ===') for p in problems: print(f' [!] {p}') if a.json: print(json.dumps(dict(products={k: {kk: vv for kk, vv in v.items() if kk != "kinds"} for k, v in ps.items()}, spans=rows, problems=problems), ensure_ascii=False, indent=1)) return 5 if (a.check and problems) else 0 def os_farm() -> str: import os return os.environ.get('WINDSCADA_FARM') or 'rudong' if __name__ == '__main__': sys.exit(main())