如何高效并行按分组将DataFrame写入Parquet文件?
问题描述
我的工作流是处理并清洗多源数据,再将数据存入多个文件夹——每个文件夹对应模型的单个输入数据。例如,若有3个数据源A、B、C和1000个数据点,文件夹结构会包含命名为1-1000的文件夹,每个文件夹内有a.parquet、b.parquet和c.parquet文件,每个文件对应大表中属于单个数据点的行。
我需要对大型DataFrame A按datapoint_id分组,将每个分组写入对应文件夹下的a.parquet,但当前的串行方法速度极慢:
for name, gdf in A.groupby(["datapoint_id"]): gdf.write_parquet(f"{name}/a.parquet")
当前速度约为25次/秒,处理含20万分组的DataFrame需数小时,且CPU利用率仅约总算力的25%,未充分利用计算资源。
高性能解决方案
1. Dask并行分组写入
Dask原生支持并行处理大型数据集,可自动利用多核CPU:
import dask.dataframe as dd import os # 将pandas DataFrame转为Dask DataFrame,分区数建议匹配CPU核心数 ddf = dd.from_pandas(A, npartitions=os.cpu_count()) def write_group(group): datapoint_id = group["datapoint_id"].iloc[0] os.makedirs(str(datapoint_id), exist_ok=True) group.to_parquet(f"{datapoint_id}/a.parquet") # 并行执行分组写入 ddf.groupby("datapoint_id").apply(write_group, meta=object).compute()
Dask会自动拆分任务到多个核心,大幅提升CPU利用率。
2. Python多进程手动实现并行
适合需要精细控制并行逻辑的场景:
import os from multiprocessing import Pool # 先将分组转为可迭代的元组列表 groups = list(A.groupby("datapoint_id")) def write_single_group(group_tuple): name, gdf = group_tuple os.makedirs(name, exist_ok=True) gdf.to_parquet(f"{name}/a.parquet") # 进程数设为CPU核心数 with Pool(processes=os.cpu_count()) as pool: pool.map(write_single_group, groups)
注意:若DataFrame过大,list(A.groupby(...))会占用较多内存,建议结合分块处理。
3. PySpark分布式处理(超大数据场景)
若数据集达到TB级,PySpark的分布式能力更适配:
from pyspark.sql import SparkSession import os spark = SparkSession.builder.appName("GroupParquetWrite").getOrCreate() spark_df = spark.createDataFrame(A) def write_parquet_group(iterator): for df in iterator: datapoint_id = df["datapoint_id"].iloc[0] os.makedirs(str(datapoint_id), exist_ok=True) df.to_parquet(f"{datapoint_id}/a.parquet") # 按分组并行写入 spark_df.groupBy("datapoint_id").mapInPandas(write_parquet_group, schema="").count()
PySpark会将任务分发到集群节点执行,适合超大规模数据处理。
4. Swifter轻量并行优化
无需额外学习分布式框架,自动适配并行:
import swifter import os def write_row_group(gdf): name = gdf.name os.makedirs(name, exist_ok=True) gdf.to_parquet(f"{name}/a.parquet") A.groupby("datapoint_id").swifter.apply(write_row_group)
Swifter会自动判断场景,底层调用Dask或多进程实现并行。
内容的提问来源于stack exchange,提问作者Udit Ranasaria
相关产品推荐
相关产品推荐

