如何避免重复执行group_by_dynamic?改写Expr实现高效聚合
解决方案
核心思路
先一次性完成固定规则的group_by_dynamic分组,将每组需要聚合的字段转为列表存储;后续所有聚合操作直接基于这些列表列,用Polars原生的列表Expr方法实现矢量化计算,彻底避免重复分组的开销。
步骤1:预处理——一次性分组并转列表
假设你的时间列为timestamp,固定分组规则为每1小时,需要聚合的字段为value1、value2、weight:
import polars as pl # 加载原始大型DataFrame(确保已按时间排序) raw_df = pl.read_parquet("large_dataset.parquet").sort("timestamp") # 仅执行一次动态分组,将目标字段转为列表 grouped_df = raw_df.group_by_dynamic( index_column="timestamp", every="1h", # 固定时间分组规则 period="1h" ).agg( pl.col("value1").list().alias("value1_list"), pl.col("value2").list().alias("value2_list"), pl.col("weight").list().alias("weight_list") )
步骤2:循环中用列表Expr实现动态聚合
在非梯度求解器的循环里,直接基于预处理后的grouped_df,用Polars列表相关的Expr完成各种参数化聚合,全程无需再次分组:
# 示例:循环处理不同聚合参数 for param in [0.3, 0.5, 0.7, 0.9]: # 用Polars Expr实现多种动态聚合 result = grouped_df.with_columns( # 带参数的加权平均值 (pl.col("value1_list") .list.zip_with(pl.col("weight_list")) .list.eval(pl.element().first() * pl.element().second() * param) .list.sum() / pl.col("weight_list").list.sum() ).alias(f"weighted_avg_param_{param}"), # 动态分位数计算 pl.col("value2_list").list.quantile(param).alias(f"quantile_{param}"), # 筛选大于param的值后的均值 pl.col("value1_list").list.filter(pl.element() > param).list.mean().alias(f"filtered_mean_{param}"), # 自定义滑动累计和(基于有序列表) pl.col("value1_list").list.cumsum().alias(f"cumsum_list") ) # 后续处理结果(如传入求解器) feed_to_solver(result)
关键注意事项
- 确保原始DataFrame已按时间列排序:
group_by_dynamic默认保留分组内的原始顺序,排序后列表的时间顺序才符合预期 - 所有列表操作均为Polars矢量化实现,性能远高于
apply或Python层面的循环 - 多字段联合聚合可通过
list.zip_with将多个列表合并为元组列表,再用list.eval对每个元组执行逻辑
内容的提问来源于stack exchange,提问作者Andrew P.
相关产品推荐
相关产品推荐

