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

如何并行分区大型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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 20:33:28