ESC
开源 5 分钟阅读

Polars 2.0 正式发布

Polars 2.0 正式发布:LazyFrame 调用 collect 时默认启用流式引擎,多数查询内存与性能大幅提升;超内存处理(溢写磁盘)默认开启;新增原生 Arrow Map 类型及专用表达式。官方基准显示 Polars 在 TPC-H/TPC-DS 上整体快于 DuckDB 与 DataFusion,192 线程下小数据查询的固定开销待修复。

来源:Hacker News

To see how we perform on typical SQL benchmarks, we ran Polars SQL on data derived from TPC-H and TPC-DS1 and ran it against the latest DuckDB release (1.5.6), DuckDB 2.0 alpha (2.0.0.dev2610011535) and the latest DataFusion release (54.0.0) on a c7a.4xlarge (16 vCPUs, 32GB RAM) and a c7a.metal (192 vCPUs, 384GB RAM). Every query ran 5 times in a hot setting, with a separate process per query and a 60 second timeout. The file cache was cleared between each engine/benchmark (not between queries). For every query we take the best of the 5 runs, and we compare engines on both the sum and the geometric mean of those query times.

The data is generated with tpcgen-cli parquet compiled from source on commit 99bedae. We looked at the default row-group sizes of tpcgen-cli and confirmed they are roughly similar to what Polars scan_csv piped through sink_parquet and Duckdb COPY produce. The SQL queries were generated with DuckDB 1.5.6’s tpch_queries() and tpcds_queries(). The data was stored on EBS.

The charts below show the runtime of each engine in seconds (lower is better), split by machine.

Polars and both DuckDB versions completed all queries. DataFusion timed out on TPC-DS q72 (and once on q67) and ran out of memory on TPC-H q18 on c7a.4xlarge; those queries are excluded from the results above for all engines.

We observe that default Polars is fastest on all but one benchmarks. Polars has a constant overhead when we scale to 192 threads, which hurts small data queries. In fact we see that Polars limited to 32 cores is competitive or winning in all benchmarks. We have diagnosed the cause on our end and will hopefully fix this problem in the next release. More information on the benchmarks can be found in the appendix. We encourage you to replicate our results and have shared a repository for this benchmark here: https://github.com/pola-rs/polars-2.0-benchmark.

This is the one of the biggest impact changes of 2.0. Calling collect on a LazyFrame will now default to the streaming engine, leading to massive memory and performance improvements on most queries. The reason this required a major version bump is that the streaming engine doesn’t guarantee row-order by default for certain operations (join, group_by, unpivot, etc.). If you require observable row-order in those operations, you can opt in to that by setting maintain_order=True.

Out-of-core (spill to disk) is now enabled by default. It starts spilling at ~80% of RAM (this may need tuning). Operations that support out-of-core at this moment (sort, window functions, many expressions) can now start spilling to disk to finish a query. The default disk budget is 64GB. In the coming time we will enable out-of-core for joins and group-by’s as well.

These two changes will make Polars much more resillient in high-memory workloads for casual data practicioners. And with out-of-core join and group-by on our roadmap, this resilience will improve even more.

Polars now supports the Arrow MapType directly as a Polars Map dtype. You can think of a Map as a Python dictionary, mapping keys to values. Before 2.0 the Arrow MapType was read in Polars as List(Struct({"key": ..., "value": ...})).

df = pl.DataFrame( { "user": ["alice", "bob", "carol"], "scores": pl.Series( [{"math": 90, "art": 75}, {"math": 60}, {}], dtype=pl.Map(pl.String, pl.Int64), ), "subject": ["art", "art", "math"], } )

shape: (3, 3)
┌───────┬─────────────────────────┬─────────┐
│ user ┆ scores ┆ subject │
│ --- ┆ --- ┆ --- │
│ str ┆ map[str, i64] ┆ str │
╞═══════╪═════════════════════════╪═════════╡
│ alice ┆ {"math": 90, "art": 75} ┆ art │
│ bob ┆ {"math": 60} ┆ art │
│ carol ┆ {} ┆ math │
└───────┴─────────────────────────┴─────────┘
# Key lookups and dictionary-like methods:
df.select(
 "user",
 pl.col("scores").map.get("math").alias("math"), # fixed key
 pl.col("scores").map.get(pl.col("subject")).alias("by_subject"), # key from another column
 pl.col("scores").map.contains_key("art").alias("has_art"),
 pl.col("scores").map.len().alias("n"),
 pl.col("scores").map.keys().alias("keys"),
 pl.col("scores").map.values().alias("values"),
)
┌───────┬──────┬────────────┬─────────┬─────┬─────────────────┬───────────┐
│ user ┆ math ┆ by_subject ┆ has_art ┆ n ┆ keys ┆ values │
│ --- ┆ --- ┆ --- ┆ --- ┆ --- ┆ --- ┆ --- │
│ str ┆ i64 ┆ i64 ┆ bool ┆ u32 ┆ list[str] ┆ list[i64] │
╞═══════╪══════╪════════════╪═════════╪═════╪═════════════════╪═══════════╡
│ alice ┆ 90 ┆ 75 ┆ true ┆ 2 ┆ ["math", "art"] ┆ [90, 75] │
│ bob ┆ 60 ┆ null ┆ false ┆ 1 ┆ ["math"] ┆ [60] │
│ carol ┆ null ┆ null ┆ false ┆ 0 ┆ [] ┆ [] │
└───────┴──────┴────────────┴─────────┴─────┴─────────────────┴───────────┘

As a supported dtype, the map type will now have dedicated expressions, like key lookups, iteration over values and other dictionary like methods.

Polars aims to be strict and fail fast. Errors should ideally raise up-front, not 20 minutes into a pipeline. Implicit behavior on data-mismatches should be opt-in, not a default, since those mismatches can hide bugs. This strictness has become even more valuable with the rise of AI-driven development. Agents can validate a query’s structure early by calling collect_schema(), which resolves types and catches schema-level mismatches without materializing any data. This ensures fast feedback, meaning agents and humans can iterate faster. Not all errors can be caught during compilation of the query plan, some depend on data. In these cases Polars defaults to stricter behavior to ensure inconsistencies are caught instead of silently producing different results. See previous posts) with some examples in where Polars has gotten more strict.

We are very excited that Polars 2.0 is out. Coming months we’ll improve on the road were in. Better out-of-core, better scaling at large CPU-counts and on Polars Cloud we aim to be the fastest distributed engine available. We also started working on GeoPolars and hope to deliver more news on this soon. If you find any problem with our new release, please open an issue: https://github.com/pola-rs/polars/issues. And finally, to help you with upgrading to 2.0, we have posted migration guide.