从 Apache Spark 迁移到 Polars:理解列式 API 与表达式组合的设计差异 从 Apache Spark 迁移到 Polars理解列式 API 与表达式组合的设计差异【免费下载链接】polarsExtremely fast Query Engine for DataFrames, written in Rust项目地址: https://gitcode.com/GitHub_Trending/po/polarsApache Spark 与 Polars 都面向大规模 DataFrame 分析但二者数据模型与编程范式有本质差别Spark 的DataFrame在概念上更接近行的集合而 Polars 的DataFrame更接近列的集合。本文以仓库文档 docs/source/user-guide/migration/spark.md 为骨架通过多个可运行对照示例剖析按列组织带来的表达能力差异说明为何 Polars 可以把在 Spark 中需要多次计算、窗口或人工 key 才能完成的逻辑压缩成一条表达式并从源码层面解释over、rolling_mean、head、shift等关键 API 的语义。读完本文你将能够把常见的 Spark 窗口/聚合代码改写成符合 Polars 惯用法的表达式代码并理解二者在行数对齐、标量广播与表达式可组合性上的边界。列式 DataFrame 与行式 DataFrame两条不同的组织逻辑文档开篇即点出两框架最根本的分歧如果说 Spark 的DataFrame类似行的集合那么 Polars 的DataFrame则更接近列的集合。在 Spark 中一行记录代表现实世界中的一个实体spark.createDataFrame传入的是元组列表而在 Polars 中构造器按列给出数据每列一个列表。请看两种写法构建的同一份样例数据import polars as pl df pl.DataFrame({ foo: [a, b, c, d, d], bar: [1, 2, 3, 4, 5], }) dfs spark.createDataFrame( [ (a, 1), (b, 2), (c, 3), (d, 4), (d, 5), ], schema[foo, bar], )这个差异不只是语法层面。行式模型要求 Spark始终保留同一行内各列之间的对应关系而列式模型下Polars 允许把任意两列交给两条互不相干的表达式去处理再按位置把结果拼成一个新表。理解这一点是读懂接下来两个例子的前提。例一把head和sum塞进同一个selectPolars 允许你在一次select中同时对不同列运行完全独立的表达式df.select( pl.col(foo).sort().head(2), pl.col(bar).filter(pl.col(foo) d).sum() )输出shape: (2, 2) ┌─────┬─────┐ │ foo ┆ bar │ │ --- ┆ --- │ │ str ┆ i64 │ ╞═════╪═════╡ │ a ┆ 9 │ ├╌╌╌╌╌┼╌╌╌╌╌┤ │ b ┆ 9 │ └─────┴─────┘这里foo列上的表达式排序后取前两行得到a、b与bar列上的表达式先按foo d过滤、再求和得到9完全无关。因为对bar的表达式只产出一个标量Polars 就把这个标量对foo表达式输出的每一行做广播——文档特意提醒a、b与那个9之间不存在任何数据上的关联它们的并置纯粹是两列拼接的结果。在 Spark 中要做到同样效果必须先单独把sum算出来再以lit字面量的形式注入from pyspark.sql.functions import col, sum, lit bar_sum ( dfs .where(col(foo) d) .groupBy() .agg(sum(col(bar))) .take(1)[0][0] ) ( dfs .orderBy(foo) .limit(2) .withColumn(bar, lit(bar_sum)) .show() )输出------ |foo|bar| ------ | a| 9| | b| 9| ------从 Polars 源码看sort与head都返回Exprselect接收多个Expr后对它们做独立的列式求值——head(n2)截取的是该列前n个元素而标量型聚合结果会被自动扩展至输出列的高度。这正是列式集合的表达红利Spark 必须物理地把子查询结果塞回行里Polars 则天然支持不同来源、不同形状标量或向量的列并行出现在同一个输出上下文中。例二把两个不同的head拼在一起既然列彼此独立Polars 自然允许在同一张 DataFrame 上叠加两个各取前两条的表达式只要它们输出相同数量的值df.select( pl.col(foo).sort().head(2), pl.col(bar).sort(descendingTrue).head(2), )输出shape: (2, 2) ┌─────┬─────┐ │ foo ┆ bar │ │ --- ┆ --- │ │ str ┆ i64 │ ╞═════╪═════╡ │ a ┆ 5 │ ├╌╌╌╌╌┼╌╌╌╌╌┤ │ b ┆ 4 │ └─────┴─────┘再次强调foo升序前两条a、b与bar降序前两条5、4是两条独立表达式的结果a↔5、b↔4的配对纯属输出列并列摆放的副产品不代表任何行级关联。Spark 行模型不允许这种无依据的配对你必须制造一个人工 key把两条各行取其序号的流 join 起来from pyspark.sql import Window from pyspark.sql.functions import col, row_number foo_dfs ( dfs .withColumn( rownum, row_number().over(Window.orderBy(foo)) ) ) bar_dfs ( dfs .withColumn( rownum, row_number().over(Window.orderBy(col(bar).desc())) ) ) ( foo_dfs.alias(foo) .join(bar_dfs.alias(bar), onrownum) .select(foo.foo, bar.bar) .limit(2) .show() )输出------ |foo|bar| ------ | a| 5| | b| 4| ------这个例子给迁移者的启示很具体当两条表达式产出的行数一致时Polars 直接用位置对齐拼接列即可而在行式模型里对齐必须显式借助row_number加 self-join。反过来也要谨记一条规则并列的列要么行数相同位置对齐要么是标量广播到每一行否则 Polars 会因形状不匹配报错——这正是列式语义对可组合所要求的纪律。表达式可组合性shiftrolling_meanover一次成型Polars 最具迁移价值的特性是把 Spark 里窗口函数必须各自独立成窗口的约束彻底解除。文档给出的典型场景是求滞后变量的滚动均值。比如要构造特征shift(price, 7)之后 7 期窗口的滚动均值按store分区、按date排序。Polars 可以把三步运算写进一条表达式再套一个over上下文df.with_columns( featurepl.col(price).shift(7).rolling_mean(7).over(store, order_bydate) )从实现看shift产生滞后列rolling_mean对其做滚动求均值window_size7表示窗口包含当前行及此前 6 个元素而over负责分组与排序语义第一个位置参数partition_by指定分组列此处为storeorder_by指定组内排序键此处为datemapping_strategy默认为group_to_rows即组内每个聚合值按其行位置映射回原表的每一行组内结果为标量时则广播到整组所有行。PySpark 做不到这一点它只允许逐元素函数 聚合函数的简单组合例如F.mean(F.abs(price)).over(window)因为abs是逐元素操作但不允许窗口函数作为另一个窗口聚合的输入因此F.mean(F.lag(price, 7)).over(window)会直接报错。要复现相同结果lag与mean必须各自拥有一个窗口还要用count保护只有滞后数据足够的行才参与求均值避免把含空值的短窗口算进去from pyspark.sql import Window from pyspark.sql import functions as F window Window().partitionBy(store).orderBy(date) rolling_window window.rowsBetween(-6, 0) ( df .withColumn(lagged_price, F.lag(price, 7).over(window)) .withColumn( feature, F.when( F.count(lagged_price).over(rolling_window) 7, F.mean(lagged_price).over(rolling_window), ), ) )对比之下能清晰看到两层差距窗口数量Spark 需要lag窗口 行数计数窗口 均值窗口共三套窗口逻辑才能自洽Polars 把滞后滚动分组排序抽象进单个Expr由引擎在内部完成窗口化求值。边界处理Polars 的rolling_mean自带有效样本数概念历史版本称min_periods自 1.21.0 起更名为min_samples可声明窗口内至少需要多少非空样本才输出结果省去手写count window_size的守卫逻辑。之所以能这样组合根本原因是 Polars 的表达式是纯列式运算的树shift、rolling_mean都只描述这一列怎么变换变换后的列再整体交给over决定分组与排序任一环节都不需要触碰行的身份。这正是文档强调的Expression Composability——表达式的可组合性来自它对行关联的彻底放弃。迁移实践建议把分步操作改写成表达式上下文综合文档的两个例题与表达式的组合原理从 Spark 迁移到 Polars 时可以遵循下面几条可操作的经验1. 主动寻找列间无关联的操作合并进同一上下文。select、with_columns等上下文可以同时接收多条Expr引擎有机会对其做并行求值与整体优化。Spark 里那种先算 A 存临时列、再拿临时列算 B的链式写法在 Polars 里往往能折叠为一条表达式缩短数据流转路径。2. 分清位置对齐与广播两种列拼接。多条向量表达式若输出行数一致Polars 按位置拼接若某表达式产出单个标量如sum则广播到其它列的每一行。千万不要想当然地把 Spark 行模型中的行级关联带到 Polars 的列并置里——那种关联需要显式的 join、over窗口或结构化数据才能表达。3. 用over取代手工窗口 join。凡是在 Spark 里要用Window.partitionBy().orderBy()加 row_number、lag 再 self-join 才能实现的逻辑优先考虑 Polars 的over(partition_by..., order_by..., mapping_strategy...)。它对窗口表达式的处理在语义上类似执行一次 group by 聚合后再 join 回原表见over的文档说明却避免了用户手写 join。4. 让平移、滚动、排序成为表达式而不是步骤。shift、rolling_mean、sort、head、filter都返回Expr天然可以继续参与下一次组合。迁移时若能坚持用返回表达式的函数组织复用逻辑而非把 DataFrame 传来传去写出的代码既更接近 Polars 的惯用法也更容易被查询优化器处理。更系统的概念差异如惰性求值、Arrow 内存格式、无索引设计可参考同目录下的 pandas 迁移指南想进一步掌握表达式上下文与窗口函数的用法可阅读仓库用户指南中的 表达式与上下文、窗口函数 与 惰性 API 使用 等章节。小结本文围绕 Apache Spark 到 Polars 迁移中最核心的思维转变展开Polars 把 DataFrame 组织为列的集合因而允许互不相关的表达式在同一上下文中组合、并按位置或广播把结果并置成新表——这是行式 Spark 做不到的。两个select例题说明独立表达式并行求值 结果拼接的语义及其形状约束表达式可组合性一节则以shift rolling_mean over与 PySpark 三窗口守卫逻辑的对比展示了 Polars 把窗口计算封装进单个表达式的设计优势。迁移者只要抓住按列思考、用表达式组合、以形状规则自律这三条就能写出既正确又地道的 Polars 代码。原始对照示例可在 docs/source/user-guide/migration/spark.md 中复现验证。【免费下载链接】polarsExtremely fast Query Engine for DataFrames, written in Rust项目地址: https://gitcode.com/GitHub_Trending/po/polars创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考