#!/usr/bin/env python3 # -*- coding: utf-8 -*- r"""输入数据自动扫描识别 —— 发现 `<安装目录>/data/raw` 下的**新增/变化**,并给出"该重跑哪几步"。 ## 为什么要有它(用户令 2026-09-19) 用户令:「对 `<安装目录>/data/raw` 目录下的接入数据处理,增加自动扫描识别机制,以能发现新增数据,并纳入重算。」 在此之前,"放了新数据要不要重算"全靠人记:`place_raw_data.py` 只管把现场包落位,`raw_data_check.py` 只管结构与命名对不对,`rebuild_all.py` 无脑全跑(或按人给的 `--skip-*` 跳过)—— **没有任何一处回答"这回新来了什么、因此必须重跑哪几步"**。本脚本补这一环: 逐族指纹 (件数/体积/最新时间/清单摘要) → 与上次快照比 → 报【新增/变化】 → 映射到链上步骤 ## 三档用法 python scripts/raw_scan.py # 只报当前状态(不写快照) python scripts/raw_scan.py --write # 写/更新快照 outputs/<场>/_raw_scan.json python scripts/raw_scan.py --check # 与快照比: rc=0 无变化 · rc=4 有新增/变化 · rc=5 没有快照 python scripts/raw_scan.py --check --write # 比完就把新快照记下(重算链第 ①b 步就是这条) python scripts/raw_scan.py --plan # 只打印"有变化时该跑哪些步" python scripts/raw_scan.py --json # 机器可读(rebuild_all 用它做"跳过了该跑的步"提醒) python scripts/raw_scan.py --deep # 变化判定加上**内容**哈希(默认只比 名字/大小/时间) python scripts/raw_scan.py --root <别的 raw> # 换个根扫(自测/多场) python scripts/raw_scan.py --selftest # 族表覆盖自检(与族表/链上步骤对账,防两处漂移) ## 口径与纪律 · **指纹**默认 = 每族逐件 `相对路径|字节数|mtime(ns)` 排序后串起来取 sha1 —— 快、且"加了一件/改了一件/ 删了一件"都会变;`--deep` 再叠一层内容 sha1(小文件才做,避免拿 153 GB 的 CMS 导出拼哈希)。 · **未归类**:`data/raw/<场>/` 下出现族表里没有的子目录(新来的数据总会有第一次)→ 如实列进 `unclassified` 并给出"没有消费者/需要认领"的话,**不猜**它是哪一族。 · 这里只**发现与建议**,不改数据、不搬文件、不自动重算(要不要跑由 `rebuild_all.py` 或运维决定; 自动重算的开关见 `configs/serve.json` 的 `raw_watch`)。rc=4 是"有变化",不是失败。 """ from __future__ import annotations import argparse import hashlib import json import os import pathlib import sys import time ROOT = pathlib.Path(__file__).resolve().parents[1] sys.path.insert(0, str(ROOT)) try: # GBK 控制台/日志重定向不再因 '²' 这类字符抛 UnicodeEncodeError sys.stdout.reconfigure(errors='replace') except Exception: pass SNAP_NAME = '_raw_scan.json' # ── 族表:raw 子目录 → (计件 glob, 链上步骤, 消费者说明) ────────────────────────────── # ★步骤标签与 `scripts/rebuild_all.py` 的计划名一一对应(人看到"③ SCADA 侧"就知道去跑哪一步)。 FAMILIES: tuple[tuple[str, tuple[str, ...], tuple[str, ...], str], ...] = ( ('故障报警', ('*.xls', '*.xlsx', '*.csv'), ('② 三门台账',), '报警事件导出 → windscada/alarms.parquet(报警集中/停机账)。缺 TimeOn/Alarmcode 的统计报表件会被摄入器跳过'), ('风机故障记录', ('*.xls', '*.xlsx', '*.csv'), ('② 三门台账',), '检修工单台账 → windscada/workorders.parquet(重复检修/停机损失)'), ('油样报告', ('*.pdf',), ('② 三门台账',), '油液化验报告 → windscada/oil_samples_index.parquet(油液面)'), ('scada_10min', ('*.csv',), ('③ SCADA 侧 10 个构建器', '④ 月度派生件', '④c 变桨面', '⑤b 事实契约 claim'), '10min 导出 → 功率曲线/损失/温度/停机/偏航/液压等 + 变桨面日粒度(温度/性能/变桨面)'), ('scada_1min', ('*.csv',), ('④c 变桨面',), '1min 导出(内部台号 01E…38B)→ 变桨面零位三口径与越线计数'), ('scada_mdb', ('*.mdb', '*.zip'), ('③ SCADA 侧 10 个构建器', '④ 月度派生件', '④c 变桨面'), '现场年度归档 Access 库(按月的 10min/1min)—— CSV 缺失时由 scada_source 回落到它'), ('windcms', ('*_decode.json',), ('④b 振动侧摄入 + CMS 报告 + 标量 z', '④d 振动在升/换件闭环', '④e 三层基线'), 'Brande TCM 导出 → 窗索引/谱库/CMS 报告/逐台页 + 在升闭环 + 三层基线(振动面)'), ('m5_cms_tcm', ('*.json', '*.docx', '*.pdf', '*.md'), ('④b 振动侧摄入 + CMS 报告 + 标量 z',), '现场给的正本(handoff_vibration_v2.json / component_history.json)与厂家报告——正本在位时优先于观澜自算'), ) # 机理层资料不在场站目录下(A2 约定): data/raw/西门子4.0技术资料 TECH_DIR, TECH_PATTERNS = '西门子4.0技术资料', ('*.xlsx', '*.xls', '*.pdf', '*.docx', '*.doc', '*.txt', '*.csv') TECH_STEPS = ('⑦ 本体: 码表/手册/文档',) TECH_NOTE = '厂商技术资料(故障处理手册/维护 WI/图纸/对译表)→ 本体对象库/检索索引/实机参数表' def _digest(entries: list[tuple[str, int, int]], deep: bool, base: pathlib.Path, files: list[pathlib.Path], deep_max: int) -> str: h = hashlib.sha1() for rel, size, mt in entries: h.update(f'{rel}|{size}|{mt}\n'.encode('utf-8', 'replace')) if deep: for p, (rel, size, _mt) in zip(files, entries): if size > deep_max: continue try: fh = hashlib.sha1() with open(p, 'rb') as f: for blk in iter(lambda: f.read(1 << 20), b''): fh.update(blk) h.update(rel.encode('utf-8', 'replace') + b'=' + fh.digest()) except OSError: pass return h.hexdigest() def fingerprint(d: pathlib.Path, patterns: tuple[str, ...], deep: bool = False, deep_max: int = 8 << 20) -> dict: """一族目录的指纹:件数/体积/最新时间/清单摘要(--deep 再叠内容哈希)。目录不在 ⇒ files=0。""" files: list[pathlib.Path] = [] if d.is_dir(): for pat in patterns: files += [p for p in d.rglob(pat) if p.is_file()] files = sorted(set(files)) entries: list[tuple[str, int, int]] = [] for p in files: try: st = p.stat() except OSError: continue entries.append((p.relative_to(d).as_posix(), st.st_size, st.st_mtime_ns)) entries.sort() newest = '' if entries: newest = time.strftime('%Y-%m-%d %H:%M:%S', time.localtime(max(e[2] for e in entries) / 1e9)) return dict(files=len(entries), bytes=sum(e[1] for e in entries), newest=newest, digest=_digest(entries, deep, d, files, deep_max), top=[e[0] for e in sorted(entries, key=lambda x: -x[2])[:3]]) def scan(raw_root: pathlib.Path, deep: bool = False, station_dir: pathlib.Path | None = None) -> dict: """扫一个 raw 根:逐族指纹 + 别处不认识的子目录(未归类)。 `station_dir` 不给时按**场配置的辨识结果**定场站目录(`raw_station_dir()`,与摄入侧同一口径); 显式 `--root` 自测时可传入 `/<场站名>` 之外的位置。 """ from src.windscada.config import raw_station_dir station = pathlib.Path(station_dir) if station_dir else raw_station_dir() fam: dict[str, dict] = {} for name, pats, steps, note in FAMILIES: fp = fingerprint(station / name, pats, deep) fp.update(steps=list(steps), note=note, where=str((station / name))) fam[name] = fp fp = fingerprint(raw_root / TECH_DIR, TECH_PATTERNS, deep) fp.update(steps=list(TECH_STEPS), note=TECH_NOTE, where=str(raw_root / TECH_DIR)) fam[TECH_DIR] = fp known = {n for n, *_ in FAMILIES} | {TECH_DIR} unclassified = [] if station.is_dir(): for sub in sorted(station.iterdir()): if sub.is_dir() and sub.name not in known: unclassified.append(dict(dir=sub.name, files=sum(1 for p in sub.rglob('*') if p.is_file()))) for sub in sorted(raw_root.iterdir()) if raw_root.is_dir() else []: if sub.is_dir() and sub.name not in known and sub != station: unclassified.append(dict(dir=sub.name, files=sum(1 for p in sub.rglob('*') if p.is_file()), note='raw 根下(非场站目录)')) return dict(at=time.strftime('%Y-%m-%d %H:%M:%S'), root=str(raw_root), station=str(station), deep=deep, families=fam, unclassified=unclassified) def diff(old: dict, new: dict) -> list[dict]: """逐族比:新增/减少/内容变化(`files` 与 `digest` 分开说,便于人看懂)。""" out = [] of = (old or {}).get('families') or {} for name, cur in (new.get('families') or {}).items(): prev = of.get(name) if not prev: if cur['files']: out.append(dict(family=name, kind='首次记录', files=cur['files'], steps=cur['steps'], note=cur['note'], newest=cur['newest'])) continue d_files = cur['files'] - prev['files'] same_list = cur['digest'] == prev['digest'] if d_files == 0 and same_list: continue kind = ('新增 %+d 件' % d_files) if d_files else '内容变化(件数不变)' out.append(dict(family=name, kind=kind, files=cur['files'], delta=d_files, bytes_delta=cur['bytes'] - prev['bytes'], newest=cur['newest'], top=cur.get('top'), steps=cur['steps'], note=cur['note'])) for u in new.get('unclassified') or []: out.append(dict(family=u['dir'], kind='未归类(族表里没有这个目录)', files=u['files'], steps=[], note='没有已知消费者 ⇒ 先认领:是新的输入族就得给它配摄入器;否则清出 data/raw')) return out def _steps_of(changes: list[dict]) -> list[str]: seen, out = set(), [] for c in changes: for s in c.get('steps') or []: if s not in seen: seen.add(s) out.append(s) return out def snap_path(raw_root: pathlib.Path) -> pathlib.Path: from src import paths as P return P.out_root() / SNAP_NAME def main() -> int: ap = argparse.ArgumentParser(description='输入数据自动扫描识别(发现新增 → 指明该重跑哪几步)') ap.add_argument('--farm', default=None) ap.add_argument('--root', default=None, help='覆盖 data/raw 根(自测/多场)') ap.add_argument('--write', action='store_true', help='把本次指纹写进快照') ap.add_argument('--check', action='store_true', help='与快照比:rc=0 无变化 · 4 有新增/变化 · 5 无快照') ap.add_argument('--plan', action='store_true', help='只打印"该跑哪些步"') ap.add_argument('--json', action='store_true', help='输出机器可读 JSON') ap.add_argument('--deep', action='store_true', help='变化判定叠内容哈希(小文件)') ap.add_argument('--selftest', action='store_true', help='族表/步骤覆盖自检(防与族表、链上步骤漂移)') a = ap.parse_args() from src import paths as P from src.windscada import config as C if a.farm: os.environ['WINDSCADA_FARM'] = a.farm # 与其它脚本同一口径(场由环境变量选) raw = pathlib.Path(a.root) if a.root else pathlib.Path(C.RAW_ROOT) # --root 换成别的 raw 根时,场站目录按"根下唯一/同名"识别,避免拿本场辨识结果去别处找 st_dir = None if a.root: from src.windscada.config import station_scan sc = station_scan(root=raw) st_dir = (raw / sc['matched']) if sc.get('matched') else None now = scan(raw, deep=a.deep, station_dir=st_dir) sp = snap_path(raw) old = None if sp.is_file(): try: old = json.loads(sp.read_text(encoding='utf-8')) except Exception: old = None changes = diff(old, now) if a.selftest: return selftest() if a.json: print(json.dumps(dict(scan=now, changes=changes, snapshot=str(sp), steps=_steps_of(changes), has_snapshot=bool(old)), ensure_ascii=False, indent=1)) if a.write: sp.parent.mkdir(parents=True, exist_ok=True) sp.write_text(json.dumps(now, ensure_ascii=False, indent=1), encoding='utf-8') if a.check: return 0 if (old and not changes) else (4 if old else 5) return 0 print(f'输入数据扫描 · {now["root"]} ({now["at"]}{" · deep" if a.deep else ""})') tot_f = sum(f['files'] for f in now['families'].values()) tot_b = sum(f['bytes'] for f in now['families'].values()) print(f' 共 {len(now["families"])} 族 · {tot_f} 件 · {tot_b / 1024 / 1024 / 1024:.1f} GB') for name, f in now['families'].items(): flag = '' for c in changes: if c['family'] == name: flag = f" ← {c['kind']}" print(f' {name:16s} {f["files"]:6d} 件 {f["bytes"] / 1024 / 1024:9.1f} MB 最新 {f["newest"] or "—"}{flag}') if now['unclassified']: print(' 未归类(族表里没有,先认领):') for u in now['unclassified']: print(f' {u["dir"]:16s} {u["files"]:6d} 件 {u.get("note", "没有已知消费者")}') if changes: print(f'\n== 本次发现 {len(changes)} 处新增/变化 ==') for c in changes: print(f' · {c["family"]}: {c["kind"]}' + (f'(体积 {c["bytes_delta"] / 1024 / 1024:+.1f} MB)' if c.get('bytes_delta') else '') + f' 最新 {c.get("newest") or "—"}') if c.get('top'): print(f' 最近落盘: ' + '、'.join(c['top'])) if c.get('steps'): print(f' 该跑: ' + ' / '.join(c['steps'])) print('\n 建议: 跑一次重算把新数据纳入产物 —— python scripts/rebuild_all.py' '(默认全跑;若你用了 --skip-*,请确认没跳过上面列出的步)') else: print('\n== 与上次快照一致,没有新数据 ==' if old else '\n== 还没有快照(跑 --write 记一次基线)==') if a.write: sp.parent.mkdir(parents=True, exist_ok=True) sp.write_text(json.dumps(now, ensure_ascii=False, indent=1), encoding='utf-8') print(f'\n快照已更新 → {P.rel(sp)}') if a.check: if not old: print(f'[?] 没有快照可比({P.rel(sp)})—— 先跑 --write 记基线; rc=5') return 5 return 4 if changes else 0 return 0 def selftest() -> int: """自检:① 族表覆盖 `products_reverse_audit.FAMILIES` 里所有 raw 输入目录 ② 步骤名都能在链上找到。""" import importlib.util import re spec = importlib.util.spec_from_file_location('_pra', ROOT / 'scripts' / 'products_reverse_audit.py') mod = importlib.util.module_from_spec(spec) spec.loader.exec_module(mod) inputs = {f['input'] for f in mod.FAMILIES if f.get('input')} mine = {n for n, *_ in FAMILIES} | {TECH_DIR} missing = sorted(inputs - mine) chain = (ROOT / 'scripts' / 'rebuild_all.py').read_text(encoding='utf-8') steps = {s for _n, _p, ss, _note in FAMILIES for s in ss} | set(TECH_STEPS) bad_steps = sorted(s for s in steps if s.split()[0] not in chain) print(f'族表覆盖: 反向呼应族表的 raw 输入 {len(inputs)} 个 · 本器覆盖 {len(mine)} 个 · 缺 {missing or "无"}') print(f'步骤标签: {len(steps)} 个 · 在 rebuild_all 里找不到的 {bad_steps or "无"}') ok = not missing and not bad_steps print('[OK] 自检通过' if ok else '[X] 自检不过 —— 族表或步骤标签与链上漂移了') return 0 if ok else 5 if __name__ == '__main__': sys.exit(main())