好的库会为用户封装掉许多复杂性。Polars 也遵循这一理念:你无需了解内部实现,写出的查询默认就应该具有良好的性能。不过,许多用户仍然希望了解其内部机制,有的是为了学习,有的是为了从查询中榨出最后一点性能。本文将从宏观上介绍 Polars 的工作方式,后续文章再深入探讨各个组件。
宏观概览
Polars 是什么?简而言之,它是“具有 DataFrame 前端的查询引擎”。即使对于概览,这个描述也过于笼统。因此,我们通过观察查询如何执行,进一步了解 DataFrame 和查询引擎这两个部分。沿着查询执行过程逐步前进,可以看到每个组件如何工作,以及它承担的角色和目的。
从宏观上看,查询执行过程如下:首先解析并验证查询,将其转换为逻辑计划。该计划描述用户打算做什么,但不描述具体如何做。接着,查询优化器多次遍历计划,消除不必要的工作,生成优化后的逻辑计划。优化阶段结束后,查询规划器把逻辑计划转换为物理计划,规定查询的实际执行方式。最终物理计划作为真正执行查询的输入,调用计算内核。

查询
与 Polars 交互时,你使用的是 DataFrame API。这个 API 从设计之初就考虑了并行执行与性能。因此,编写 Polars 查询,相当于用 Polars 设计的领域专用语言(DSL)编写一个小程序,或者说查询。这种 DSL 有自己的规则,用来规定哪些查询有效、哪些无效。
本文使用著名的纽约出租车行程数据集1。下面的示例按区域计算总费用超过25美元的行程每分钟的平均费用。这个案例足够简单,便于理解,同时也有足够的内容来展示查询引擎的作用。
import polars as pl
query = (
pl.scan_parquet("yellow_tripdata_2023-01.parquet")
.join(pl.scan_csv("taxi_zones.csv"), left_on="PULocationID", right_on="LocationID")
.filter(pl.col("total_amount") > 25)
.group_by("Zone")
.agg(
(pl.col("total_amount") /
(pl.col("tpep_dropoff_datetime") - pl.col("tpep_pickup_datetime")).dt.total_minutes()
).mean().alias("cost_per_minute")
).sort("cost_per_minute",descending=True)
)
上面的查询类型为 LazyFrame。纽约出租车行程数据集超过300万行,这段语句却立即返回,这是为什么?因为语句只是定义了查询,还没有执行它。这种机制称为惰性求值,是 Polars 的重要优势之一。查看 Rust 端的数据结构,可以看到它包含两个部分:logical_plan,以及优化器配置标志 opt_state。
pub struct LazyFrame {
pub logical_plan: LogicalPlan,
pub(crate) opt_state: OptState,
}
逻辑计划是一棵树:数据源是叶节点,变换是其他节点。计划描述查询的结构,以及查询包含的表达式。
pub enum LogicalPlan {
/// Filter on a boolean mask
Selection {
input: Box<LogicalPlan>,
predicate: Expr,
},
/// Column selection
Projection {
expr: Vec<Expr>,
input: Box<LogicalPlan>,
schema: SchemaRef,
options: ProjectionOptions,
},
/// Join operation
Join {
input_left: Box<LogicalPlan>,
input_right: Box<LogicalPlan>,
schema: SchemaRef,
left_on: Vec<Expr>,
right_on: Vec<Expr>,
options: Arc<JoinOptions>,
},
...
}
将查询转换为逻辑计划时,一个重要步骤是验证。Polars 预先知道数据的模式,因此可以验证变换是否正确,以免查询执行到一半才遇到错误。例如,定义一个选择不存在列的查询,会在执行前返回错误:
pl.LazyFrame([]).select(pl.col("does_not_exist"))
polars.exceptions.ColumnNotFoundError: column_does_not_exist
Error originated just after this operation:
DF []; PROJECT */0 COLUMNS; SELECTION: "None"
在 LazyFrame 上调用 show_graph,可以查看逻辑计划:
query.show_graph(optimized=False)
查询优化
查询优化器的目标是优化 LogicalPlan 的性能。它通过遍历树结构,修改、添加或删除节点来实现优化。许多优化都可以加快执行,例如调整操作顺序。一般来说,应尽早执行 filter,以便丢弃不会使用的数据,避免不必要的工作。对于本例,可以用同一个 show_graph 函数展示优化后的逻辑计划:
query.show_graph()
乍看之下,优化前后的计划似乎没有区别。但实际上已经发生了两项重要优化:投影下推(Projection pushdown)与谓词下推(Predicate pushdown)。
Polars 分析查询后发现,只需要少量列:行程数据需要4列,区域数据需要2列。读取整个数据集会浪费资源,因为其余列完全用不到。因此,投影下推通过分析查询,可以显著加快数据读取。在叶节点的 π 4/19 和 π 2/4 标记中,可以看到这一优化。
谓词下推让 Polars 尽可能在靠近数据源的位置过滤数据,避免读取随后会被查询丢弃的数据。过滤节点已经移动到带 σ 标记的 Parquet 读取器中,表示读取器会立即删除不符合过滤条件的行。输入数据量减少后,接下来的连接操作也会快得多。
Polars 支持一系列优化,详见优化文档。
查询执行
逻辑计划优化完毕后,就该执行了。逻辑计划描述的是用户想执行什么,而不是如何执行,这正是物理计划发挥作用的地方。一个简单的实现方式是只提供一种连接算法和一种排序算法,这样便可以直接执行逻辑计划。但这样会付出很大的性能代价,因为了解数据特征和运行环境,能让 Polars 选择更有针对性的算法。因此,连接算法不止一种,而是有多种,各有不同的特点与性能。查询规划器把 LogicalPlan 转换为 PhysicalPlan,为查询选择最合适的算法,然后由计算引擎执行操作。本文不详细介绍引擎的执行模型或其高性能的原因,这些将留待后续讨论。
比较优化前后计划的性能,可以看到本例约有4倍提升。这就是惰性执行以及使用查询引擎的价值:它不必按顺序立即计算每个表达式,而能先优化,避免不必要的工作。用户只需写出查询,无需为这些优化付出额外工作;复杂性全部封装在查询引擎内部。
%%time
query.collect(no_optimization=True);
CPU times: user 2.45 s, sys: 1.18 s, total: 3.62 s
Wall time: 544 ms
%%time
query.collect();
CPU times: user 616 ms, sys: 54.2 ms, total: 670 ms
Wall time: 135 ms
结语
本文介绍了 Polars 的主要组件。希望你现在能更好地理解,从 API 到实际执行,Polars 是如何工作的。接下来的文章将深入探讨每个组件,敬请关注。











暂无评论内容