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

使用Multiprocessing加速df.apply()无效,内核超时问题求助

问题分析与解决方案

核心问题

你的多进程方案失效的主要原因有两个:

  1. 大DataFrame的进程间拷贝开销:每个子进程都会完整复制flight_info_df,内存占用暴增,且序列化/反序列化的时间远超过并行带来的收益。
  2. 行级apply本身的低效:get_info_previous_flight里的逐行loc查询是Python循环,属于pandas的反模式,就算并行也无法掩盖这种低效。

单进程3分钟的耗时,本质是行级循环的问题,不是CPU利用率的问题——多进程反而因为数据拷贝拖慢了整体速度。

最优解决方案:用向量化操作替代行级循环

直接用groupby + shift实现需求,这是pandas原生的高效操作,不需要多进程就能把耗时压缩到几秒级别:

import pandas as pd

# 先按飞机和日期分组,对n_flight_of_day排序(确保顺序正确)
flight_info_df = flight_info_df.sort_values(['ac_registration', 'dep_sched_date', 'n_flight_of_day'])

# 定义映射关系
previous_flight_info = {
    'change_reason_code_previous_flight': 'change_reason_code',
    'dep_delay_previous_flight': 'dep_delay',
    'act_trans_time_previous_flight': 'trans_time',
    'sched_trans_time_previous_flight': 'sched_trans_time',
    'act_groundtime_previous_flight': 'Act Groundtime',
    'sched_groundtime_previous_flight': 'Sched Groundtime'
}

# 分组后shift获取前一班航班的数据
grouped = flight_info_df.groupby(['ac_registration', 'dep_sched_date'])
for new_col, original_col in previous_flight_info.items():
    # shift(1)表示取同组的上一行数据
    flight_info_df[new_col] = grouped[original_col].shift(1)
    # 对n_flight_of_day=1的行填充0(和原逻辑一致)
    flight_info_df.loc[flight_info_df['n_flight_of_day'] == 1, new_col] = 0

# 如果需要同步leg_no的前值,同样处理
flight_info_df['leg_no_previous'] = grouped['leg_no'].shift(1)
flight_info_df.loc[flight_info_df['n_flight_of_day'] == 1, 'leg_no_previous'] = 0

若仍需多进程优化(仅当数据量极大时)

如果数据量达到千万级以上,向量化操作后还需要并行,要解决数据拷贝问题:

  1. 用共享内存传递DataFrame:使用pandas.DataFrame.share_memory()让子进程共享同一份内存数据,避免拷贝。
  2. 拆分任务为数据块而非列:把DataFrame按分组拆分成多个子块,分配给不同进程处理,而不是按列处理。

示例代码(简化版):

import multiprocessing as mp
import pandas as pd

def process_chunk(chunk, previous_flight_info):
    chunk = chunk.sort_values(['ac_registration', 'dep_sched_date', 'n_flight_of_day'])
    grouped = chunk.groupby(['ac_registration', 'dep_sched_date'])
    for new_col, original_col in previous_flight_info.items():
        chunk[new_col] = grouped[original_col].shift(1)
        chunk.loc[chunk['n_flight_of_day'] == 1, new_col] = 0
    return chunk

if __name__ == '__main__':
    # 拆分DataFrame为多个块
    chunks = [flight_info_df[i:i+10000] for i in range(0, len(flight_info_df), 10000)]
    previous_flight_info = {
        'change_reason_code_previous_flight': 'change_reason_code',
        # ... 其他映射
    }
    
    # 使用共享内存(可选,进一步减少拷贝)
    flight_info_df.share_memory()
    
    with mp.Pool() as pool:
        results = pool.starmap(process_chunk, [(chunk, previous_flight_info) for chunk in chunks])
    
    # 合并结果
    final_df = pd.concat(results)

为什么你的原方案不行?

  • tasks里每个元组都包含flight_info_df,每个进程都会复制整个DataFrame,内存占用是单进程的N倍(N是进程数),导致系统卡顿甚至超时。
  • apply(axis=1)本身是逐行循环,就算并行,每个进程里的循环还是慢,且多进程的调度开销抵消了收益。

内容的提问来源于stack exchange,提问作者KaiHoogenhoud

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 15:32:51