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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 10:28:35