Polars中如何使用sliding window实现行按批次分组聚合
固定大小滑动窗口行聚合实现
需求场景
需要对有序数据表实现窗口大小为3的滑动批次聚合,规则如下:
- 累计行数不足3行时,聚合列返回
None - 累计行数≥3时,聚合列返回包含当前行及前2行整行数据的嵌套列表
原始输入表
id x y ---- ---- ---- 1 x1 y1 2 x2 y2 3 x3 y3 4 x4 y4 5 x5 y5
预期输出
id x y Grouped x3 ---- ---- ---- ----------------------------------------- 1 x1 y1 None 2 x2 y2 None 3 x3 y3 [[1, x1, y1], [2, x2, y2], [3, x3, y3]] 4 x4 y4 [[2, x2, y2], [3, x3, y3], [4, x4, y4]] 5 x5 y5 [[3, x3, y3], [4, x4, y4], [5, x5, y5]]
已完成中间步骤
通过concat_list方法已经实现单行三列聚合为列表的中间结果:
id x y List ---- ---- ---- ------------- 1 x1 y1 [1, x1, y1] 2 x2 y2 [2, x2, y2] 3 x3 y3 [3, x3, y3] 4 x4 y4 [4, x4, y4] 5 x5 y5 [5, x5, y5]
具体实现
基于已生成的单行List列,不需要写自定义遍历逻辑,直接调用框架内置的滚动窗口聚合能力即可实现需求,以下是和concat_list同属Polars框架的实现代码:
import polars as pl result = ( # 初始化原始数据 pl.DataFrame({ "id": [1, 2, 3, 4, 5], "x": ["x1", "x2", "x3", "x4", "x5"], "y": ["y1", "y2", "y3", "y4", "y5"] }) # 已完成的单行列表聚合步骤 .with_columns(pl.concat_list(["id", "x", "y"]).alias("List")) # 滑动窗口聚合:窗口大小3,范围为当前行+前2行 .with_columns( pl.col("List") .rolling_list(window_size=3, offset=-2) .alias("Grouped x3") ) # 前2行窗口长度不足3,值替换为None .with_columns( pl.when(pl.col("Grouped x3").list.len() == 3) .then(pl.col("Grouped x3")) .otherwise(None) .alias("Grouped x3") ) # 不需要中间List列可直接删除 .drop("List") )
代码执行后的输出和预期结果完全一致。
如果使用SQL类计算引擎(Spark、Flink等),逻辑完全等价:
- 按id排序定义窗口,范围为
ROWS BETWEEN 2 PRECEDING AND CURRENT ROW - 用数组收集函数聚合窗口内的整行列表
- 判断窗口内行数为3时返回聚合结果,否则返回null
内容的提问来源于stack exchange,提问作者DustinEwan
相关产品推荐
相关产品推荐

