使用Polars懒API聚合多Parquet文件时按文件添加日期列
用Polars懒API实现按文件统一时间列的聚合
完全可以用Polars懒API实现需求,不需要提前修改原始Parquet文件,全程基于懒加载处理大数据量。核心思路是给每个文件添加标识列,再通过窗口函数提取该文件的首个date值并广播到整列,具体代码如下:
( pl.scan_parquet('data/data-16828*.parquet') # 为每行添加所属文件的路径标识 .map_batches(lambda df: df.with_columns(pl.lit(df.meta.source).alias("file_path"))) # 用当前文件第一行的date值填充整列,生成统一的时序标识列 .with_columns(pl.first("date").over("file_path").alias("file_date")) # 按目标维度 + 统一时序列分组聚合 .groupby(['type_id', 'location_id', 'file_date']) .agg([ pl.min('n').alias('n_min'), pl.max('n').alias('n_max') ]) .collect() )
关键步骤说明:
- 添加文件标识列:通过
map_batches结合df.meta.source获取每个批次(对应单个Parquet文件)的路径,新增file_path列作为文件唯一标识。 - 生成统一时序列:使用窗口函数
over("file_path"),提取每个文件分组下的第一个date值,自动广播到该文件的所有行,得到file_date列——这就是你需要的、每个文件统一的时序标识。 - 时序聚合:基于
file_date和原有维度分组,完成聚合计算。
整个流程全程采用懒执行模式,Polars会在collect()阶段才实际处理数据,且自动并行处理多个文件,完全适配超内存的大数据场景。
内容的提问来源于stack exchange,提问作者basse
相关产品推荐
相关产品推荐

