#!/usr/bin/env python3 # -*- coding: utf-8 -*- r"""输入数据放置体检 + 增量放置指导 (2026-09-17 用户令 3)。 用户要的是两件事: ① **检查**放进去的输入数据(特别是**增量**放的)是否符合规则; ② 给出**放置指导**(该放哪、命名怎么写、放完跑什么)。 ## 规则 (每条都能指到具体文件, 依据 = 数据目录约定 + 各摄入脚本的真实读法) R1 顶层只许 场站目录 / 共享资料(西门子4.0技术资料) / 说明文件(*.txt|*.md); 常见错误: 把 10 个 csv 直接解到 data/raw 根下, 或把整包 zip 丢在根上。 R2 场站目录下只许约定子目录 (src/windscada/config.py::STATION_SUBDIRS); 常见错误: 解压时多套一层包名目录 (如 `10分钟数据/10分钟数据/WTG01.csv`)。 R3 各源类的结构与命名: scada_10min/ scada_1min/ 平铺 `<机组>.csv`, 机组集合要与场配置一致 故障报警/ .xls/.xlsx/.xml, 不许多套层 风机故障记录/ 年目录 + .xls/.xlsx/.rar/.jpg 油样报告/ 两级 `台号/部件/*.pdf` (维护页按此分目录) windcms/ 含 CMS 原始导出: `*_decode.json` 或 `<包名>/measurement/**` m5_cms_tcm/ 含 handoff 正本 (handoff_vibration_v2.json / component_history.json) + 厂家报告/ R4 可读性抽样: CSV 能按 utf-8/gbk 解出非空表头; TCM `*_decode.json` 能解析出场站/机组字段。 (`--deep` 才做, 默认只查结构与命名, 免得在 GB 级目录上耗时) R5 **增量语义** (给了 --src 时): 逐件对比"源包会落到哪"与"现在有什么", 出三清单: 新增 / 相同(将跳过) / **冲突**(同名不同大小 ⇒ 会覆盖已摄入的数据, 必须人来定夺)。 R6 时间空洞: scada_* 的文件名里若带年月, 检查是否缺月/重月 (数据层页面的时间轴会缺一段)。 R7 体量提示: 单文件 > 2 GB、目录 > 60 GB、或混进 zip/rar (应先解压再放)。 R8 放置指导: 按"新增了哪类源"给出该跑的重算命令 (台账/SCADA/振动/全量)。 ## 退出码 0 合规(可能带提示) · 5 结构或命名违例 · 6 增量冲突(同名不同大小) · 7 缺源类/时间空洞 · 8 仅有提示 ## 用法 python scripts/raw_data_check.py # 体检现有 data/raw/<场站> python scripts/raw_data_check.py --deep # 加可读性抽样 python scripts/raw_data_check.py --src <现场包目录> # 增量体检: 出新增/相同/冲突三清单 python scripts/raw_data_check.py --src <现场包> --json python scripts/raw_data_check.py --write-doc # 把规则表与结论写进放置指导文档 """ from __future__ import annotations import argparse import json import os import pathlib import re import sys import zipfile ROOT = pathlib.Path(__file__).resolve().parents[1] sys.path.insert(0, str(ROOT)) from src import paths as P # noqa: E402 RC_OK, RC_STRUCT, RC_CLASH, RC_GAP, RC_NOTE = 0, 5, 6, 7, 8 TOP_OK_DIRS = ('西门子4.0技术资料',) # data/raw 下与"场站目录"并列的共享资料 TOP_OK_FILES = ('.txt', '.md', '.pdf') # 说明类文件 ZIP_MOD = {'.zip', '.rar', '.7z', '.tar', '.gz', '.zst'} # Access 读写 .mdb 时留下的锁/临时文件 (不是源数据; 2026-09-19 实逮: MDB 重建期间 scada_mdb 下出现 .laccdb) ACCESS_TMP_SUFFIX = {".laccdb", ".ldb", ".tmp"} BIG_FILE = 2 * 1024 ** 3 BIG_DIR = 60 * 1024 ** 3 # 各类源的期望 (后缀集合, 说明) —— 与摄入脚本的真实读法一致 (见 place_raw_data.py 的映射表注释) EXPECT = { 'scada_10min': ({'.csv'}, '平铺 <机组>.csv (每台一个文件, 不许再套层)'), 'scada_1min': ({'.csv'}, '平铺 <机组>.csv'), '故障报警': ({'.xls', '.xlsx', '.xml'}, '年度/季度报警导出 (.xls) + XML 报警导出'), '风机故障记录': ({'.xls', '.xlsx', '.rar', '.jpg', '.png'}, '年目录 + 月度汇总表/现场照片'), '油样报告': ({'.pdf'}, '两级: <台号>/<部件>/*.pdf'), 'windcms': ({'.json', '.pdf', '.docx'}, 'CMS 原始导出 *_decode.json(可套 <包名>/measurement/ 层) + 厂家报告/'), 'm5_cms_tcm': ({'.json', '.docx', '.pdf', '.md'}, 'handoff 正本 + 厂家报告/'), # 2026-09-17 用户令: 现场交来的年度 MDB 归档(25年.zip / 26年.zip)就地放在这里 —— 它就是 # 原始收资(Access 按月分类),取代此前的"已导出 CSV"两个 zip(10分钟数据.zip / 1分钟数据.zip)。 # 允许 .zip(原样归档)+ .mdb(已解出时);解出的目录层次 `<年>/<月>/` 由 R2 的子目录规则放行。 'scada_mdb': ({'.zip', '.mdb'}, '年度归档 <年>年.zip(内层 <年>/<月>/<年-月-类>.zip → .mdb)'), } DATE_RE = re.compile(r'(20\d{2})[-_年]?(0[1-9]|1[0-2])') def station_dir(farm: str | None = None) -> pathlib.Path: """当前场的原始件目录: data/raw/<场站目录名>。""" from src.windscada import config as C try: d = C.raw_station_dir(farm) except BaseException: d = P.RAW_ROOT / (farm or P.farm()) return pathlib.Path(d) def sizes_map(base: pathlib.Path) -> dict: """{相对路径: 字节数} —— 只算文件。""" out = {} if base.is_dir(): for f in base.rglob('*'): if f.is_file(): out[f.relative_to(base).as_posix()] = f.stat().st_size return out def check_tree(st: pathlib.Path, deep=False) -> tuple[list, list, dict]: """规则体检 → (问题列表[(level, 项, 说明, rc)], 提示列表, 统计)。""" res, notes = [], [] stat = dict(exists=st.is_dir(), files=0, bytes=0, subdirs={}, units=[]) # R1 data/raw 顶层 if P.RAW_ROOT.is_dir(): for p in sorted(P.RAW_ROOT.iterdir()): if p.is_dir(): if p.name != st.name and p.name not in TOP_OK_DIRS: res.append(('!', f'data/raw/{p.name}', f'顶层目录未在约定内 (场站目录应为 {st.name}/; ' f'共享资料放 {TOP_OK_DIRS})', RC_STRUCT)) elif p.suffix.lower() not in TOP_OK_FILES: res.append(('!', f'data/raw/{p.name}', 'data/raw 顶层只放场站目录与说明文件; ' '源件要放进 <场站>/<源类>/', RC_STRUCT)) if not st.is_dir(): notes.append(('i', f'{P.rel(st)}', '尚未放置输入数据 —— 页面会显示"无产物/无数据" (属预期); ' '放置见 docs/输入数据放置指导_v0.1.md', RC_OK)) return res, notes, stat # R2 场站内子目录 + R3 各源类 from src.windscada import config as C allowed = set(C.STATION_SUBDIRS) for d in sorted(st.iterdir()): if d.is_file(): if d.suffix.lower() in ZIP_MOD: res.append(('!', P.rel(d), '压缩包直接放在场站目录下 —— 先解压再按源类落位', RC_NOTE)) else: notes.append(('i', P.rel(d), '场站目录下的散装文件 (多半是现场随手拷的; 建议归入对应源类子目录)', RC_OK)) continue if d.name not in allowed: res.append(('!', f'{P.rel(d)}', f'未约定的子目录 (约定: {sorted(allowed)})' ' —— 常见成因: 解压时多套了一层包名', RC_STRUCT)) continue files = [f for f in d.rglob('*') if f.is_file()] nbytes = sum(f.stat().st_size for f in files) stat['subdirs'][d.name] = dict(files=len(files), bytes=nbytes) stat['files'] += len(files) stat['bytes'] += nbytes exts, bad = {}, [] for f in files: e = f.suffix.lower() exts[e] = exts.get(e, 0) + 1 want = EXPECT.get(d.name, (None, ''))[0] # ★2026-09-19: Access 的**锁文件/临时文件**(.laccdb/.ldb/~$*) 不是源数据 —— # 读/写 .mdb 时 Access 会在同目录留下它(实测: MDB 重建期间 scada_mdb 下出现 .laccdb), # 按"后缀不在约定内"报结构问题会把人引到错的地方。忽略并如实注明。 if e in ACCESS_TMP_SUFFIX or f.name.startswith('~$'): notes.append(('i', P.rel(f), 'Access 锁/临时文件 (读写 .mdb 时产生), 不算源件', RC_OK)) continue if want and e not in want and e not in ZIP_MOD: bad.append(f) if bad: sample = ', '.join(sorted({f.suffix.lower() for f in bad})[:4]) res.append(('!', P.rel(d), f'{len(bad)} 个文件后缀不在该源类的约定内 ({sample}); ' f'约定: {sorted(EXPECT[d.name][0])}', RC_STRUCT)) if not files: notes.append(('i', P.rel(d), '目录是空的 (放了数据才有用)', RC_OK)) # 各类专项 if d.name in ('scada_10min', 'scada_1min'): names = {f.name for f in files} nested = [f for f in files if len(f.relative_to(d).parts) > 1] if nested: res.append(('!', P.rel(d), f'{len(nested)} 个文件在子目录里 —— 摄入按 <源类>/<机组>.csv 平铺读, ' f'多一层就读不到', RC_STRUCT)) try: want = set(C.farm(None)['turbines']) n_want = len(want) except BaseException: want, n_want = set(), 0 got = {f.stem for f in files if f.parent == d} # ★ 两类命名都对: scada_10min 用显示机组号 (WTG01.csv); scada_1min 的现场导出用的是 # **内部台号** (01E.csv…37B.csv, 见 place_raw_data.py 的映射说明) —— 后者按"台数一致"判完整, # 不再按名字判 (2026-09-17 实测: 按 WTG 名字判会把 38/38 齐全的 1min 目录误报成"缺 38 台")。 if d.name == 'scada_10min': miss = sorted(want - got) if want else [] extra = sorted(got - want) if want else [] stat['units'].append((d.name, len(got))) if miss: res.append(('!', P.rel(d), f'缺 {len(miss)} 台机组的数据: {miss[:6]}' f'{" …" if len(miss) > 6 else ""} (页面按台取数, 缺台就是空白)', RC_GAP)) if extra: res.append(('!', P.rel(d), f'{len(extra)} 个文件名不是本场机组: {extra[:6]}', RC_STRUCT)) else: stat['units'].append((d.name, len(got))) if n_want and len(got) != n_want: res.append(('!', P.rel(d), f'台数不符: 现有 {len(got)} 个 <台号>.csv, 场配置是 {n_want} 台' f' (1min 导出用内部台号 01E/37B 这类, 按台数判完整)', RC_GAP)) # R6 时间空洞 (文件名带年月时才判) months = sorted({m.group(0) for m in (DATE_RE.search(f.stem) for f in files) if m}) if months: stat['months'] = months if d.name == '油样报告': flat = [f for f in files if len(f.relative_to(d).parts) < 2] if flat: res.append(('!', P.rel(d), f'{len(flat)} 个 PDF 直接躺在源类目录下 —— 约定是 <台号>/<部件>/*.pdf ' f'(维护页按此分组)', RC_STRUCT)) if d.name == 'windcms': js = [f for f in files if f.suffix.lower() == '.json'] dec = [f for f in js if f.name.endswith('_decode.json')] # TCM 导出本体 if not dec: res.append(('!', P.rel(d), '没有 *_decode.json (CMS 原始导出) —— 振动侧就没有可重算的源件', RC_GAP)) if deep and dec: # 抽样核对"摄入脚本真正要的结构" —— 依据 scripts/rudong_tcm_index.py: # doc['body']['body'][] = [ {Record: {...}}, … ], 机组取 Record.Location / Record.LocationName # (顶层还有 method/controller/buildId 等信封字段, 别拿它们当判据 —— 2026-09-17 实测踩过这个坑) import random as _rnd _rnd.seed(0) sample = _rnd.sample(dec, min(12, len(dec))) ok_n, bad_n, shapes = 0, 0, {} for f in sample: try: o = json.loads(f.read_text(encoding='utf-8', errors='replace')) body = ((o or {}).get('body') or {}).get('body') if isinstance(o, dict) else None if not isinstance(body, dict) or not body: shapes[str(sorted(o)[:4]) if isinstance(o, dict) else type(o).__name__] = \ shapes.get(str(sorted(o)[:4]) if isinstance(o, dict) else type(o).__name__, 0) + 1 bad_n += 1 continue recs = next((v for v in body.values() if isinstance(v, list) and v), None) keys = sorted((recs[0] or {}).get('Record', {}) or {}) if recs else [] if recs and (('Location' in keys) or ('LocationName' in keys)): ok_n += 1 else: shapes[f'Record 缺 Location/LocationName: {keys[:6]}'] = \ shapes.get(f'Record 缺 Location/LocationName: {keys[:6]}', 0) + 1 bad_n += 1 except Exception as e: shapes[f'解析失败 {type(e).__name__}'] = shapes.get(f'解析失败 {type(e).__name__}', 0) + 1 bad_n += 1 stat['tcm_sample'] = dict(ok=ok_n, bad=bad_n, shapes=shapes) if bad_n: res.append(('!', P.rel(d), f'TCM 抽样 {len(sample)} 件里 {bad_n} 件结构对不上摄入脚本的取法 ' f'(body.body[时间戳]=[{{Record:…}}], 机组取 Record.Location/LocationName): ' f'{list(shapes.items())[:3]}', RC_STRUCT)) else: stat['tcm_ok'] = True if d.name == 'm5_cms_tcm': hand = [f for f in files if f.name in ('handoff_vibration_v2.json', 'component_history.json')] if not hand: notes.append(('i', P.rel(d), '缺 handoff 正本 (handoff_vibration_v2.json / component_history.json) —— ' '融合面判级与 /cms/ 只能用随包快照 (docs §7 的已知缺口)', RC_OK)) # R4 可读性抽样 (CSV 编码/表头) if deep: for name in ('scada_10min', 'scada_1min'): d = st / name if d.is_dir(): f = next(iter(sorted(x for x in d.glob('*.csv'))), None) if f: head = None for enc in ('utf-8-sig', 'gbk'): try: with f.open(encoding=enc, errors='strict') as fh: head = fh.readline().strip() stat[f'{name}_enc'] = enc break except Exception: continue if not head: res.append(('!', P.rel(f), 'CSV 表头读不出来 (utf-8/gbk 都不行)', RC_STRUCT)) elif len(head.split(',')) < 5: res.append(('!', P.rel(f), f'表头列数异常 ({len(head.split(","))} 列): {head[:60]}', RC_STRUCT)) # R7 体量 if stat['bytes'] > BIG_DIR: notes.append(('?', P.rel(st), f'源件合计 {stat["bytes"]/2**30:.1f} GB —— 摄入与重算会很慢, 建议分批', RC_NOTE)) return res, notes, stat def increment_findings(plan: dict, src_name: str = '') -> list: """把增量三清单里的**冲突**变成一条可拦人的问题 (供本模块与 place_raw_data 共用)。""" out = [] if plan.get('clash'): c = plan['clash'][0] out.append(('!', f'--src {src_name}'.strip(), f'{len(plan["clash"])} 件**同名不同大小**: 直接放会覆盖已摄入的源件 ' f'(示例 {c["path"]}: 现有 {c.get("now", 0):,} B → 源包 ' f'{c["size"]:,} B) —— 先确认哪一份是对的, 再决定 --force 覆盖还是改名并存', RC_CLASH)) return out def increment_src(src: pathlib.Path, st: pathlib.Path) -> tuple[list, dict]: """R5 增量三清单: 遍历源包 zip 成员/散件, 算出它们会落到哪, 与现有树对比。""" res, plan = [], dict(new=[], same=[], clash=[], unmapped=[]) import importlib.util spec = importlib.util.spec_from_file_location('_place', ROOT / 'scripts' / 'place_raw_data.py') place = importlib.util.module_from_spec(spec) try: spec.loader.exec_module(place) except SystemExit: pass rules = list(getattr(place, 'RULES', [])) + list(getattr(place, 'RULES_MECH', [])) + list(getattr(place, 'RULES_VIB', [])) loose = list(getattr(place, 'RULES_VIB_LOOSE', [])) have = sizes_map(st) have_root = sizes_map(P.RAW_ROOT) def classify(rel_path: str, size: int, base: pathlib.Path): target = base / rel_path cur = have.get(rel_path) if base == st else have_root.get(rel_path) item = dict(path=rel_path, size=size) if cur is None: plan['new'].append(item) elif cur == size: plan['same'].append(item) else: item['now'] = cur plan['clash'].append(item) for zname, prefix, target, strip, root_kind in rules: zp = src / zname if src.is_dir() else None if not zp or not zp.is_file(): continue base = st if root_kind == 'station' else P.RAW_ROOT try: with zipfile.ZipFile(zp) as z: for info in z.infolist(): if info.is_dir(): continue nm = place.gbk_name(info) if prefix and not nm.startswith(prefix): continue tail = nm[len(prefix):] if prefix else nm if strip: tail = tail.split('/', 1)[1] if '/' in tail else tail rel = f'{target}/{tail}' if prefix and nm == prefix.rstrip('/'): continue classify(rel, info.file_size, base) except zipfile.BadZipFile: res.append(('!', zp.name, '不是合法 zip (现场包常被截断; 重新拷一份)', RC_STRUCT)) for fn, target in loose: fp = src / fn if fp.is_file(): classify(f'{target}/{fn}', fp.stat().st_size, st) return res, plan def guidance(plan: dict, stat: dict) -> list[str]: """R8 放置指导: 按新增源类给出该跑的重算命令与预期变化。""" out = [] new = [x['path'] for x in plan.get('new', [])] kinds = {p.split('/', 1)[0] for p in new} if not new: out.append('没有新增源件: 不必重算 (若刚换过数据, 用 scripts/rebuild_all.py --dry-run 看计划)') return out out.append(f'新增 {len(new)} 件, 涉及源类: {sorted(kinds)}') if 'scada_10min' in kinds or 'scada_1min' in kinds: out.append('· SCADA 侧变了 → `python scripts/rebuild_all.py` (含 10 个构建器, 约 15 分钟); ' '台账/报警/油样一起换时同一条命令即可') if kinds & {'故障报警', '风机故障记录', '油样报告'}: out.append('· 台账类变了 → `python scripts/rebuild_all.py --skip-scada` (约 2 分钟)') if kinds & {'windcms', 'm5_cms_tcm'}: out.append('· 振动侧变了 → 重算链第 ④b 步会自动摄入 (`scripts/vib_raw_build.py`); ' '注意它只做索引/谱, 不会重生成报告 (本包缺六层链, 见 docs §7)') out.append('· 放完先自检: `python scripts/raw_data_check.py --deep` → 再看 `/detail/` 左栏「系统维护」两屏') out.append('· 重算后核对锚点: 报警 39211 · 工单 5876 · 油样 404 · temp_monthly 19494 · 本体 9702 对象') return out def main() -> int: ap = argparse.ArgumentParser() ap.add_argument('--farm', default=None) ap.add_argument('--src', default=None, help='现场数据包目录 (给了就做增量体检: 新增/相同/冲突三清单)') ap.add_argument('--deep', action='store_true', help='加可读性抽样 (CSV 编码/表头, TCM json 字段)') ap.add_argument('--json', action='store_true') ap.add_argument('--write-doc', action='store_true', help='把结论写进 docs/输入数据放置指导_v0.1.md') a = ap.parse_args() st = station_dir(a.farm) res, notes, stat = check_tree(st, deep=a.deep) plan = {} if a.src: src = pathlib.Path(a.src) if not src.is_dir(): print(f'[X] --src 不是目录: {src}') return RC_STRUCT r2, plan = increment_src(src, st) res += r2 + increment_findings(plan, src.name) guide = guidance(plan, stat) if a.json: print(json.dumps(dict(station=P.rel(st), results=[dict(level=x[0], item=x[1], note=x[2], rc=x[3]) for x in res], notes=[dict(level=x[0], item=x[1], note=x[2]) for x in notes], stat=stat, plan=plan, guidance=guide), ensure_ascii=False, indent=1, default=str)) else: print(f'== 输入数据放置体检 · {P.rel(st)} ==') if stat.get('exists'): print(f' {stat["files"]:,} 件 / {stat["bytes"]/2**30:.2f} GB') for k, v in sorted(stat['subdirs'].items()): print(f' {k:14s} {v["files"]:7,d} 件 {v["bytes"]/2**20:9.1f} MB') for level, item, note, rc in res: print(f' [{level}] {item}: {note}') seen = set() for level, item, note, rc in notes: if note in seen: continue seen.add(note) print(f' [{level}] {item}: {note}') if plan: print('\n== 增量放置三清单 (源包 → data/raw) ==') print(f' 新增 {len(plan["new"]):,} 件 (将拷入) · 相同 {len(plan["same"]):,} 件 (同尺寸, 跳过) · ' f'冲突 {len(plan["clash"]):,} 件 (同名不同大小, **需人确认**)') for k, label in (('new', '新增'), ('clash', '冲突')): for x in plan[k][:5]: extra = f' (现有 {x.get("now"):,} B)' if k == 'clash' else '' print(f' [{label}] {x["path"]} {x["size"]:,} B{extra}') if len(plan[k]) > 5: print(f' … 另有 {len(plan[k]) - 5} 件{label}') print('\n== 放置指导 ==') for g in guide: print(f' {g}') rc_map = {x[3] for x in res if x[0] in ('X', '!') and x[3]} rc = max(rc_map) if rc_map else (RC_NOTE if notes else RC_OK) if not a.json: print(f'\n结论: {"合规" if rc == 0 else "见上"}; 退出码 {rc}') if a.write_doc: write_doc(res, notes, stat, plan, guide) return rc DOC = P.ROOT / 'docs' / '输入数据放置指导_v0.1.md' BEGIN, END = '', '' def write_doc(res, notes, stat, plan, guide) -> None: body = [BEGIN, '### 体检结论(自动生成,勿手改)', '', f'· 源件目录 `{P.rel(station_dir())}`:{stat.get("files", 0):,} 件 / {stat.get("bytes", 0)/2**30:.2f} GB'] for k, v in sorted((stat.get('subdirs') or {}).items()): body.append(f' - `{k}/`:{v["files"]:,} 件 / {v["bytes"]/2**20:.1f} MB') if res: body += ['', '| 级别 | 项 | 说明 |', '|---|---|---|'] body += [f'| `{x[0]}` | `{x[1]}` | {x[2]} |' for x in res[:40]] else: body += ['', '· 结构与命名检查:**无不一致**(提示见下)'] for x in notes[:12]: body.append(f'· `{x[0]}` {x[1]}:{x[2]}') if plan: body += ['', f'· 增量(源包 → data/raw):新增 {len(plan["new"]):,} · 相同 {len(plan["same"]):,} · ' f'冲突 {len(plan["clash"]):,}'] body += ['', '**放置后该做什么**:'] + [f'· {g}' for g in guide] + ['', END] block = '\n'.join(body) if DOC.is_file(): t = DOC.read_text(encoding='utf-8') if BEGIN in t and END in t: t = re.sub(re.escape(BEGIN) + r'.*?' + re.escape(END), lambda m: block, t, flags=re.S) else: t = t.rstrip() + '\n\n' + block + '\n' else: t = '# 输入数据放置指导\n\n' + block + '\n' DOC.write_text(t, encoding='utf-8') print(f' 已写入 {P.rel(DOC)}') if __name__ == '__main__': from src import console console.soft() sys.exit(main())