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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 16:02:43