如何将Parquet分区数据集直接读取到Polars进行聚合计算
解决方案:直接用Polars处理大型分区Parquet数据集(无需转Pandas)
针对你的需求,这里提供两种高效的解决方案,完全避免将数据转换为Pandas格式,同时支持流式处理以应对内存不足的场景:
方案1:直接用Polars扫描Parquet数据集(推荐)
Polars原生支持读取分区Parquet数据集,pl.scan_parquet()会创建懒加载的LazyFrame,仅在执行collect()时才会处理数据,配合streaming=True可分批次处理数据,无需加载全量数据到内存。
修改后的完整代码:
import polars as pl # Version 0.20.3 import pyarrow as pa # Version 11.0.0 import pyarrow.parquet as pq # 生成测试数据并写入分区Parquet pl_df = pl.DataFrame({ "Name": ["ABC","DEF","GHI",'JKL'], "date": ["2024-01-01","2024-01-10","2023-01-29","2023-01-29"], "price":[1000,1500,1800,2100] , }) pl_df = pl_df.with_columns(date= pl.col("date").cast(pl.Date)) pq.write_to_dataset( pl_df.to_arrow(), root_path=r"C:\Users\desktop PC\Downloads\test_pl", partition_cols=["date"], compression='gzip', existing_data_behavior='overwrite_or_ignore' ) # 直接扫描分区Parquet并完成聚合 df = ( pl.scan_parquet(r"C:\Users\desktop PC\Downloads\test_pl", partition_cols=["date"]) .group_by(["date"]) .agg( pl.col("price").sum().alias("grouped_sum"), pl.col("price").count().alias("grouped_count"), ) .collect(streaming=True) ) print(df)
- 若需要指定schema,可添加
schema参数,例如schema={"date": pl.Date, "Name": pl.String, "price": pl.Int64},无需通过Pandas生成schema。 - 流式处理模式会自动拆分数据块,大幅降低内存占用。
方案2:从Arrow Table直接转Polars(若已用pyarrow读取)
如果已经通过pyarrow获取了pq_df(Arrow Table对象),可直接用pl.from_arrow()转换为Polars对象,跳过Pandas转换步骤:
# 读取Parquet数据集到Arrow Table pq_df = pq.read_table(r"C:\Users\desktop PC\Downloads\test_pl", schema=pd_df_schema) # 直接从Arrow Table转Polars并聚合 df = ( pl.from_arrow(pq_df).lazy() .group_by(["date"]) .agg( pl.col("price").sum().alias("grouped_sum"), pl.col("price").count().alias("grouped_count"), ) .collect(streaming=True) ) print(df)
pl.from_arrow()支持直接转换Arrow Table/RecordBatch,性能远高于Pandas中转,且不会额外占用内存。- 同样配合懒加载和流式处理,适合大数据场景。
内容的提问来源于stack exchange,提问作者Klllmmm
相关产品推荐
相关产品推荐

