《Python编程实战》11.2 Polars 与大规模数据

用同一份 50 万行真实数据实测 Polars 2.0:LazyFrame 惰性求值与查询计划、投影与谓词下推、多线程并行与流式引擎、表达式 API 与 pandas 的对照迁移,并给出同一任务的耗时对比,最后讲清什么时候该换、什么时候别换。

本节目标:用与上一节完全相同的 50 万行数据,实测 Polars 2.0 的列式内存模型、LazyFrame 惰性求值与查询计划(.explain())、投影/谓词下推、流式引擎与分区落盘,并给出与 pandas 的同任务耗时对比,最后说清「什么时候该从 pandas 换到 Polars」。
适用版本:Python 3.12+(实测 3.14.6);polars 2.0.0、numpy 2.5.3

11.2 Polars 与大规模数据

上一节我们看到 pandas 单线程、先物化再计算的路径在数据量上来后会撞墙。Polars 是另一条路线:Rust 实现、列式存储(Apache Arrow 内存模型)、多线程并行、查询先规划再执行。它和 pandas 不是「新旧替代」,而是两种取舍——这一节用实测把取舍摆清楚。

数据仍是 11.1 的那份 50 万行订单表,先写成 CSV 供两侧公平读取:

import polars as pl
df = pl.read_csv("orders.csv")
print(df.height)   # 500000

11.2.1 列式内存模型

pandas 的行式块存的是「一整块连续内存」,Polars 按列存:每一列独立成一块 Arrow 数组,同类型、连续、可直接交给 SIMD 和多线程。这带来两个直接好处:只读需要的列时,其他列的内存根本不用碰;同一列的运算可整列并行。

Polars 自带内存估算,不用 deep=True 那种递归扫描:

print(df.estimated_size("mb"))   # -> 15.26 MB

50 万行 5 列 15.26 MB,和 pandas 的 16.69 MB 量级一致(都是紧凑数值存储)。列式的收益不体现在「同样数据更省内存」,而体现在「只碰需要的那几列」。

11.2.2 LazyFrame:先规划,再执行

这是 Polars 和 pandas 最大的心智差异。pl.scan_csv(...) 返回一个 LazyFrame——它不读数据,只记下「你要做什么」,直到你调 .collect() 才真正执行。中间这一段时间,查询优化器会重写你的查询计划。

q = (pl.scan_csv("orders.csv")
     .filter(pl.col("amount") > 500)
     .group_by("city")
     .agg(pl.col("amount").sum().alias("total"), pl.len().alias("n"))
     .sort("total", descending=True))
print(q.explain())

真跑的查询计划(实测输出):

SORT BY [descending: [true]] [col("total")]
  AGGREGATE[maintain_order: false]
    [col("amount").sum().alias("total"), len().alias("n")] BY [col("city")]
    FROM
    Csv SCAN [orders.csv]
    PROJECT 2/5 COLUMNS
    SELECTION: col("amount") > 500.0
    ESTIMATED ROWS: 616263

两个优化看得见:

  • 投影下推(Projection Pushdown):PROJECT 2/5 COLUMNS——查询只用 city 和 amount 两列,CSV 解析时只解析这两列,另外三列连读都不读。
  • 谓词下推(Predicate Pushdown):SELECTION: col("amount") > 500.0 被推到了扫描层,在读入阶段就过滤,进聚合的数据量直接变小。

pandas 没有这一层——你写 df = pd.read_csv(...) 的那一刻,全表已经进了内存,后面的 filter 是事后过滤。数据量越大,这个差距越致命。养成习惯:Polars 里能用 scan_* 就别用 read_*。

11.2.3 表达式 API:一列即一个表达式

Polars 的列操作叫表达式(pl.col("x")),表达式可组合、可并行:

df.with_columns(
    (pl.col("amount") * pl.col("qty")).alias("revenue"),
    pl.col("amount").rank(descending=True).alias("amount_rank"),
)

和 pandas 的关键差异:Polars 没有索引。一切靠列名和表达式,select/filter/with_columns 都返回新表,天然不可变——不存在 11.1 里那些链式赋值的坑。pandas 常见写法到 Polars 的对照:

任务pandasPolars
选列df[["a", "b"]]df.select("a", "b")
过滤df[df["a"] > 1]df.filter(pl.col("a") > 1)
新增列df["c"] = df["a"] * 2df.with_columns((pl.col("a") * 2).alias("c"))
分组聚合df.groupby("k")["v"].sum()df.group_by("k").agg(pl.col("v").sum())
连接df.merge(other, on="k")df.join(other, on="k", how="left")
排序df.sort_values("v")df.sort("v")

分组聚合的写法值得注意:Polars 里 group_by 必须跟 agg,且一次 agg 里可以放多个不同表达式(求和、计数、去重、分位一次算完),不像 pandas 要 agg(["sum", "mean"]) 那样用字符串列表。

11.2.4 并行执行与流式引擎

.collect() 默认用内存引擎(把结果物化进内存),但多线程并行是自动的——同一个 50 万行 group_by 实测:

polars eager group_by: 1.8 ms
polars lazy  group_by: 1.7 ms

对比 11.1 里 pandas 的同一操作 7.1 ms,约 4 倍。差别主要来自:pandas 的 groupby 是单线程哈希聚合,Polars 按分区并行 + 列式。

当数据大过内存时,用流式引擎边读边算,不把全量放进内存:

q = (pl.scan_csv("orders.csv")
     .filter(pl.col("amount") > 500)
     .with_columns((pl.col("amount") * pl.col("qty")).alias("revenue")))
q.sink_parquet("out_stream.parquet")   # 流式落盘,峰值内存受控

实测耗时 27.5 ms。也可以显式指定引擎:

q.collect(engine="streaming")   # 流式;大表用
q.collect(engine="in-memory")   # 内存;小表快

sink_* 是流式专用的落盘方法(sink_parquet/sink_csv),和 collect 的区别是它分块写、不物化全表——这是处理「读不下的数据集」的正解。

11.2.5 同一任务的实测对比

把「读 CSV → 过滤 amount>500 → 按 city 求和 → 降序排」这条完整链路在两侧各跑一遍(各 3 次取最快):

# pandas 全量读
d = pd.read_csv("orders.csv")
d[d["amount"] > 500].groupby("city", observed=True)["amount"].sum().sort_values(ascending=False)

# polars 惰性
(pl.scan_csv("orders.csv").filter(pl.col("amount") > 500)
   .group_by("city").agg(pl.col("amount").sum()).sort("amount", descending=True).collect())
pandas 全量读 :  128.8 ms
pandas 分块读 :  131.4 ms
polars 惰性   :   12.5 ms

约 10 倍差距,而且两侧结果一致(广州 求和最高,约 2.3717e7)。差距的来源不是「Rust 比 C 快」,而是三件事叠加:只解析两列(投影下推)、读入即过滤(谓词下推)、多线程聚合。pandas 要等价地省,得手动 usecols=["city", "amount"] + 手动分块过滤,写起来啰嗦且仍拿不到多线程。

11.2.6 Parquet 与分区落盘

Polars 原生支持 Parquet(Rust 实现,不依赖 pyarrow)。50 万行实测:

CSV 大小: 13841 KB   Parquet 大小: 3904 KB
写 parquet: 24.0 ms  读 parquet: 31.0 ms

Parquet 是列式 + 压缩 + 带 schema,体积约 CSV 的 1/3.5,且读的时候能只读需要的列。它还能按分区写盘:

df.write_parquet("partitions/", partition_by="city")

实测生成 8 个 hive 风格分区目录(city=上海/00000000.parquet、city=北京/... 等,中文目录名会被 URL 编码)。之后用通配符把整个分区目录当一个逻辑表读:

pl.scan_parquet("partitions/**/*.parquet").group_by("city").agg(pl.col("revenue").sum()).collect()

分区落盘的价值:下游查询若带 city 过滤条件,可以只扫命中的分区目录,跳过其余数据——这是数据仓库里「分区裁剪」的基本功,也是下一节 ETL 落盘的核心手段。

⚠️ 一处未实测说明:pandas 的 df.to_parquet() 需要 pyarrow 或 fastparquet 引擎,本机未安装(实测报 ImportError: Unable to find a usable engine; tried using: 'pyarrow', 'fastparquet')。所以本文所有 Parquet 操作都是 Polars 原生实现的实测结果;pandas 侧的 Parquet 读写未跑通,不贴任何「预期输出」。

11.2.7 窗口函数:over() 是第二把利器

group_by 会把数据「压扁」成每组一行;但很多计算要保留原行、同时带上组内信息(组内占比、组内排名、组内累计)。这在 pandas 里靠 groupby().transform() 或 groupby().apply(),写起来容易绕;Polars 用 .over("key") 一行表达:

import polars as pl
df = pl.DataFrame({
    "city": ["北京", "北京", "上海", "上海", "上海"],
    "amount": [100.0, 300.0, 200.0, 50.0, 250.0],
})
out = df.with_columns(
    (pl.col("amount") / pl.col("amount").sum().over("city")).alias("share"),
    pl.col("amount").cum_sum().over("city").alias("cum"),
)
print(out)

实测输出:

shape: (5, 4)
┌──────┬────────┬───────┬───────┐
│ city ┆ amount ┆ share ┆ cum   │
╞══════╪════════╪═══════╪═══════╡
│ 北京 ┆ 100.0  ┆ 0.25  ┆ 100.0 │
│ 北京 ┆ 300.0  ┆ 0.75  ┆ 400.0 │
│ 上海 ┆ 200.0  ┆ 0.4   ┆ 200.0 │
│ 上海 ┆ 50.0   ┆ 0.1   ┆ 250.0 │
│ 上海 ┆ 250.0  ┆ 0.5   ┆ 500.0 │
└──────┴────────┴───────┴───────┘

行数不变(5 行),每行多带了「占本市总额的比例」和「组内累计」。over() 的关键价值是:它和 with_columns 组合,把「分组」变成列上的一个属性,而不是一次性的聚合动作——排名、占比、移动平均、组内差分全靠它,且表达式内部自动并行。

11.2.8 什么时候该换,什么时候别换

Polars 不是无脑更优,选型看场景:

场景建议
数据几十万行以上、要反复算Polars,惰性 + 并行收益明显
只读 CSV/Parquet 的一部分列Polars,投影下推省 I/O
数据大到内存放不下Polars 流式引擎 sink_*
深度依赖 pandas 生态(statsmodels、sklearn 接口、大量 .plot())留在 pandas
一次性小脚本、几十行数据无所谓,pandas 写法更顺手
需要 pandas 特有索引语义(多重索引、.loc 切片)留在 pandas

两者可以混用:Polars 里 df.to_pandas()、pandas 里 pl.from_pandas(df)(后者零拷贝,共享 Arrow 内存)。别为了「先进」把整个项目从 pandas 重写——按热点换,收益最大的是 I/O 密集和大聚合。

延伸阅读

小结

  • Polars 的收益来自列式存储 + 多线程 + 查询规划三件事叠加,不是单纯的「语言更快」。
  • scan_* + .collect() 是默认姿势:.explain() 能看见投影下推(PROJECT 2/5 COLUMNS)和谓词下推(SELECTION),pandas 没有这一层。
  • Polars 无索引、不可变,表达式 API 天然避开链式赋值坑;group_by 必须跟 agg,一次聚合多个表达式。
  • 同一「读 CSV → 过滤 → 聚合 → 排序」链路实测 12.5 ms vs pandas 128.8 ms(约 10 倍);group_by 单项 1.8 ms vs 7.1 ms(约 4 倍)。
  • Parquet 用 Polars 原生实现:体积约 CSV 的 1/3.5,支持 partition_by 分区落盘与分区裁剪;pandas 的 Parquet 因本机无 pyarrow 未实测。
  • 数据大到内存放不下用 sink_parquet 流式落盘(collect(engine="streaming"));小脚本和重生态场景不必换,两者可 to_pandas()/from_pandas() 混用。

到这里,「用哪个引擎、怎么快」讲完了,但真实的 ETL 不是一个脚本——它是一串有依赖关系、要能重跑、要能落盘分区的任务。下一节我们把 Polars 的落盘能力和一个标准库写的 DAG 调度器拼起来,做出最小可用的数据管道。

阅读导航:上一节:pandas 数据处理与性能陷阱 · 下一节:数据管道与 ETL 编排 。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「python」更多文章

  1. 《Python高级编程》目录
  2. 《Python高级编程》11.3 PEP 流程与版本迁移策略
  3. 《Python高级编程》11.2 嵌入式与自由线程运行时