You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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的进程管理特性,容易出现任务调度死锁或进程僵死。

优化方案

  1. 分批次提交任务
    用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)
    
  2. 优化任务拆分粒度
    按周或按月拆分任务,减少总任务数量,降低进程调度开销:

    # 按周拆分示例
    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]
    
  3. 规避Jupyter环境干扰
    Jupyter主进程的IO和状态管理逻辑复杂,容易和多进程调度冲突。建议将并行代码封装成独立.py脚本,在终端直接运行:

    python your_processing_script.py
    
  4. 添加任务超时监控
    用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])
    
  5. 优化数据传递方式
    原代码提前生成所有子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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.16 05:02:39