Python多进程处理多年分时DataFrame时程序挂起问题求助
多进程处理海量时间序列数据挂起问题排查与解决
问题背景
我有一个包含多年分钟级时间序列记录的大型DataFrame,需要对每分钟执行复杂计算,因此采用多进程按日历日拆分任务来加速处理。
现有代码
主进程代码
from multiprocessing import set_start_method set_start_method('spawn') import multiprocessing import pandas as pd from lib.targetcalculations import calculate_target # 获取df中的唯一日期列表 date_range = df.reset_index().datetime.map(lambda t: t.date()).unique() # 执行进程 - 每天一个任务 if __name__ == '__main__': n_processes = 8 p = multiprocessing.Pool(n_processes) results = p.map(calculate_target, [df.loc[str(p)] for p in date_range]) new_df = pd.concat(results)
计算函数代码(lib.targetcalculations)
def calculate_target(df): # 执行计算并返回单列DataFrame # 为简化此处返回模拟数据 return pd.DataFrame(data = df['column1'], index=df.index)
问题现象
- 处理1-2年的小数据切片时,代码可成功执行,每年耗时约50秒,所有核心均被充分利用。
- 处理四年全量数据时,程序会挂起:初始核心全负载,随后核心使用率下降,程序却未完成,且内存使用未失控。
- 单独处理每一年的数据均可成功,说明计算函数无数据错误。
- 运行环境为Jupyter,
targetcalculations函数位于独立.py文件中,已尝试修改启动方式为spawn,问题依旧。
排查与解决方案
核心问题分析
全量数据下任务数量过多(四年约1460个日期任务),multiprocessing.Pool.map会一次性将所有任务提交到进程池,导致进程池的任务队列过载,加上Jupyter的进程管理特性,容易出现任务调度死锁或进程僵死。
优化方案
分批次提交任务
用imap或imap_unordered替代map,分批迭代任务,避免一次性加载所有子DataFrame到内存,同时缓解进程池调度压力:if __name__ == '__main__': n_processes = 8 with multiprocessing.Pool(n_processes) as p: results = [] # chunksize设为总任务数/进程数的2-3倍,四年数据可设30-50 for res in p.imap_unordered(calculate_target, [df.loc[str(p)] for p in date_range], chunksize=30): results.append(res) new_df = pd.concat(results)优化任务拆分粒度
按周或按月拆分任务,减少总任务数量,降低进程调度开销:# 按周拆分示例 date_range = df.reset_index().datetime.map(lambda t: t.isocalendar()[:2]).unique() weekly_slices = [df.loc[str(week[0]) + '-' + str(week[1])] for week in date_range]规避Jupyter环境干扰
Jupyter主进程的IO和状态管理逻辑复杂,容易和多进程调度冲突。建议将并行代码封装成独立.py脚本,在终端直接运行:python your_processing_script.py添加任务超时监控
用apply_async配合超时机制,定位可能僵死的任务:if __name__ == '__main__': n_processes = 8 with multiprocessing.Pool(n_processes) as p: results = [] for date in date_range: # 给每个任务设置超时,比如60秒 res = p.apply_async(calculate_target, args=(df.loc[str(date)],), timeout=60) results.append(res) new_df = pd.concat([res.get() for res in results])优化数据传递方式
原代码提前生成所有子DataFrame会产生额外内存拷贝,改为在计算函数中按需读取当日数据(如果原始数据存在文件中):# 修改计算函数,传入日期字符串而非子DataFrame def calculate_target(date_str): df_day = pd.read_parquet('your_data.parquet', filters=[('datetime', '>=', date_str), ('datetime', '<', f"{date_str} 23:59:59")]) return pd.DataFrame(data=df_day['column1'], index=df_day.index) # 主进程仅传递日期字符串 results = p.map(calculate_target, [str(p) for p in date_range])
内容的提问来源于stack exchange,提问作者Mike Tomaino
相关产品推荐
相关产品推荐

