SCOOP框架下如何让子进程等待主进程完成计算后再执行
SCOOP框架下主进程预处理后子进程再执行的解决方案
问题背景
- 环境:Python 3.6 + SCOOP框架,无法升级版本
- 需求:所有子进程先完成一段前置计算,等待主进程在
if __name__ == '__main__':块内完成耗时的DataFrame预处理后,再用主进程生成的最终DataFrame执行后续计算 - 当前问题:SCOOP启动子进程后,子进程会直接执行
if __name__ == '__main__':外部的所有全局代码,此时主进程还没完成DataFrame预处理,导致子进程使用未更新的DataFrame报错 - 尝试过的方法:
scoop.futures.map、scoop.futures.supply、multiprocessing.managers、multiprocessing.Barrier(8).wait()均未解决,了解scoop.futures.wait()但不知如何获取futures参数
原代码示例
import pandas as pd import genetic_algorithm from scoop import futures df = pd.read_csv('database.csv') # 数据量太大,不想每次给worker传副本,希望每个worker都有一份 if __name__ == '__main__': df = add_new_columns(df) # 耗时计算,只想执行一次,不想所有worker都跑 df = computation_using_new_columns(df) # <--- 报错:这行在add_new_columns完成前就被子进程执行了 def fitness_function(): ... # 所有worker都要用这个函数,放if块里会报错 if __name__ == '__main__': results = list(futures.map(genetic_algorithm, df))
执行命令:python3 -m scoop script.py
解决方案
核心思路是让主进程先完成DataFrame预处理,再将最终数据同步给所有子进程,同时控制子进程在拿到最终数据后才执行依赖代码。
修改后的代码
import pandas as pd import genetic_algorithm from scoop import futures # 全局变量存储最终处理后的DataFrame,子进程共用 final_df = None def init_worker(processed_df): """子进程初始化:接收主进程传递的最终DataFrame,同时执行前置计算""" global final_df final_df = processed_df # 这里写子进程需要先执行的前置计算逻辑 pre_computation() def pre_computation(): """子进程的前置计算""" # 例如:预加载某些资源、执行独立于final_df的计算 pass def computation_using_new_columns(): """依赖final_df的后续计算,必须在子进程拿到final_df后执行""" global final_df # 替换成你的实际计算逻辑 return final_df.apply(lambda x: x['new_col'] * 2) def fitness_function(): """子进程使用的适应度函数,依赖final_df和后续计算结果""" processed_data = computation_using_new_columns() # 替换成你的实际适应度计算逻辑 return processed_data.sum() if __name__ == '__main__': # 1. 初始读取原始数据 raw_df = pd.read_csv('database.csv') # 2. 主进程单独执行耗时的预处理 processed_df = add_new_columns(raw_df) # 3. 将最终DataFrame分发到所有子进程,完成初始化和前置计算 # supply会确保所有子进程完成初始化后,才处理后续任务 init_futures = futures.supply(init_worker, processed_df) # 等待所有子进程初始化完成(可选,确保万无一失) futures.wait(init_futures) # 4. 启动遗传算法任务,子进程此时已拥有final_df # 注意:这里不要直接传processed_df,改用索引或轻量参数,避免重复传递大数据 results = list(futures.map(genetic_algorithm, range(len(processed_df))))
关键修改说明
- 避免全局代码提前执行:把原来全局执行的
computation_using_new_columns(df)改成函数,依赖全局变量final_df,只有子进程拿到最终数据后才会调用执行 - 子进程初始化同步:用
init_worker函数接收主进程传递的最终DataFrame,同时在里面完成子进程的前置计算;通过futures.supply分发初始化任务,确保所有子进程都同步到最新数据 - 控制任务执行顺序:用
futures.wait(init_futures)等待所有子进程完成初始化,再启动后续的遗传算法任务,彻底避免时序问题 - 优化数据传递:
futures.map不再直接传递大DataFrame,改用索引等轻量参数,子进程通过全局的final_df获取数据,减少内存消耗
补充注意事项
- 由于Python 3.6的序列化限制,传递大DataFrame时确保数据可被pickle序列化;如果数据量极大,可以考虑主进程将处理后的DataFrame写入临时文件,子进程在
init_worker中读取文件加载数据 fitness_function必须定义在全局作用域,SCOOP子进程需要能找到并导入该函数,不能放在if __name__ == '__main__':块内- 如果
futures.supply传递数据仍有问题,可以尝试SCOOP的shared模块(若可用)实现数据共享,进一步减少内存开销
内容的提问来源于stack exchange,提问作者João Bravo
相关产品推荐
相关产品推荐

