Pandarallel首个Worker运行极慢问题排查求助
Pandarallel首个Worker任务进度过慢的解决办法
问题场景
使用pandarallel对pandas DataFrame应用自定义函数时,8个Worker中首个运行异常缓慢,其余7个可在1分钟内完成任务,但首个Worker耗时数小时,无阻塞无报错,导致其余Worker闲置等待。运行日志如下:
INFO: Pandarallel will run on 8 workers. INFO: Pandarallel will use standard multiprocessing data transfer (pipe) to transfer data between the main process and workers. 17.14% :::::: | 170 / 992 | 100.00% :::::::::::::::::::::::::::::::::::::::: | 992 / 992 | 100.00% :::::::::::::::::::::::::::::::::::::::: | 992 / 992 | 100.00% :::::::::::::::::::::::::::::::::::::::: | 992 / 992 | 100.00% :::::::::::::::::::::::::::::::::::::::: | 992 / 992 | 100.00% :::::::::::::::::::::::::::::::::::::::: | 992 / 992 | 100.00% :::::::::::::::::::::::::::::::::::::::: | 991 / 991 | 100.00% :::::::::::::::::::::::::::::::::::::::: | 991 / 991 |
相关代码:
import pandas as pd from pandarallel import pandarallel pandarallel.initialize(nb_workers=8, progress_bar=True) def test_function(): # 自定义函数逻辑 pass df_temp = pd.DataFrame({'page': range(1, 5)}) df = pd.DataFrame(list(df_temp.page.parallel_apply(test_function)))
解决办法
手动拆分任务,均衡分配
Pandarallel默认的任务拆分可能出现分配不均,可手动将数据拆分为等量的N份(N为Worker数量),再并行处理后合并结果:import pandas as pd import numpy as np from multiprocessing import Pool def test_function(page): # 自定义函数逻辑 pass df_temp = pd.DataFrame({'page': range(1, 5)}) # 手动拆分数据为8份 splits = np.array_split(df_temp.page, 8) # 用Pool并行处理拆分后的数据集 with Pool(8) as pool: results = [] for split in splits: results.extend(pool.map(test_function, split)) df = pd.DataFrame(results)更换数据传输模式
日志显示使用pipe传输数据,换成共享内存模式可降低传输开销,适合大数据量场景:pandarallel.initialize(nb_workers=8, progress_bar=True, use_memory_fs=True)提前加载Worker依赖资源
若自定义函数需要加载模型、配置等资源,首个Worker可能因首次加载耗时久,可通过初始化钩子让所有Worker启动时提前加载:def init_worker(): # 全局加载资源,每个Worker启动时执行一次 global resource resource = load_your_resource() pandarallel.initialize(nb_workers=8, progress_bar=True, init_worker=init_worker)调整Worker数量
尝试减少Worker数量(比如设为4),过多的Worker可能增加调度开销,尤其在数据量不大时,反而会导致任务分配失衡。
内容的提问来源于stack exchange,提问作者diggi2395
相关产品推荐
相关产品推荐

