second_scada_collation.py 9.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247
  1. import logging
  2. import os
  3. import shutil
  4. import time
  5. from collections import defaultdict
  6. from datetime import datetime
  7. from multiprocessing import Pool # 新增
  8. from pathlib import Path
  9. import pandas as pd
  10. FILE_MODIFY_MIN_SECONDS = 60
  11. # ====================== 配置项 ======================
  12. # base_dir = "/data/wind-turbine"
  13. base_dir = "/data/wind_files/headquarter/scada/second"
  14. BASE_TMP_DIR = Path(f"{base_dir}/tmp")
  15. BASE_TMP_BAK_DIR = Path(f"{base_dir}/tmp_bak")
  16. ERROR_DATA_DIR = Path(f"{base_dir}/error_data")
  17. LOG_DIR = Path("/data/logs/second_scada_data_py")
  18. LOG_FILE = LOG_DIR / f"process_{datetime.now().strftime('%Y%m%d')}.log"
  19. LOG_LEVEL = logging.INFO
  20. # 多进程配置(根据服务器CPU核心数调整)
  21. # 服务器:64核/126GB,需给其他任务留资源,取约1/3~3/8
  22. MAX_WORKERS = 8
  23. # ====================== 日志 ======================
  24. def init_logger():
  25. LOG_DIR.mkdir(parents=True, exist_ok=True)
  26. logging.basicConfig(
  27. level=LOG_LEVEL,
  28. format="%(asctime)s - %(levelname)s - %(message)s",
  29. handlers=[
  30. logging.FileHandler(LOG_FILE, encoding="utf-8"),
  31. logging.StreamHandler()
  32. ]
  33. )
  34. return logging.getLogger(__name__)
  35. logger = init_logger()
  36. # ====================== 工具方法 ======================
  37. def is_file_valid(file_path: Path) -> bool:
  38. try:
  39. pd.read_parquet(file_path)
  40. return True
  41. except Exception as e:
  42. logger.error(f"文件损坏,跳过:{file_path},异常:{str(e)}")
  43. return False
  44. def get_target_files(base_dir: Path) -> list:
  45. target_files = []
  46. now = time.time()
  47. for root, _, files in os.walk(base_dir):
  48. for file in files:
  49. if not file.endswith(".parquet"):
  50. continue
  51. file_path = Path(root) / file
  52. if now - file_path.stat().st_mtime <= FILE_MODIFY_MIN_SECONDS:
  53. continue
  54. target_files.append(file_path)
  55. return target_files
  56. def group_files_by_date(tmp_files: list) -> dict:
  57. """按 data_date 分组,返回 {date_str: [file_path, ...]}"""
  58. grouped = defaultdict(list)
  59. for f in tmp_files:
  60. # tmp/<data_date>/<model>/<wind_farm_code>/<wind_turbine_code>.parquet
  61. data_date = f.parts[-4]
  62. grouped[data_date].append(f)
  63. return grouped
  64. from functools import wraps
  65. def retry_on_exception(max_retries=3, delay=1, backoff=2):
  66. """重试装饰器"""
  67. def decorator(func):
  68. @wraps(func)
  69. def wrapper(*args, **kwargs):
  70. retries = 0
  71. current_delay = delay
  72. while retries < max_retries:
  73. try:
  74. return func(*args, **kwargs)
  75. except Exception as e:
  76. retries += 1
  77. if retries == max_retries:
  78. raise
  79. logger.warning(f"操作失败,{retries}/{max_retries} 次重试,错误: {e}")
  80. time.sleep(current_delay)
  81. current_delay *= backoff
  82. return None
  83. return wrapper
  84. return decorator
  85. @retry_on_exception(max_retries=3, delay=0.5)
  86. def safe_read_parquet(file_path):
  87. """带重试机制的 Parquet 读取"""
  88. return pd.read_parquet(file_path)
  89. def merge_rolling_days_data(model: str, wind_farm_code: str, wind_turbine_code: str,
  90. df_oneday: pd.DataFrame, ROLLING_DAYS: int) -> None:
  91. BASE_ROLLING_DAYS_DIR = Path(f"{base_dir}/{ROLLING_DAYS}days")
  92. output_path = BASE_ROLLING_DAYS_DIR / model / wind_farm_code / f"{wind_turbine_code}.parquet"
  93. output_path.parent.mkdir(parents=True, exist_ok=True)
  94. df_list = []
  95. if output_path.exists():
  96. df_list.append(safe_read_parquet(output_path))
  97. df_list.append(df_oneday)
  98. merged_df = pd.concat(df_list, ignore_index=True)
  99. merged_df['localtime'] = pd.to_datetime(merged_df['localtime'], errors='coerce')
  100. merged_df = merged_df.dropna(subset=['localtime'])
  101. merged_df.drop_duplicates(subset=['localtime'], keep='last', inplace=True)
  102. merged_df.sort_values(by=['localtime'], inplace=True)
  103. # 按自然日 00:00 保留最近 ROLLING_DAYS
  104. latest_time = merged_df['localtime'].max()
  105. latest_day = latest_time.floor('D')
  106. thirty_days_ago = latest_day - pd.Timedelta(days=ROLLING_DAYS - 1)
  107. merged_df = merged_df[merged_df['localtime'] > thirty_days_ago]
  108. # 合并后整体按1秒颗粒度补齐,向后填充
  109. # 范围对齐到自然日:从最早日 00:00:00 到最新日 23:59:59
  110. if not merged_df.empty:
  111. day_start = merged_df['localtime'].min().floor('D')
  112. day_end = merged_df['localtime'].max().floor('D') + pd.Timedelta(days=1) - pd.Timedelta(seconds=1)
  113. full_index = pd.date_range(start=day_start, end=day_end, freq='1s')
  114. merged_df = merged_df.set_index('localtime').reindex(full_index).ffill()
  115. merged_df.index.name = 'localtime'
  116. merged_df.reset_index(inplace=True)
  117. merged_df.to_parquet(output_path, engine="pyarrow", index=False)
  118. logger.info(f"{ROLLING_DAYS}天滚动文件生成成功:{output_path},行数:{len(merged_df)}")
  119. # ====================== 单个文件处理逻辑(抽成独立函数)======================
  120. def process_single_file(tmp_file):
  121. # 20251020/CCWE1500-82.DF/EmlMGcty/zKCujSuK.parquet
  122. parts = tmp_file.parts
  123. data_date = parts[-4]
  124. model = parts[-3]
  125. wind_farm_code = parts[-2]
  126. wind_turbine_code = tmp_file.stem
  127. # 处理完成后移动到 tmp_bak,避免下次扫描再次拾取
  128. bak_path = BASE_TMP_BAK_DIR / data_date / model / wind_farm_code / f"{wind_turbine_code}.parquet"
  129. bak_path.parent.mkdir(parents=True, exist_ok=True)
  130. try:
  131. if not is_file_valid(tmp_file):
  132. target_path = ERROR_DATA_DIR / data_date / model / wind_farm_code / f"{wind_turbine_code}.parquet"
  133. target_path.parent.mkdir(parents=True, exist_ok=True)
  134. shutil.copy2(tmp_file, bak_path)
  135. shutil.move(tmp_file, target_path)
  136. return "fail"
  137. # 读取 + 打标签
  138. df = safe_read_parquet(tmp_file)
  139. for col in df.columns:
  140. if col != "localtime":
  141. df[col] = pd.to_numeric(df[col], errors="coerce")
  142. # 按1秒颗粒度补齐,向后填充
  143. # 范围对齐到自然日:从当天 00:00:00 到 23:59:59
  144. df['localtime'] = pd.to_datetime(df['localtime'], errors='coerce')
  145. df = df.dropna(subset=['localtime'])
  146. df.drop_duplicates(subset=['localtime'], keep='last', inplace=True)
  147. df.sort_values(by=['localtime'], inplace=True)
  148. day_start = df['localtime'].min().floor('D')
  149. day_end = day_start + pd.Timedelta(days=1) - pd.Timedelta(seconds=1)
  150. full_index = pd.date_range(start=day_start, end=day_end, freq='1s')
  151. df = df.set_index('localtime').reindex(full_index).ffill()
  152. df.index.name = 'localtime'
  153. df.reset_index(inplace=True)
  154. # 14天数据合并
  155. merge_rolling_days_data(model, wind_farm_code, wind_turbine_code, df, 14)
  156. shutil.move(tmp_file, bak_path)
  157. return "success"
  158. except Exception as e:
  159. logger.error(f"处理失败:{tmp_file},异常:{str(e)}", exc_info=True)
  160. target_path = ERROR_DATA_DIR / data_date / model / wind_farm_code / f"{wind_turbine_code}.parquet"
  161. target_path.parent.mkdir(parents=True, exist_ok=True)
  162. shutil.move(tmp_file, target_path)
  163. return "fail"
  164. # ====================== 主流程(多进程)======================
  165. def process_wind_data():
  166. logger.info(f"临时目录:{BASE_TMP_DIR}")
  167. tmp_files = get_target_files(BASE_TMP_DIR)
  168. if not tmp_files:
  169. logger.info("无待处理文件,任务结束")
  170. return
  171. logger.info(f"待处理文件总数:{len(tmp_files)}")
  172. # 按日期分组,串行处理每一天(避免同风机跨天并发写同一 14days 文件)
  173. grouped = group_files_by_date(tmp_files)
  174. sorted_dates = sorted(grouped.keys())
  175. logger.info(f"待处理日期:{sorted_dates}")
  176. total_success = 0
  177. total_fail = 0
  178. # 整个任务共用一个 Pool;maxtasksperchild 控制工作进程回收,防止内存累积
  179. with Pool(processes=MAX_WORKERS, maxtasksperchild=20) as pool:
  180. for data_date in sorted_dates:
  181. date_files = grouped[data_date]
  182. logger.info(f"开始处理日期 {data_date},文件数:{len(date_files)}")
  183. results = pool.map(process_single_file, date_files)
  184. success_count = sum(1 for res in results if res == "success")
  185. fail_count = sum(1 for res in results if res == "fail")
  186. total_success += success_count
  187. total_fail += fail_count
  188. logger.info(f"日期 {data_date} 完成 | 成功:{success_count} | 失败:{fail_count}")
  189. logger.info(f"全部任务完成 | 成功:{total_success} | 失败:{total_fail}")
  190. logger.info("=" * 50)
  191. if __name__ == "__main__":
  192. process_wind_data()