The Polars 2.0 release has arrived, and it quietly makes a case for SQL being built into the language rather than treated as an extra feature added later. The update includes out-of-core spill-to-disk support, a streaming engine, and a new Map data type. It also brings new benchmark results showing Polars performing better than DuckDB and DataFusion across common SQL tests.
A version bump was not meant to signal a major feature release, per the announcement post. Yet it has become one anyway. This release delivers plenty of reason for excitement.
What 2.0 Actually Brings
The highlights of the release break down into five parts:
- Initial out-of-core (spill-to-disk) support, now enabled by default
- Core performance improvements across the board
- Full SQL support treated as a first-class citizen
- A new Map data type
- Stricter handling of data types and explicitness, leading to faster feedback and faster AI iteration
Polars has added two major features: out-of-core support and a streaming engine. Out-of-core support means the library spills data to disk when memory is exhausted, instead of failing or grinding to a halt. The streaming engine alters how queries are gathered, so that a call to collect on a LazyFrame now defaults to using it. That change brings substantial gains in memory use and speed across most queries.
SQL as First-Class Citizen
Over the past few years, Polars has built a strong engine. Now, in 2.0, the team’s goal is for that engine to handle more kinds of work, including SQL. Making that fast required shipping many improvements to the optimizer and engine.
Polars has made SQL’s inner workings part of its standard package rather than a separate feature set. The key improvements include join reordering, improved common-subplan-elimination, and dynamic predicates and bloom filters. These are the tools that make SQL run quickly.
The Benchmark Numbers
The team tested Polars by running benchmarks on data sourced from TPC-H and TPC-DS1, comparing its performance against several other projects. Those comparisons included the current DuckDB release (1.5.6), a pre-release version of DuckDB (2.0 alpha, 2.0.0.dev2610011535), and the most recent release of DataFusion (54.0.0).
| Machine | vCPUs | RAM |
|---|---|---|
| c7a.4xlarge | 16 | 32GB |
| c7a.metal | 192 | 384GB |
Each query ran through five iterations in a warm state, with a distinct process dedicated to each query and a 60-second limit set. The file cache was emptied after each engine/benchmark pair, though not between individual queries.
The team selected the top run from each of the five attempts for comparison, measuring engines against both the total and the geometric average of those query times. The data came from tpcgen-cli built from source at commit 99bedae, with the SQL queries produced using DuckDB 1.5.6’s tpch_queries() and tpcds_queries(). The storage sat on EBS.
On c7a.4xlarge, DataFusion encountered a timeout while processing TPC-DS q72 and once on q67, and it exhausted its memory on TPC-H q18. These three queries are omitted from the results for every engine listed above.
What the Charts Show
The findings stand out. By default, Polars is the quickest on every benchmark except one. When scaled up to 192 threads, Polars carries a constant overhead that hurts small data queries. In fact, limiting Polars to just 32 core makes it competitive or winning across all benchmarks.
The root of the issue was identified by the team on their side, and they plan to address it in the upcoming release. Further details about the tests appear in the appendix, where the team also invites others to confirm the findings. A separate collection holding the test material sits at https://github.com/pola-rs/polars-2.0-benchmark.
Streaming Engine and OOC as Default
The shift in default behavior for collect on a LazyFrame is one of the largest changes in 2.0. Rather than running on the legacy engine, it now defaults to the streaming engine. This switch brings about significant gains in both memory efficiency and query performance across most use cases.
For some operations, such as join, group_by, and unpivot, the streaming engine doesn’t ensure row-order by default. Should you need to observe row-order in these operations, you can enable it by setting maintain_order=True.
Spilling to disk is turned on by default now. A query begins moving data to storage when it has used around 80% of available memory, a threshold that might require adjustment. Several operations currently allow this behavior: sorting, window functions, and many expressions, all of which can now start writing intermediate results to disk to complete their work.
The starting limit on disk space is 64GB. The team plans to turn on out-of-core support for joins and group-by operations soon. Once both of these changes are in place, Polars will be able to handle high-memory workloads much better for people who use it casually.
The New Map Data Type
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”: …})).
Here is how the new data type looks in practice:
python
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"],
})
This arrangement has a form of (3), 3, and the presentation of its contents is plain to see. It permits searches by label and dictionary-style operations upon the column.
python
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"),
)
The outcome contains six columns: math, by_subject, has_art, n, keys, and values. Keys produces a list of strings, while values generates a list of integers.
Why This Matters
The Polars 2.0 release represents a notable leap for a project that has developed quietly over time. The benchmark results make that case most directly, yet the announcement also points to a change in how Polars approaches SQL.
Polars’ streaming engine and out-of-core support strengthen its performance under heavy memory demands. The addition of the Map data type closes a gap in what the library can do. Strict type checking also speeds up development, with the team citing faster feedback and faster AI iteration as the benefits.
This is not a bold marketing stunt. It is a technical improvement that serves the people who depend on the tool day to day.
The Next Release
The scaling problem on their end has been diagnosed, and the team is hoping to address it in the next release. That represents a promise worth keeping an eye on. Meanwhile, the out-of-core join and group-by roadmap is noteworthy, since once those features are enabled, the case for resilience will be even more compelling.
Polars has been building a solid engine for the last couple of years. In 2.0, it is finally treating SQL as a first-class citizen. The benchmarks demonstrate the improvements.
This update is a factual report about the technology itself, not a promotional campaign. It describes a tool that has been improved, with tests backing up the claim. On almost every test, Polars comes out ahead as the quickest default choice.
Source material: “Release of Polars 2.0,” pola.rs.
Get the Notebook.
The day's best stories and every fresh verdict, in plain English, in your inbox by seven. One email a day, no more.

