| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247 |
- 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/<data_date>/<model>/<wind_farm_code>/<wind_turbine_code>.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()
|