如何基于另一列动态设置rolling_sum的window_size参数?
问题:Polars中如何基于另一列动态设置rolling_sum的window_size?
现有如下Polars DataFrame,其中calculate列是期望结果——它通过滚动求和得到,且滚动窗口的大小取自periods列(而非固定值):
import polars as pl df = pl.DataFrame( { "date": pl.date_range( pl.date(2023, 2, 1), pl.date(2023, 2, 5), interval="1d", eager=True), "periods": [2, 2, 2, 1, 1], "quantity": [10, 12, 14, 16, 18], "calculate": [22, 26, 30, 16, 18] } )
当设置固定window_size=2时,执行以下代码可正常运行:
df.select(pl.col("quantity").rolling_sum(window_size=2))
但尝试将window_size设为pl.col("periods")时,出现错误:
TypeError: argument 'window_size': 'Expr' object cannot be converted to 'PyString'
解决方案
Polars原生的rolling系列方法不支持动态窗口大小(即每行使用不同的窗口尺寸),可以通过以下两种方式实现需求:
方法1:逐行处理(简单直观)
使用map_elements结合列表切片,针对每行计算对应窗口内的求和:
df = df.with_columns( pl.struct(["periods", "quantity"]) .map_elements( lambda x: sum(df["quantity"][max(0, x.index - x["periods"] + 1):x.index + 1]), return_dtype=pl.Int64 ) .alias("dynamic_rolling_sum") ) # 验证结果 print(df[["calculate", "dynamic_rolling_sum"]])
方法2:向量化实现(性能更优)
通过交叉连接+过滤的向量化操作实现,避免逐行处理,适合大数据量场景:
# 添加行索引 df = df.with_row_index("idx") # 交叉连接后筛选符合窗口范围的行,再分组求和 result_df = df.join(df, how="cross", suffix="_right").filter( pl.col("idx_right").is_between( pl.col("idx") - pl.col("periods") + 1, pl.col("idx") ) & (pl.col("idx_right") >= 0) ).group_by("idx").agg( pl.col("quantity_right").sum().alias("dynamic_rolling_sum") ).drop("idx").join(df, how="right").drop("idx") # 验证结果 print(result_df[["calculate", "dynamic_rolling_sum"]])
两种方法最终得到的dynamic_rolling_sum列都会和期望的calculate列完全一致。
内容的提问来源于stack exchange,提问作者anerjee
相关产品推荐
相关产品推荐

