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