如何在Python中并行处理Polars partition_by生成的DataFrame列表?
加速Polars分片DataFrame的并行处理方案
方案1:使用concurrent.futures.ProcessPoolExecutor
适合CPU密集型任务,基于Python标准库实现,无需额外安装依赖。Polars的DataFrame分片支持进程间序列化传递,直接用map方法批量处理分组即可。
示例代码:
import polars as pl import numpy as np from concurrent.futures import ProcessPoolExecutor def processing_df(df): # 替换为你的复杂处理逻辑 return np.array(df.select(pl.col("value").sum()).to_numpy()) # 示例源DataFrame df_sourse = pl.DataFrame({ "group_col": ["A", "A", "B", "B", "C"], "value": [1, 2, 3, 4, 5] }) df_gr = df_sourse.partition_by(groups="group_col", maintain_order=True) # 并行执行分组处理 with ProcessPoolExecutor(max_workers=4) as executor: results = list(executor.map(processing_df, df_gr)) print(results)
方案2:使用multiprocessing.Pool
和上述方案逻辑一致,适合习惯使用multiprocessing模块的场景,Windows系统必须加if __name__ == "__main__":判断,避免进程启动异常。
示例代码:
import polars as pl import numpy as np import multiprocessing as mp def processing_df(df): # 替换为你的复杂处理逻辑 return np.array(df.select(pl.col("value").sum()).to_numpy()) if __name__ == "__main__": df_sourse = pl.DataFrame({ "group_col": ["A", "A", "B", "B", "C"], "value": [1, 2, 3, 4, 5] }) df_gr = df_sourse.partition_by(groups="group_col", maintain_order=True) with mp.Pool(processes=mp.cpu_count()) as pool: results = pool.map(processing_df, df_gr) print(results)
进程未正常运行的排查要点
- Windows系统必须加
if __name__ == "__main__"::Windows的进程启动机制会重新导入模块,不加该判断会导致无限创建子进程或函数无法被正确调用。 - 函数序列化问题:
processing_df内部不能引用无法被pickle序列化的对象(如未定义在模块级别的内部函数、特殊自定义类实例)。若存在此类情况,将相关逻辑移至模块级别,或安装cloudpickle库替换默认序列化器。 - Polars版本兼容:确保使用0.18.x及以上版本的Polars,旧版本可能存在DataFrame进程间传递的bug。
- 进程数限制:若分组数量极大,不要开启过多进程,建议设置
max_workers为CPU核心数(避免资源耗尽)。
额外优化建议
- 如果
processing_df返回的结果可合并为Polars DataFrame,处理完成后可直接拼接:
results_df = pl.concat([pl.from_numpy(res) for res in results])
- 若处理逻辑为IO密集型(如分片读写文件),改用
ThreadPoolExecutor可减少进程切换开销。
内容的提问来源于stack exchange,提问作者Михаил Губаренко
相关产品推荐
相关产品推荐

