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

如何使用Python对Parquet文件读取的DataFrame按指定字段分组聚合求和

解决方案

针对大型Parquet文件的分组聚合需求,可根据文件大小选择对应实现方案:

PySpark实现(适配TB级超大型文件,无需全量加载内存)

from pyspark.sql import SparkSession
from pyspark.sql.functions import sum

# 初始化Spark会话
spark = SparkSession.builder.appName("ParquetAggregate").getOrCreate()

# 读取源Parquet文件
df = spark.read.parquet("path/to/your/source_file.parquet")

# 按指定字段分组聚合
agg_df = df.groupBy("client_identifier_1", "client_identifier_2", "product_identifier") \
           .agg(sum("value_a").alias("value_a_total"), sum("value_b").alias("value_b_total"))

# 输出聚合结果
agg_df.write.parquet("path/to/your/aggregated_result.parquet")

spark.stop()

Pandas实现(适配文件大小不超过可用内存70%的场景)

import pandas as pd

# 可指定columns参数仅读取需要的字段,降低内存占用
df = pd.read_parquet(
    "path/to/your/source_file.parquet",
    columns=["client_identifier_1", "client_identifier_2", "product_identifier", "value_a", "value_b"]
)

# 分组求和
agg_df = df.groupby(
    ["client_identifier_1", "client_identifier_2", "product_identifier"],
    as_index=False
)[["value_a", "value_b"]].sum()

# 保存结果
agg_df.to_parquet("path/to/your/aggregated_result.parquet")

DuckDB实现(内存占用远低于Pandas,适配中等大小的文件)

import duckdb

# 直接通过SQL查询聚合,无需全量加载数据
conn = duckdb.connect()
agg_df = conn.execute("""
    SELECT 
        client_identifier_1,
        client_identifier_2,
        product_identifier,
        SUM(value_a) AS value_a_total,
        SUM(value_b) AS value_b_total
    FROM read_parquet('path/to/your/source_file.parquet')
    GROUP BY 1,2,3
""").df()

agg_df.to_parquet("path/to/your/aggregated_result.parquet")

内容的提问来源于stack exchange,提问作者Igor Eulálio

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 10:24:04