编程 Polars 1.41 深度拆解:当 Rust 把 DataFrame 从 pandas 手里抢过来——从 Arrow 内存模型、惰性优化器到流式引擎的全链路实战

2026-08-19 01:14:11 +0800 CST views 4

Polars 1.41 深度拆解:当 Rust 把 DataFrame 从 pandas 手里抢过来——从 Arrow 内存模型、惰性优化器到流式引擎的全链路实战

选题来源:最新开源项目 GitHub Trending / 编程工具链最新发布。Polars 1.41 于 2026 年 8 月 19 日发布,带来更快的 Parquet 元数据解码等性能优化;项目累计下载量已破 6.75 亿次,GitHub 39k+ stars。本文从内存模型、查询优化器、多线程引擎到流式执行,配可运行代码与生产优化清单,彻底讲透这个正在重写 Python 数据栈的库。

一、背景介绍:pandas 的「优雅」,和它扛不住的真实世界

如果你写过任何数据分析、ETL、特征工程或者机器学习前置处理的代码,那么你几乎一定写过 import pandas as pd。pandas 用一套 DataFrame 抽象统一了整个 Python 数据生态,它的 API 设计之直觉、文档之丰富,至今仍是教科书级别的存在。

但如果你处理过「真正的大数据」——一个 50GB 的 CSV、一张上亿行的宽表、一次需要连续 join 五张表的清洗任务——你就会撞上 pandas 的三堵墙:

第一堵墙:单线程。 pandas 的核心循环大多跑在单个 CPU 核心上。你的机器有 16 核,pandas 只用 1 个。数据量上升时,你看到的不是线性变慢,而是指数级的绝望。

第二堵墙:行式内存 + 临时副本。 pandas 底层是 NumPy 的 object 数组,字符串、混合类型都得走 Python 对象,意味着频繁的内存分配、拷贝、以及 Python GIL 的反复进出。一次 df[df.a > 0][["b","c"]],中间可能生成好几个完整的中间 DataFrame。

第三堵墙:没有查询优化器。 pandas 是「命令式」的:你写的每一行代码,它就老老实实执行一遍。你先 filter 再 select,它就真的先全量 filter 再全量 select;你写了个后面根本没用到的列计算,它就真的算了一遍。它不会「读懂」你的意图去重排执行计划。

正是这三堵墙,催生了 Polars。

Polars 由 Python 数据生态中一位核心工程师在 2020 年发起,核心完全用 Rust 重写,内存模型直接建立在 Apache Arrow 之上。它的设计哲学和 pandas 正好相反:

  • 列式(columnar)存储:数据按列连续存放,CPU 缓存命中率高,天然适配向量化(SIMD)。
  • 惰性(lazy)优先:你描述「想要什么」,优化器决定「怎么算最快」。
  • 多线程即默认:开箱即用 Rayon 工作窃取调度,把 16 核全用满,且全程不碰 GIL。
  • 零拷贝集成:Arrow 内存可被其他 Arrow 生态工具(DuckDB、Arrow Flight、Spark、DataFusion)直接读取,无需序列化。

官方在 TPC-H 派生基准上的结论是:相比 pandas 可获得 超过 30 倍 的性能提升,部分场景号称可达 50 倍。而本文要做的,是把这些「倍数」背后的工程原理,一层层拆开给你看。

二、核心概念:Series、DataFrame 与「表达式」这套新语法

在深入引擎之前,必须先建立 Polars 的三个核心心智模型,否则后面的优化器、流式执行你都会觉得「飘」。

2.1 Series 与 DataFrame

Series 是带有名字的一列同质数据(同一种数据类型),DataFrame 是若干 Series 组成的表。和 pandas 最大的不同:Polars 的 Series 直接映射到 Arrow 的 Array,列内存是物理连续的。

import polars as pl

df = pl.DataFrame({
    "name":       ["Alice", "Bob", "Charlie", "Diana"],
    "department": ["Eng", "Sales", "Eng", "Sales"],
    "salary":     [120_000, 95_000, 110_000, 105_000],
    "hire_year":  [2019, 2021, 2018, 2022],
})
print(df)
# shape: (4, 4)
# ┌───────┬────────────┬─────────┬───────────┐
# │ name  ┆ department ┆ salary  ┆ hire_year │
# └───────┴────────────┴─────────┴───────────┘

pl.DataFrame(...) 走的是 eager(立即执行) 路径:你给的数据,立刻被物化(materialize)成内存里的 Arrow 列。

2.2 表达式(Expression):Polars 的灵魂

pandas 的思维是「操作整张表」:df["salary"] * 1.1。Polars 的杀手锏是 表达式——pl.col("salary") * 1.1 不是一个立即求值的结果,而是一段「计算配方」。它描述「对名为 salary 的列做什么」,只有在被 select / with_columns / filter 等「上下文」包裹时才会真正执行。

# 表达式本身不执行,只描述意图
expr = pl.col("salary") * 1.1

# 放入上下文才执行
df_with_bonus = df.with_columns(
    expr.alias("salary_with_bonus"),
    (pl.col("hire_year").sub(2026).abs()).alias("tenure"),
)

为什么要多此一举?因为表达式是可组合的、可被优化器重排的「声明式」描述。下面这个例子最能说明问题:

result = (
    df
    .lazy()                                  # 进入惰性模式
    .filter(pl.col("salary") > 100_000)
    .select("name", "salary")
    .collect()                               # 此时才真正执行
)

注意 .lazy() 之后,.filter.select 都只是往查询计划(LogicalPlan)里「登记」了一个节点,直到 .collect() 被调用,优化器才会出手。这中间的「延迟」,就是 Polars 性能的全部秘密所在。

2.3 eager vs lazy:不是风格选择,是性能分水岭

维度eager(.DataFramelazy(.lazy()
执行时机每步立即物化攒成计划,最后一次性优化执行
优化器不参与全程参与(下推、剪枝、重排)
IO 下推不支持支持(filter/select 推到 Parquet/CSV 读取层)
内存峰值高(每步可能留中间表)低(流式 + 剪枝)
适用场景交互式探索小数据ETL、大文件、生产管道

结论先行:生产环境里,几乎没有理由不用 lazy。

三、架构分析:优化器、多线程引擎与流式执行

这一节是全文的硬核。Polars 的速度不是「Rust 比 Python 快」这么一句废话能解释的,它是一整套系统工程。

3.1 查询优化器(Query Optimizer)

当你在 lazy 模式下链式调用一堆表达式,Polars 在 .collect() 前会跑一遍优化器,应用一系列 规则(rule)。理解这些规则,你才能写出「被优化」的代码。

(1)谓词下推(Predicate Pushdown)

如果你写了 .filter(pl.col("x") > 10).select("y"),优化器会把 filter 提前到 select 之前执行——先在更少的行上做列投影,减少后续计算量。更进一步,当数据源是 Parquet 时,这个 filter 会直接下推到文件读取层,Parquet 的 Row Group 统计信息(min/max)能让引擎跳过根本不可能命中条件的整个 Row Group,磁盘都不用读。

(2)列投影下推(Projection Pushdown)

.select("y") 会让引擎只读取 y 这一列。在 Parquet 这种列式格式下,意味着其他列的物理数据块完全不被加载进内存。一张 200 列的宽表,你只读 3 列,内存和 IO 直接降到 1.5%。

(3)切片下推(Slice Pushdown)

.head(100) 本是「取前 100 行」,但配合排序/过滤,优化器能把它下推到扫描阶段,让引擎只读够 100 行就停。处理几十 GB 日志只想要「最近的 100 条」时,这个优化让耗时从分钟级降到毫秒级。

(4)公共子计划消除(Common Subplan Elimination)

如果你对同一份数据做了两个分支计算,优化器会识别出重复的子计划,只算一次。

(5)表达式简化(Simplify Expression) & 类型强转(Type Coercion)

pl.col("a") + 0 会被化简掉;混合类型运算会被静态插入类型转换节点,避免运行时反复 infer。

你可以用一行代码「偷看」优化器到底干了什么:

q = (
    pl.scan_parquet("data/*.parquet")           # 惰性扫描,不立即读
      .filter(pl.col("ts") > "2026-01-01")
      .filter(pl.col("amount") > 100)
      .select("user", "amount", "ts")
      .group_by("user")
      .agg(pl.col("amount").sum().alias("total"))
)
print(q.explain(optimized=True))                # 打印优化后的物理计划

explain() 输出的计划里,你会看到 SELECTION / PROJECTION / AGGREGATION 节点,以及诸如 PUSHED {user, amount, ts} 这样的下推标注。读懂它,你就从「用 Polars」进阶到了「调 Polars」。

3.2 多线程向量化引擎(基于 Rayon)

Polars 的执行层默认使用 Rust 的 Rayon 数据并行库。Rayon 提供的是「工作窃取(work-stealing)」调度:每个 CPU 核心有一个本地任务队列,核心空闲时会去「偷」别的核心的活干。这避免了传统线程池的锁竞争,扩展到几十核依然线性。

更关键的是 向量化:Polars 对 Series 的操作通常不是「逐元素 Python 循环」,而是用 Rust 写死的、一次处理 1024 行的 chunk(Arrow 的 Array 天然分块),并在底层尽量触发 CPU 的 SIMD 指令(单指令多数据)。一个 salary * 1.1,在 100 万行上,执行的是「一条 CPU 指令批量算 8/16/32 个 float」,而不是 100 万次 Python 函数调用。

这就是为什么「永远优先用原生表达式,不要写 Python 循环」:

# ❌ 慢:map_elements 把每行拉回 Python,GIL + 解释器开销
df.with_columns(
    pl.col("salary").map_elements(lambda s: s * 1.1, return_dtype=pl.Float64)
      .alias("bad_bonus")
)

# ✅ 快:纯 Rust 向量化,零 Python 介入
df.with_columns(
    (pl.col("salary") * 1.1).alias("good_bonus")
)

map_elements(以及 apply)是「逃生舱」,不是「日常座舱」。它把数据从 Arrow 内存「降维」成 Python 对象,性能立刻回到 pandas 级别。能用 pl.col() + 原生算子的,绝不用 map

3.3 流式引擎(Streaming / Out-of-Core)

当数据量超过内存(RAM),eager 模式会直接 OOM 崩溃。Polars 的杀手锏是 流式执行引擎:它把查询切成一系列的「流水线阶段(pipeline stage)」,每个阶段以「批次(batch)」为单位处理数据——读一批、算一批、写出一批,内存里始终只保留几个 batch。

关键接口是 sink_*

(
    pl.scan_csv("huge_50gb.csv")                # 惰性扫描,不进内存
      .filter(pl.col("event") == "purchase")
      .group_by("user_id")
      .agg(pl.col("price").sum().alias("spent"))
      .sort("spent", descending=True)
      .head(1000)
      .sink_parquet("top_spenders.parquet")     # 流式写出,边算边落盘
)

这条管道可以处理远超内存的数据集,而且流式引擎同样享受谓词/列下推与多线程。.sink_parquet() 触发的是纯流式路径,不会在 .collect() 时把全量结果物化到内存。

工程经验:只要数据可能超过可用内存的 60%,就直接写 lazy + sink,别等 OOM 了再改。

3.4 Arrow 内存模型:零拷贝生态的门票

Polars 的所有数据都活在 Arrow 格式里。这意味着:

  • 读取 Parquet(Arrow 的列式兄弟)几乎是「直接映射」,极少转换;
  • 与 DuckDB、Apache DataFusion、Spark、Arrow Flight 之间可以零拷贝传递数据;
  • 多进程/多语言(Python↔Rust↔Node)共享同一块内存布局,没有序列化损耗。

这是 Polars 能融入「现代数据栈」而非成为又一个孤岛的根本原因。

四、代码实战:从入门到生产管道

光能讲原理不够,下面是一套可以直接跑、可以直接搬进项目的实战。

4.1 环境准备

pip install polars
# 如需 Arrow 数据库直连 / Excel / 云存储等扩展
# pip install polars[database,excel,cloud]
import polars as pl
print(pl.__version__)   # 1.41.x

4.2 表达式的「组合魔法」

Polars 表达式可以像乐高一样拼接。下面这个例子展示「条件聚合 + 窗口函数 + 多列计算」一次写完:

df = pl.DataFrame({
    "user":  ["u1", "u1", "u2", "u2", "u3"],
    "ts":    ["2026-01-01", "2026-03-15", "2026-02-10", "2026-05-20", "2026-04-01"],
    "price": [100, 200, 50, 300, 150],
    "city":  ["BJ", "SH", "BJ", "GZ", "SH"],
})

out = df.with_columns(
    # 窗口函数:每个用户的价格排名(无需 group_by 物化)
    pl.col("price").rank().over("user").alias("rank_in_user"),
    # 条件聚合:只统计该用户 SH 城市的累计金额
    pl.col("price").filter(pl.col("city") == "SH")
      .sum().over("user").alias("sh_total"),
    # 字符串 + 日期表达式
    pl.col("ts").str.slice(0, 4).alias("year"),
).filter(pl.col("price") > 80)

print(out)

注意 .over("user") 是窗口上下文,它在不折叠行的前提下做分组计算——这正是 pandas 里要用 transform 才能别扭实现的场景,在 Polars 里是原生一等公民。

4.3 IO 下推:让磁盘替你干活

这是最容易拿到「免费性能」的地方。两个写法结果相同,性能天差地别

# ❌ 先把 20GB 全读进内存,再在内存里 filter
df = pl.read_parquet("events/*.parquet")
small = df.filter(pl.col("day") == "2026-08-19")

# ✅ 惰性扫描,filter 下推到 Parquet 读取层,只加载命中的 Row Group
small = (
    pl.scan_parquet("events/*.parquet")
      .filter(pl.col("day") == "2026-08-19")
      .select("user_id", "action", "latency_ms")
      .collect()
)

后者在 Parquet 的 min/max 统计加持下,可能只读取了物理文件的 5%。只要你的数据来自 Parquet/CSV 且后面有 filter/select,永远用 scan_* + lazy。

4.4 真实 ETL 管道示例

假设我们要从一批原始点击日志里,算出「每个用户过去 7 天的日均停留时长,并标记异常长会话」:

q = (
    pl.scan_parquet("clicks/*.parquet")
      .filter(pl.col("event_time") >= (pl.col("event_time").max() - pl.duration(days=7)))
      .filter(pl.col("duration_ms").is_not_null())
      .with_columns(
          (pl.col("duration_ms") / 1000.0).alias("duration_s"),
          pl.col("duration_ms").gt(pl.col("duration_ms").quantile(0.99).over("user_id"))
            .alias("is_outlier"),
      )
      .group_by("user_id")
      .agg(
          pl.col("duration_s").mean().alias("avg_dur_s"),
          pl.col("is_outlier").sum().alias("outlier_cnt"),
          pl.len().alias("session_cnt"),
      )
      .filter(pl.col("avg_dur_s") > 30)
      .sort("avg_dur_s", descending=True)
)

report = q.collect()   # 小结果进内存
# 或者 q.sink_parquet("user_dur_report.parquet") 处理超大结果

这一整条从「原始日志」到「用户级日报」的管道,Polars 会自动完成下推、剪枝、并行与(可选)流式执行。等效的 pandas 代码往往是它的数倍行数、数十倍耗时。

4.5 自定义算子(UDF):Python 逃生舱与 Rust 插件

当原生表达式真的覆盖不了你的算法(比如一个复杂的字符串状态机),有两层逃生舱。

第一层:Python map_elements(慢,但零门槛):

def normalize_phone(s: str) -> str:
    return "".join(ch for ch in s if ch.isdigit())[-11:]

df = df.with_columns(
    pl.col("phone").map_elements(normalize_phone, return_dtype=pl.String)
      .alias("phone_norm")
)

第二层:Rust 插件(生产级性能)。Polars 支持用 Rust 编写编译进引擎的自定义表达式,通过 #[polars_expr] 宏注册,运行时和原生算子一样快、一样并行:

// my_plugin/src/lib.rs —— 概念示例(需 polars 插件工具链)
use polars::prelude::*;
use pyo3_polars::derive::polars_expr;

#[polars_expr(output_type=Float64)]
fn pl_sigmoid(inputs: &[Series]) -> PolarsResult<Series> {
    let ca = inputs[0].f64()?;
    let out: Float64Chunked = ca.apply_values(|v| 1.0 / (1.0 + (-v).exp()));
    Ok(out.into_series())
}

插件经过编译后,Python 侧以原生速度调用,既不退出 Arrow 内存,也不碰 GIL。对于热路径上的复杂计算,这是把 Polars 性能榨干的终极手段。

4.6 多线程调优:把核用满

Polars 默认吃满所有核心。在容器/受限环境里,可以用环境变量或配置约束:

import polars as pl

# 方式一:环境变量(进程启动前)
# export POLARS_MAX_THREADS=4

# 方式二:代码内设置
pl.Config.set_num_threads(4)

# 查看当前线程配置
print(pl.Config.get_num_threads())

多进程并行(比如对多个文件分片处理)时,建议「每进程 1~2 个线程 × 进程数 ≈ 物理核数」,避免线程数远超核数导致的上下文切换开销。

五、性能优化清单(可直接贴进团队 Wiki)

把上面所有原理浓缩成一份「生产 checklist」:

  1. 永远 lazy + scan_*:读 Parquet/CSV/JSON 一律用 pl.scan_* 进入惰性模式,让下推和流式生效。
  2. 避免 map_elements / apply:90% 的场景都能用 pl.col() 原生表达式表达;实在不行再考虑 Rust 插件。
  3. 列投影即省钱select 只挑要用的列,Parquet 下推会直接少读磁盘。
  4. filter 尽量早:越早 filter,下游数据越少,下推收益越大。
  5. 大文件用 sink_* 流式落盘:超过内存 60% 的数据,直接 lazy + .sink_parquet(),别 .collect()
  6. 优先 Parquet,而非 CSV:Parquet 列式 + 统计信息 + 压缩,读取快、体积小、还能下推;CSV 是性能黑洞。
  7. explain(optimized=True) 自查:上线前打印执行计划,确认下推是否真的生效。
  8. 控制线程数:容器里设 POLARS_MAX_THREADS,避免和同机其他服务抢核。
  9. 窗口计算用 .over():替代 pandas 的 groupby().transform(),原生且并行。
  10. 读写直连云存储:Polars 原生支持 S3 / Azure Blob,无需先 aws s3 cp 到本地。

5.1 一个可复现的微基准

import time
import polars as pl

# 构造 1000 万行数据
n = 10_000_000
df = pl.DataFrame({
    "a": pl.int_range(0, n, dtype=pl.Int64),
    "b": pl.int_range(0, n, dtype=pl.Int64),
}).with_columns((pl.col("a") * 2 + pl.col("b")).alias("c"))

# Polars 原生向量化
t0 = time.perf_counter()
r1 = df.filter(pl.col("c") > n).select(pl.col("c").mean())
print("polars:", time.perf_counter() - t0, r1)

把同样的 filter + mean 用 pandas 逐行实现,你会发现 Polars 在千万行量级轻松快出一个数量级以上——而这还只是单表单算子的「开胃菜」,多表 join + 聚合的复杂管道里差距更夸张。

六、Polars Cloud 与生态:从单机到分布式

单机的 Polars 已经很强,但数据团队的终极诉求是「同一套 API,从笔记本到生产集群无缝伸缩」。这正是 Polars Cloud 的定位:它把 Polars 的查询计划编译成分布式执行,在云端/本地可水平扩展,且对用户的代码零改动——你写的那段 lazy 查询,既可以 .collect() 在本机跑,也可以提交给 Polars Cloud 跑在成百上千核上。

它的意义在于打破了「pandas 写原型、Spark 写生产」的长期割裂:现在同一套表达式,原型期在笔记本跑,上线后交给云直接放大,没有重写、没有 API 断层。配合 Arrow 生态,Polars 还能与 DuckDB(OLAP)、DataFusion(查询引擎)、Arrow Flight(数据传输)组成一套「全 Arrow、零序列化」的现代数据管道。

七、总结与展望

回看开头 pandas 的三堵墙——单线程、行式拷贝、无优化器——Polars 用一个 Rust 内核 + Arrow 内存 + 惰性优化器 + 流式引擎的组合拳,几乎逐一拆掉:

  • 多线程向量化引擎解决了「算得慢」;
  • Arrow 列式内存解决了「存得散、传不动」;
  • 惰性查询优化器解决了「不会聪明地算」;
  • 流式执行解决了「内存装不下」;
  • Polars Cloud解决了「单机到分布式的断层」。

对工程师的实操建议很明确:新项目直接上 Polars;老 pandas 项目从数据读取和重计算的热点函数开始迁移,先用 pl.from_pandas() 桥接,再逐步把 apply 重写成表达式。你会立刻感受到迭代速度的质变。

展望 2026 下半年:随着 1.41 在 Parquet 元数据解码、流式引擎稳定性上的持续打磨,以及 Polars Cloud 的成熟,Polars 正从「pandas 的高性能替代品」进化为「Python 数据栈的新默认底座」。对于每天和表格数据打交道的人而言,现在花一个下午吃透它的表达式与优化器,可能是今年性价比最高的一次技术投资。


关键 takeaway:Polars 快,不是因为「Rust 比 Python 快」这种车轱辘话,而是因为它把你那句 df[df.x>0][['y']] 翻译成了一个经过下推、剪枝、并行、向量化重排的执行计划。理解优化器,比背诵 API 重要一百倍。

推荐文章

随机分数html
2025-01-25 10:56:34 +0800 CST
程序员茄子在线接单