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

如何避免重复执行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.

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 02:51:19