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

如何将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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 20:35:21