使用Multiprocessing加速df.apply()无效,内核超时问题求助
问题分析与解决方案
核心问题
你的多进程方案失效的主要原因有两个:
- 大DataFrame的进程间拷贝开销:每个子进程都会完整复制
flight_info_df,内存占用暴增,且序列化/反序列化的时间远超过并行带来的收益。 - 行级
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
若仍需多进程优化(仅当数据量极大时)
如果数据量达到千万级以上,向量化操作后还需要并行,要解决数据拷贝问题:
- 用共享内存传递DataFrame:使用
pandas.DataFrame.share_memory()让子进程共享同一份内存数据,避免拷贝。 - 拆分任务为数据块而非列:把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
相关产品推荐
相关产品推荐

