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

如何高效并行按分组将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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 17:55:12