如何使用polars.Expr.map_batches传递多个Series实现自定义滚动均值?
解决Polars中
map_batches()传递多个Series实现加权滚动均值的问题 核心要点:正确构造exprs序列并编写对应处理逻辑
你忽略的关键是**map_batches()的第一个参数exprs需要传入多列的表达式/列名序列**,对应的处理函数会按顺序接收包含这些列对应Series的序列。以下是具体实现步骤:
1. 编写自定义加权滚动均值处理函数
函数需接收数值列和权重列两个Series,对每个滚动窗口计算加权平均:
import polars as pl from polars import Series def weighted_rolling_mean(nums: Series, weights: Series, window_size: int) -> Series: # 利用Polars原生rolling窗口逐段计算加权平均 return nums.rolling(window_size).apply( lambda win_num, win_wt: (win_num * win_wt).sum() / win_wt.sum(), exprs=[weights] )
2. 正确调用map_batches()
将数值列、权重列的表达式组成序列传入exprs,处理函数按顺序提取两个Series:
# 示例数据集 df = pl.DataFrame({ "value": [1, 2, 3, 4, 5], "weight": [0.1, 0.2, 0.3, 0.4, 0.5] }) window_size = 3 # 调用map_batches生成加权滚动均值列 result = df.with_columns( pl.map_batches( exprs=[pl.col("value"), pl.col("weight")], function=lambda series_pair: weighted_rolling_mean(series_pair[0], series_pair[1], window_size), return_dtype=pl.Float64 ).alias("weighted_rolling_mean") ) print(result)
3. 简化写法:直接在回调中嵌入逻辑
也可以省去单独的函数,把窗口计算逻辑直接写在map_batches的回调里:
result = df.with_columns( pl.map_batches( [pl.col("value"), pl.col("weight")], lambda s: s[0].rolling(window_size).apply( lambda win_num, win_wt: (win_num * win_wt).sum() / win_wt.sum(), exprs=[s[1]] ), return_dtype=pl.Float64 ).alias("weighted_rolling_mean") )
关键注意事项
exprs必须是列名字符串序列或列表达式序列,不能仅传入单个列- 处理函数的参数是
Sequence[Series],元素顺序和exprs中的列顺序完全对应 - 窗口大小可通过闭包变量传入处理函数,无需额外作为
map_batches的参数传递
内容的提问来源于stack exchange,提问作者Samuel Hapak
相关产品推荐
相关产品推荐

