如何使用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
相关产品推荐
相关产品推荐

