首页
/ Window Functions in Polars: Grouped Ranking, Sorting, and Aggregation with `over`

Window Functions in Polars: Grouped Ranking, Sorting, and Aggregation with `over`

2026-09-08 11:29:08作者:龚格成

本指南以 Polars 官方用户手册中 window-functions.md 为主体,系统讲解 Polars(Rust 实现的极速 DataFrame 查询引擎)中通过 .over() 实现的分组窗口表达式:在 select / with_columns 上下文中按组做排名、排序、聚合并映射回原表行,涵盖 group_to_rows / explode / join 三种结果映射策略。读完你将掌握按组打分的 over 写法、与 group_by 的取舍、三种 mapping_strategy 的语义与性能差异,以及 .over() 在 Python API 与 Rust 引擎中的实现脉络。

仓库相关文档与源码(本文所述均为当前仓库内实际内容):用户指南 expressions/index.md、Python 窗口函数参考 py-polars/docs/source/reference/expressions/window.rst

什么是 Window Function:不压缩行数的分组计算

Window function 是一类“带超能力”的表达式。普通的 group_by 聚合会把多行折叠成每组一行;而窗口函数允许你在 select 上下文里对分组执行聚合,同时保留每一行的位置,把计算结果映射回对应的行。这是它区别于普通聚合最核心的一点:

  • group_by 通常产出一个行数等于“分组个数”的 DataFrame;
  • over 通常产出一个与原始 DataFrame 行数相同的 DataFrame。

这与 SQL 中 PostgreSQL 的 window function 语义一致——Polars Python API 对 over 的 docstring 中也直接把它与 PostgreSQL 的窗口函数对照(见 py-polars/src/polars/expr/expr.pyover 方法说明)。

由于 .over() 这类窗口表达式不改变行数,它非常适合在 同一张表的列上边算分组结果边保留明细,典型场景是:组内排名、组内排序取 Top-N、组内均值/求和并广播回每行。

数据准备:加载宝可梦数据集

先加载一个宝可梦数据集。为后续演示“按类型分组”的语义,示例将 "Type 1""Type 2" 两列 cast 为 Enum 类型:

import polars as pl

types = (
    "Grass Water Fire Normal Ground Electric Psychic Fighting Bug Steel "
    "Flying Dragon Dark Ghost Poison Rock Ice Fairy".split()
)
type_enum = pl.Enum(types)
# then let's load some csv data with information about pokemon
pokemon = pl.read_csv(
    "docs/assets/data/pokemon.csv",
).cast({"Type 1": type_enum, "Type 2": type_enum})
print(pokemon.head())

该示例的完整可执行脚本位于 docs/source/src/python/user-guide/expressions/window.py;对应的 Rust 版本示例位于 docs/source/src/rust/user-guide/expressions/window.rs(Rust 侧通过 reqwest 从远程读取该 CSV)。

按组运算:对 “Speed” 列做组内排名

假设我们要给宝可梦的 “Speed”(速度)排名,但不希望做全局排名,而是希望在每个由 “Type 1” 定义的类型组内部排名。做法是:先写出对 “Speed” 列排名的表达式,再追加 .over("Type 1"),指明“在列 ‘Type 1’ 的每个唯一值上执行”:

result = pokemon.select(
    pl.col("Name", "Type 1"),
    pl.col("Speed").rank("dense", descending=True).over("Type 1").alias("Speed rank"),
)

print(result)

Python 侧 .over() 签名见 py-polars/src/polars/expr/expr.py。Rust 侧对应实现是 Expr::over,完整签名(含 mapping_strategy 的等价参数)见下文“源码层面的三种映射”一节。

背后的执行直觉:可以把 Polars 想象成先选出 “Type 1” 列取值相同的那些行子集,仅对这个子集计算排名表达式,然后把该组的结果投影回原始行;Polars 对所有存在的组重复这一过程。下图高亮了 “Type 1” 为 “Grass” 的那组宝可梦的排名计算过程:

speed_rank_by_type 示意图:按 Type 1 分组后对 Speed 计算组内排名并映射回原行

注意一个容易误解的例子:宝可梦 “Golbat” 的 “Speed” 值是 90,比 “Venusaur” 的 80 更大,但 Venusaur 却排名第 1——原因正是 Golbat 与 Venusaur 的 “Type 1” 列取值不同,二者处于不同分组,各自在组内排名。

使用多个列进行更细粒度的分组

over 接受任意数量的表达式/列名作为分组键。上面的排名也可以改为按 “Type 1” 与 “Type 2” 的组合分组,得到更细粒度的组内排名:

result = pokemon.select(
    pl.col("Name", "Type 1", "Type 2"),
    pl.col("Speed")
    .rank("dense", descending=True)
    .over("Type 1", "Type 2")
    .alias("Speed rank"),
)

print(result)

多个分组键之间是“组合并集”的关系(等价于 SQL 中 PARTITION BY "Type 1", "Type 2"),组内计算互不影响。

overgroup_by + explode 的关系与形状差异

从一般意义上说,.over() 能拿到的结果,也可以用先聚合、再 explode 实现,只不过行序会不同:

result = (
    pokemon.group_by("Type 1")
    .agg(
        pl.col("Name"),
        pl.col("Speed").rank("dense", descending=True).alias("Speed rank"),
    )
    .select(pl.col("Name"), pl.col("Type 1"), pl.col("Speed rank"))
    .explode("Name", "Speed rank")
)

print(result)

对比上面两条代码路径,可以总结二者的结果形态差异:

  • group_by 通常产出行数等于分组个数的 DataFrame(每组一行,聚合值以 List 形式存在列里);
  • over 通常产出与原始 DataFrame 行数一致的 DataFrame(聚合结果映射回每一行)。

两者的使用取舍要看意图:group_by 用于“报表型”汇总,over 用于“保留明细同时附加分组统计量”的特征工程或窗口计算。

需要特别说明的是:over 并不总是保证输出与原始 DataFrame 行数一致——这正是下面要展开的 mapping_strategy 参数所控制的。

将结果映射回 DataFrame 行:mapping_strategy

over 接受一个参数 mapping_strategy,它决定分组表达式的结果如何被映射回 DataFrame 的行。为了讲清三种策略的差异,先构造一个运动员数据框:

athletes = pl.DataFrame(
    {
        "athlete": list("ABCDEF"),
        "country": ["PT", "NL", "NL", "PT", "PT", "NL"],
        "rank": [6, 1, 5, 4, 2, 3],
    }
)
print(athletes)

数据共 6 名运动员(A–F),分属两个国家:PT(葡萄牙,3 人)与 NL(荷兰,3 人)。

group_to_rows(默认)

默认策略是 "group_to_rows":组内表达式的计算结果长度应与该组行数一致,结果按原行位置映射回该组的每一行。

下面按国籍内部对运动员的排名排序。荷兰运动员原本位于第 2、3、6 行——执行后它们仍留在这些位置;变化的是运动员姓名的顺序,从 “B”、“C”、“F” 变为 “B”、“F”、“C”:

result = athletes.select(
    pl.col("athlete", "rank").sort_by(pl.col("rank")).over(pl.col("country")),
    pl.col("country"),
)

print(result)

下图的左列展示了按国家排序前后的原始行位置,右列则对应结果——group_to_rows 不改变各组的整体行位置,只重排组内内容:

athletes_over_country 示意图:group_to_rows 策略将组内排序结果按原行位置映射回原表

explode

如果把 mapping_strategy 设为 "explode",则同一国家的运动员被集中到一起,但最终行的顺序(就国家而言)与原始顺序不再一致,正如下图所示:

athletes_over_country_explode 示意图:explode 策略不保持原行位置,同组成员被集中到相邻行

因为 Polars 无需跟踪每个组内行在原表中的位置,"explode" 通常比 "group_to_rows" 更快。但使用它需要更小心:它意味着我们想保留的其它列也必须一并重排,否则会出现列间错位。因此这里的示例是对 pl.all() 整体做 sort + over,把所有列同步重排:

result = athletes.select(
    pl.all()
    .sort_by(pl.col("rank"))
    .over(pl.col("country"), mapping_strategy="explode"),
)

print(result)

join

mapping_strategy 的另一个可选值是 "join":它先把每组聚合的结果收集成一个 List,然后把该 List 重复广播到该组的每一行

result = athletes.with_columns(
    pl.col("rank").sort().over(pl.col("country"), mapping_strategy="join"),
)

print(result)

注意此时 rank 列会被替换为一个 List 列,每行都持有其所属国家的完整排序结果。Python docstring 对此有明确警告:该策略可能非常耗费内存(见 py-polars/src/polars/expr/expr.py),因为要在每一行复制整组数据。

源码层面的三种映射

三种策略在 Rust 引擎中有精确对应。WindowMapping 枚举定义于 crates/polars-plan/src/dsl/options/mod.rs

  • GroupsToRows(默认值):把组内值映射回组内对应位置;
  • Explode:把聚合出的 List 展开并直接做横向拼接(hstack)而不是 join;要求各组有序结果才有意义;
  • Join:把各组以 List<group_dtype> 形式 join 到各行的位置上——源码注释明确警告“这可能非常消耗内存”。

物理执行侧,窗口表达式的计算结构 WindowExpr 定义在 crates/polars-expr/src/expressions/window.rs,其中记录了分组键 group_by、排序键 order_by、要施加窗口函数的列 apply_columns、物理函数与映射方式 mapping。而在该文件内部,还根据组内结果类型选择 MapStrategycrates/polars-expr/src/expressions/window.rs):Join(按 key join,对归约聚合而言最贵)、Explode(直接展开)、Map(用一次 arg_sort 把结果映射回原位置)。这正是文档所讲“explode 更快、join 最贵、group_to_rows 需保持位置”三种结论的引擎实现依据。

Rust 侧调用 mapping_strategy="explode" 的对应写法是 Expr::over_with_options(Some(partition), None, WindowMapping::Explode),示例见 docs/source/src/rust/user-guide/expressions/window.rs

窗口化聚合表达式:标量结果的广播

如果应用到某个组上的表达式最终产出的是标量值(例如均值、求和、最大值),则这个标量会被广播到该组的所有行。这是最常见的特征工程用法:为每个类型组计算平均速度,并把均值写回组内每一行:

result = pokemon.select(
    pl.col("Name", "Type 1", "Speed"),
    pl.col("Speed").mean().over(pl.col("Type 1")).alias("Mean speed in group"),
)

print(result)

此时输出行数与原始一致,新增列 “Mean speed in group” 在属于同一 “Type 1” 的行上是同一个均值。

更贴近真实建模的组合写法:同一张表里同时计算“按类型平均攻击力”、“按类型组合平均防御力”以及“全局平均攻击力”三种粒度的统计量——这正是 over 相比逐步 join 的优势所在:

result = pokemon.select(
    "Type 1",
    "Type 2",
    pl.col("Attack").mean().over("Type 1").alias("avg_attack_by_type"),
    pl.col("Defense")
    .mean()
    .over(["Type 1", "Type 2"])
    .alias("avg_defense_by_type_combination"),
    pl.col("Attack").mean().alias("avg_attack"),
)
print(result)

Rust 侧等价的均值窗口表达式是 col("Speed").mean().over(["Type 1"])(见 docs/source/src/rust/user-guide/expressions/window.rs),说明 over 在 Python 与 Rust API 之间的语义是一致的。

更多示例:分组 Top-N 组合练习

最后做几道综合练习,把“排序 + head + over + explode”串起来。下面一次性计算了 5 个窗口结果:

  • 按类型排序所有宝可梦(输出 “Type 1” 组内有序的头部元素);
  • 每种 “Type 1” 类型取前 3 只;
  • 类型内按速度降序取前 3,命名 "fastest/group"
  • 类型内按攻击力降序取前 3,命名 "strongest/group"
  • 类型内按名字排序取前 3,命名 "sorted_by_alphabet"
result = pokemon.sort("Type 1").select(
    pl.col("Type 1").head(3).over("Type 1", mapping_strategy="explode"),
    pl.col("Name")
    .sort_by(pl.col("Speed"), descending=True)
    .head(3)
    .over("Type 1", mapping_strategy="explode")
    .alias("fastest/group"),
    pl.col("Name")
    .sort_by(pl.col("Attack"), descending=True)
    .head(3)
    .over("Type 1", mapping_strategy="explode")
    .alias("strongest/group"),
    pl.col("Name")
    .sort()
    .head(3)
    .over("Type 1", mapping_strategy="explode")
    .alias("sorted_by_alphabet"),
)
print(result)

关键点在于:外层先 sort("Type 1") 保证同一类型相邻,随后每个分组表达式把 head(3) 限定的行以新行的形式产出,配合 mapping_strategy="explode" 输出“每种类型前 3 名”的紧凑结果。这正是“窗口函数负责分组内取值、explode 负责改变行数”的典型配合。sort_by / head / over 三种算子叠加时,Rust 侧等价实现需要对 head 结果再显式 .explode(...)(见 docs/source/src/rust/user-guide/expressions/window.rs),这一点体现了 Python API 对窗口展开路径的自动化封装。

补充:over 的完整参数与排序支持

Python 侧 over 完整签名(py-polars/src/polars/expr/expr.py)为:

Expr.over(
    partition_by=None,          # 分组键,支持列名或表达式,可迭代
    *more_exprs,                # 更多分组键(位置参数)
    order_by=None,              # 在每个分区内先按此排序再计算
    descending=False,           # order_by 的排序方向
    nulls_last=False,           # 排序时 null 是否放最后
    mapping_strategy="group_to_rows",  # "group_to_rows" | "join" | "explode"
)

其中 order_by 参数对顺序敏感的窗口运算(如 cum_sumdiff)尤其有用:它可以先在每个分组内部排序,再施加窗口表达式,省去显式 sort_by 的书写。结合前面所述的行数语义,可以总结使用建议:

  • 想要不改变行数、只附加分组统计量,用默认 "group_to_rows"(标量结果自动广播);
  • 想要组内排序并改变行数、让同组相邻,用 "explode" 获得最佳性能;
  • 想要每一行都携带整组的聚合列表,用 "join",但需留意其内存开销;
  • 组内结果长度与组大小不一致时,"group_to_rows" 不适用,应选用 "explode""join"

这些窗口行为的回归测试可在 crates/polars/tests/it/lazy/expressions/window.rs 中找到(仓库 Rust 集成测试目录),适合作为进一步阅读与验证的实现参考。

延伸阅读

登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.14 K
2.75 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
857
1.35 K
docsdocs
暂无描述
Markdown
898
5.82 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
531
596
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
921
1.84 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.8 K
1.02 K
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.36 K
1.46 K
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
1.02 K
519
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
548
391