如何并行分区大型Polars DataFrame并保存为CSV文件?
分组并行保存Polars子DataFrame的方案
Polars原生并行支持?
目前Polars没有提供直接的原生API来并行完成分组后保存子DataFrame的操作,partition_by方法返回的是串行迭代器,默认只能逐个处理分组。
Python原生并行实现
可以利用Python标准库的concurrent.futures模块,通过进程池或线程池实现并行保存,以下是两种常用方案:
方案1:进程池(ProcessPoolExecutor)
适合CPU或内存密集型场景,避免GIL限制,每个进程拥有独立内存空间,减少内存竞争:
import polars as pl from concurrent.futures import ProcessPoolExecutor def save_group(group_df): # 获取分组标识 group1 = group_df[0, "group1"] group2 = group_df[0, "group2"] # 保存到CSV group_df.write_csv(f"~/{group1}_{group2}.csv") if __name__ == "__main__": # 加载大型DataFrame df = pl.read_parquet("your_large_data.parquet") # 预先将分组转为列表(适合内存足够的情况) grouped_dfs = list(df.partition_by(["group1", "group2"])) # 启动进程池并行处理 with ProcessPoolExecutor() as executor: executor.map(save_group, grouped_dfs)
方案2:线程池(ThreadPoolExecutor)
适合IO密集型场景(如写文件),无需提前加载所有分组到内存,节省内存资源:
import polars as pl from concurrent.futures import ThreadPoolExecutor def save_group(group_df): group1 = group_df[0, "group1"] group2 = group_df[0, "group2"] group_df.write_csv(f"~/{group1}_{group2}.csv") if __name__ == "__main__": df = pl.read_parquet("your_large_data.parquet") # 直接使用partition_by返回的迭代器,逐个处理分组 grouped_iterator = df.partition_by(["group1", "group2"]) # 启动线程池,可根据IO能力调整max_workers with ThreadPoolExecutor(max_workers=8) as executor: executor.map(save_group, grouped_iterator)
注意事项
- Windows环境下必须将主逻辑放在
if __name__ == "__main__":块内,否则多进程会报错 - 线程池的
max_workers可根据磁盘IO性能调整,一般设置为8-16即可 - 确保分组标识组合唯一,避免不同分组写入同一文件导致数据覆盖
- 如果分组数量极大且单分组数据量也大,优先选择线程池,减少内存占用
内容的提问来源于stack exchange,提问作者lebesgue
相关产品推荐
相关产品推荐

