import logging import os import shutil import time from collections import defaultdict from datetime import datetime from multiprocessing import Pool # 新增 from pathlib import Path import pandas as pd FILE_MODIFY_MIN_SECONDS = 60 # ====================== 配置项 ====================== # base_dir = "/data/wind-turbine" base_dir = "/data/wind_files/headquarter/scada/second" BASE_TMP_DIR = Path(f"{base_dir}/tmp") BASE_TMP_BAK_DIR = Path(f"{base_dir}/tmp_bak") ERROR_DATA_DIR = Path(f"{base_dir}/error_data") LOG_DIR = Path("/data/logs/second_scada_data_py") LOG_FILE = LOG_DIR / f"process_{datetime.now().strftime('%Y%m%d')}.log" LOG_LEVEL = logging.INFO # 多进程配置(根据服务器CPU核心数调整) # 服务器:64核/126GB,需给其他任务留资源,取约1/3~3/8 MAX_WORKERS = 8 # ====================== 日志 ====================== def init_logger(): LOG_DIR.mkdir(parents=True, exist_ok=True) logging.basicConfig( level=LOG_LEVEL, format="%(asctime)s - %(levelname)s - %(message)s", handlers=[ logging.FileHandler(LOG_FILE, encoding="utf-8"), logging.StreamHandler() ] ) return logging.getLogger(__name__) logger = init_logger() # ====================== 工具方法 ====================== def is_file_valid(file_path: Path) -> bool: try: pd.read_parquet(file_path) return True except Exception as e: logger.error(f"文件损坏,跳过:{file_path},异常:{str(e)}") return False def get_target_files(base_dir: Path) -> list: target_files = [] now = time.time() for root, _, files in os.walk(base_dir): for file in files: if not file.endswith(".parquet"): continue file_path = Path(root) / file if now - file_path.stat().st_mtime <= FILE_MODIFY_MIN_SECONDS: continue target_files.append(file_path) return target_files def group_files_by_date(tmp_files: list) -> dict: """按 data_date 分组,返回 {date_str: [file_path, ...]}""" grouped = defaultdict(list) for f in tmp_files: # tmp////.parquet data_date = f.parts[-4] grouped[data_date].append(f) return grouped from functools import wraps def retry_on_exception(max_retries=3, delay=1, backoff=2): """重试装饰器""" def decorator(func): @wraps(func) def wrapper(*args, **kwargs): retries = 0 current_delay = delay while retries < max_retries: try: return func(*args, **kwargs) except Exception as e: retries += 1 if retries == max_retries: raise logger.warning(f"操作失败,{retries}/{max_retries} 次重试,错误: {e}") time.sleep(current_delay) current_delay *= backoff return None return wrapper return decorator @retry_on_exception(max_retries=3, delay=0.5) def safe_read_parquet(file_path): """带重试机制的 Parquet 读取""" return pd.read_parquet(file_path) def merge_rolling_days_data(model: str, wind_farm_code: str, wind_turbine_code: str, df_oneday: pd.DataFrame, ROLLING_DAYS: int) -> None: BASE_ROLLING_DAYS_DIR = Path(f"{base_dir}/{ROLLING_DAYS}days") output_path = BASE_ROLLING_DAYS_DIR / model / wind_farm_code / f"{wind_turbine_code}.parquet" output_path.parent.mkdir(parents=True, exist_ok=True) df_list = [] if output_path.exists(): df_list.append(safe_read_parquet(output_path)) df_list.append(df_oneday) merged_df = pd.concat(df_list, ignore_index=True) merged_df['localtime'] = pd.to_datetime(merged_df['localtime'], errors='coerce') merged_df = merged_df.dropna(subset=['localtime']) merged_df.drop_duplicates(subset=['localtime'], keep='last', inplace=True) merged_df.sort_values(by=['localtime'], inplace=True) # 按自然日 00:00 保留最近 ROLLING_DAYS latest_time = merged_df['localtime'].max() latest_day = latest_time.floor('D') thirty_days_ago = latest_day - pd.Timedelta(days=ROLLING_DAYS - 1) merged_df = merged_df[merged_df['localtime'] > thirty_days_ago] # 合并后整体按1秒颗粒度补齐,向后填充 # 范围对齐到自然日:从最早日 00:00:00 到最新日 23:59:59 if not merged_df.empty: day_start = merged_df['localtime'].min().floor('D') day_end = merged_df['localtime'].max().floor('D') + pd.Timedelta(days=1) - pd.Timedelta(seconds=1) full_index = pd.date_range(start=day_start, end=day_end, freq='1s') merged_df = merged_df.set_index('localtime').reindex(full_index).ffill() merged_df.index.name = 'localtime' merged_df.reset_index(inplace=True) merged_df.to_parquet(output_path, engine="pyarrow", index=False) logger.info(f"{ROLLING_DAYS}天滚动文件生成成功:{output_path},行数:{len(merged_df)}") # ====================== 单个文件处理逻辑(抽成独立函数)====================== def process_single_file(tmp_file): # 20251020/CCWE1500-82.DF/EmlMGcty/zKCujSuK.parquet parts = tmp_file.parts data_date = parts[-4] model = parts[-3] wind_farm_code = parts[-2] wind_turbine_code = tmp_file.stem # 处理完成后移动到 tmp_bak,避免下次扫描再次拾取 bak_path = BASE_TMP_BAK_DIR / data_date / model / wind_farm_code / f"{wind_turbine_code}.parquet" bak_path.parent.mkdir(parents=True, exist_ok=True) try: if not is_file_valid(tmp_file): target_path = ERROR_DATA_DIR / data_date / model / wind_farm_code / f"{wind_turbine_code}.parquet" target_path.parent.mkdir(parents=True, exist_ok=True) shutil.copy2(tmp_file, bak_path) shutil.move(tmp_file, target_path) return "fail" # 读取 + 打标签 df = safe_read_parquet(tmp_file) for col in df.columns: if col != "localtime": df[col] = pd.to_numeric(df[col], errors="coerce") # 按1秒颗粒度补齐,向后填充 # 范围对齐到自然日:从当天 00:00:00 到 23:59:59 df['localtime'] = pd.to_datetime(df['localtime'], errors='coerce') df = df.dropna(subset=['localtime']) df.drop_duplicates(subset=['localtime'], keep='last', inplace=True) df.sort_values(by=['localtime'], inplace=True) day_start = df['localtime'].min().floor('D') day_end = day_start + pd.Timedelta(days=1) - pd.Timedelta(seconds=1) full_index = pd.date_range(start=day_start, end=day_end, freq='1s') df = df.set_index('localtime').reindex(full_index).ffill() df.index.name = 'localtime' df.reset_index(inplace=True) # 14天数据合并 merge_rolling_days_data(model, wind_farm_code, wind_turbine_code, df, 14) shutil.move(tmp_file, bak_path) return "success" except Exception as e: logger.error(f"处理失败:{tmp_file},异常:{str(e)}", exc_info=True) target_path = ERROR_DATA_DIR / data_date / model / wind_farm_code / f"{wind_turbine_code}.parquet" target_path.parent.mkdir(parents=True, exist_ok=True) shutil.move(tmp_file, target_path) return "fail" # ====================== 主流程(多进程)====================== def process_wind_data(): logger.info(f"临时目录:{BASE_TMP_DIR}") tmp_files = get_target_files(BASE_TMP_DIR) if not tmp_files: logger.info("无待处理文件,任务结束") return logger.info(f"待处理文件总数:{len(tmp_files)}") # 按日期分组,串行处理每一天(避免同风机跨天并发写同一 14days 文件) grouped = group_files_by_date(tmp_files) sorted_dates = sorted(grouped.keys()) logger.info(f"待处理日期:{sorted_dates}") total_success = 0 total_fail = 0 # 整个任务共用一个 Pool;maxtasksperchild 控制工作进程回收,防止内存累积 with Pool(processes=MAX_WORKERS, maxtasksperchild=20) as pool: for data_date in sorted_dates: date_files = grouped[data_date] logger.info(f"开始处理日期 {data_date},文件数:{len(date_files)}") results = pool.map(process_single_file, date_files) success_count = sum(1 for res in results if res == "success") fail_count = sum(1 for res in results if res == "fail") total_success += success_count total_fail += fail_count logger.info(f"日期 {data_date} 完成 | 成功:{success_count} | 失败:{fail_count}") logger.info(f"全部任务完成 | 成功:{total_success} | 失败:{total_fail}") logger.info("=" * 50) if __name__ == "__main__": process_wind_data()