place_raw_data.py 8.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173
  1. #!/usr/bin/env python3
  2. # -*- coding: utf-8 -*-
  3. """把现场给的数据包按 A2 约定落位到 data/raw/<场站名称>/ 下 (2026-09-11)。
  4. ## 为什么要有这个脚本
  5. A2 定了四项数据源的落位: data/raw/<场站名称>/{scada_10min, 故障报警, 风机故障记录, 油样报告}。
  6. 现场拿到的却是几个 GBK 名的大压缩包 (10分钟数据.zip / 如东风场数据.zip), 内部目录名与
  7. 落位目录名并不一一对应 (例如报警在 `报警数据/`, 工单在 `工作/风机故障记录/`,
  8. 油样埋在 `数据收集/更新/(8)…/2025年油样/`)。手工拖拽容易拖错一层、也说不清依据什么。
  9. 本脚本把映射**写成表**, 于是"哪个包里的哪个目录去了哪"是可复核的, 重复跑也不会走样。
  10. ## 落位映射 (target 相对 data/raw/<场站名称>/)
  11. scada_10min/ ← 10分钟数据.zip 全部 38 个 WTG*.csv (平铺; src/windscada/data.py 按
  12. <src_10min>/<turbine>.csv 读, 所以必须是文件本身不是再套一层目录)
  13. 故障报警/ ← 如东风场数据.zip: 报警数据/** (年度/季度 .xls, XML 报警导出)
  14. 如东风场数据.zip: 数据收集/更新/(2)故障记录(首发故障有标识)(2025.1-至今)/故障记录/**
  15. (2025年全年故障记录.xls / 2026年至今.xls — 正是页面写的 *年至今.xls)
  16. 风机故障记录/ ← 如东风场数据.zip: 工作/风机故障记录/** (2021…2026 年故障记录/ + 业主统计故障/ + 年度 .rar)
  17. 如东风场数据.zip: 数据收集/更新/(3)现场检修记录(2025.1-至今)/大部件维修记录.*.xlsx
  18. 油样报告/ ← 如东风场数据.zip: 数据收集/更新/(8)风机振动数据、油液分析记录(2025.1-至今)/2025年油样/**
  19. (去掉"2025年油样"这一层, 直接落成 <台号>/<部件>/<pdf>,
  20. 与维护页说明"按台号/部件分目录"一致)
  21. ## 故意不做的事
  22. · 不落 `(3)现场检修记录/2025年检修记录` 与 `2026年检修记录`: 与 工作/风机故障记录/2025年故障记录、
  23. 2026年故障记录 **逐件同名同大小**(19/19 件), 是同一批月度汇总表的副本。落两遍会让"按年目录"
  24. 摄入看到两份重复台账 (2026-08-31 长停台账虚高 35 倍那类事故的同款成因: 快照重复必须归并, 不能叠加)。
  25. · 不落 `(8)…/振动分析报告/`: 属振动线 (CMS), 不在 A2 四项数据层内。
  26. · 不落 1分钟数据.zip / scada数据(如东)/ / 西门子4.0技术资料/ 等: A3 本轮范围是 A2 那四类。
  27. 用法:
  28. python scripts/place_raw_data.py --dry-run # 只报要落什么, 不写盘
  29. python scripts/place_raw_data.py # 真落位 (已存在则覆盖)
  30. python scripts/place_raw_data.py --src D:\\别的现场数据目录
  31. """
  32. from __future__ import annotations
  33. import argparse
  34. import pathlib
  35. import shutil
  36. import sys
  37. import zipfile
  38. ROOT = pathlib.Path(__file__).resolve().parents[1]
  39. sys.path.insert(0, str(ROOT))
  40. # 现场数据包所在目录: 用 --src 指定, 或 env GUANLAN_PLACE_SRC; 默认取安装目录下 data/_incoming
  41. DEFAULT_SRC = pathlib.Path(os.environ.get('GUANLAN_PLACE_SRC') or (ROOT / 'data' / '_incoming'))
  42. ZIP_10MIN = '10分钟数据.zip'
  43. ZIP_FARM = '如东风场数据.zip'
  44. # (压缩包, 包内前缀, 目标子目录, 是否去掉前缀这一层)
  45. RULES = [
  46. ('10分钟数据.zip', '', 'scada_10min', True),
  47. ('如东风场数据.zip', '如东风场数据/报警数据/', '故障报警', True),
  48. ('如东风场数据.zip', '如东风场数据/数据收集/更新/(2)故障记录(首发故障有标识)(2025.1-至今)/故障记录/', '故障报警', True),
  49. ('如东风场数据.zip', '如东风场数据/工作/风机故障记录/', '风机故障记录', True),
  50. ('如东风场数据.zip', '如东风场数据/数据收集/更新/(3)现场检修记录(2025.1-至今)/大部件维修记录.20240619143912557.xlsx', '风机故障记录', True),
  51. ('如东风场数据.zip', '如东风场数据/数据收集/更新/(8)风机振动数据、油液分析记录(2025.1-至今)/2025年油样/', '油样报告', True),
  52. ]
  53. # 说明性的"故意不落", 只在报告里列出来
  54. SKIPPED = [
  55. ('如东风场数据.zip', '如东风场数据/数据收集/更新/(3)现场检修记录(2025.1-至今)/2025年检修记录/',
  56. '与 工作/风机故障记录/2025年故障记录 逐件同名校验相同 (副本)'),
  57. ('如东风场数据.zip', '如东风场数据/数据收集/更新/(3)现场检修记录(2025.1-至今)/2026年检修记录/',
  58. '与 工作/风机故障记录/2026年故障记录 逐件同名校验相同 (副本)'),
  59. ('如东风场数据.zip', '如东风场数据/数据收集/更新/(8)风机振动数据、油液分析记录(2025.1-至今)/振动分析报告/',
  60. '属振动线 (CMS), 不在 A2 四项数据层内'),
  61. ]
  62. def gbk_name(info: zipfile.ZipInfo) -> str:
  63. """zip 条目名 → 真名。现场包用 GBK 存中文名, zipfile 按 cp437 解出来是乱码。"""
  64. raw = info.filename
  65. try:
  66. return raw.encode('cp437').decode('gbk')
  67. except Exception:
  68. try:
  69. return raw.encode('cp437').decode('utf-8')
  70. except Exception:
  71. return raw
  72. def human(n: float) -> str:
  73. for unit in ('B', 'KB', 'MB', 'GB'):
  74. if n < 1024 or unit == 'GB':
  75. return f'{n:.1f} {unit}'
  76. n /= 1024
  77. return f'{n:.1f} GB'
  78. def plan(src: pathlib.Path, station: pathlib.Path):
  79. """→ [(zip, 目标文件, 包内条目名, 解压后大小)]; 只读压缩包目录, 不解压。"""
  80. out = []
  81. for zname, prefix, target, strip in RULES:
  82. zp = src / zname
  83. if not zp.exists():
  84. raise SystemExit(f'缺压缩包: {zp}')
  85. with zipfile.ZipFile(zp) as zf:
  86. for info in zf.infolist():
  87. if info.is_dir():
  88. continue
  89. name = gbk_name(info).replace('\\', '/')
  90. if not name.startswith(prefix):
  91. continue
  92. rest = name[len(prefix):] if strip else name
  93. if not rest:
  94. continue
  95. out.append((zname, station / target / rest, name, info.file_size))
  96. return out
  97. def main() -> int:
  98. ap = argparse.ArgumentParser()
  99. ap.add_argument('--src', default=str(DEFAULT_SRC), help='现场数据目录 (默认 %(default)s)')
  100. ap.add_argument('--dry-run', action='store_true', help='只报计划, 不写盘')
  101. a = ap.parse_args()
  102. src = pathlib.Path(a.src)
  103. if not src.is_dir():
  104. raise SystemExit(f'现场数据目录不存在: {src}')
  105. from src.windscada.config import farm, raw_station_dir
  106. station = pathlib.Path(raw_station_dir())
  107. cfg = farm()
  108. print(f'场站目录: {station} (来自 src/windscada/config.py raw_station={cfg.get("raw_station")!r})')
  109. print(f'现场数据: {src}\n')
  110. items = plan(src, station)
  111. by_target = {}
  112. for zname, dst, entry, size in items:
  113. t = dst.relative_to(station).parts[0]
  114. d = by_target.setdefault(t, [0, 0])
  115. d[0] += 1
  116. d[1] += size
  117. print('== 计划落位 ==')
  118. for t in ('scada_10min', '故障报警', '风机故障记录', '油样报告'):
  119. n, sz = by_target.get(t, (0, 0))
  120. print(f' {t:12s} {n:5d} 件 {human(sz):>10s} → {station / t}')
  121. print(f' 合计 {len(items)} 件, {human(sum(i[3] for i in items))}')
  122. print('\n== 故意不落 ==')
  123. for zname, prefix, why in SKIPPED:
  124. print(f' {prefix}\n 理由: {why}')
  125. if a.dry_run:
  126. print('\n(dry-run, 未写盘)')
  127. return 0
  128. print('\n== 落位 ==')
  129. done = 0
  130. for zname, dst, entry, size in items:
  131. dst.parent.mkdir(parents=True, exist_ok=True)
  132. with zipfile.ZipFile(src / zname) as zf:
  133. info = next(i for i in zf.infolist() if gbk_name(i).replace('\\', '/') == entry)
  134. with zf.open(info) as fsrc, open(dst, 'wb') as fdst:
  135. shutil.copyfileobj(fsrc, fdst, 1024 * 1024 * 4)
  136. got = dst.stat().st_size
  137. if got != size:
  138. raise SystemExit(f'写出大小不符: {dst} 期望 {size} 实得 {got}')
  139. done += 1
  140. if done % 10 == 0 or size > 50 * 1024 * 1024:
  141. print(f' [{done}/{len(items)}] {human(size):>10s} {dst.relative_to(station)}', flush=True)
  142. print(f'\n完成: {done} 件写入 {station}')
  143. return 0
  144. if __name__ == '__main__':
  145. sys.exit(main())