· via Hacker News – Front Page (hnrss.org)
Polars 2.0 makes streaming the default, adds out-of-core spilling and first-class SQL
Polars 2.0 ships with the streaming engine on by default and out-of-core spilling enabled, promotes SQL to a first-class interface with benchmark wins over DuckDB and DataFusion, and adds a native Map dtype.

The Polars project has shipped version 2.0 of its dataframe library, turning what was planned as a compatibility-focused version bump into a substantial feature release. According to the announcement post on pola.rs, which reached the Hacker News front page on October 6, the highlights are default out-of-core execution, broad performance work across the core engine, SQL as a first-class interface, a new Map datatype, and stricter handling of dtypes.
Streaming and spill-to-disk become the default
The change the Polars team itself calls the highest-impact one is that calling collect on a LazyFrame now runs through the streaming engine by default, which the post says brings large memory and performance improvements on most queries. It is also the main reason the release needed a major version number: the streaming engine does not guarantee row order for certain operations — joins, group_by and unpivot are listed — so code that depends on observable ordering must explicitly opt back in with maintain_order=True.
Out-of-core execution, meaning spill-to-disk under memory pressure, is switched on by default as well. Per the post, spilling starts at roughly 80% of available RAM — a threshold the team acknowledges may need tuning — with a default disk budget of 64GB. Sorting, window functions and many expressions can currently spill; joins and group-bys are expected to gain out-of-core support later, which the team argues will make the library more resilient for memory-heavy workloads.
SQL as a first-class citizen, backed by vendor benchmarks
Polars 2.0 marks the point where SQL is treated as a first-class way to query data with the library. SQL coverage has grown sharply over recent months, and the post credits optimizer work — join reordering, much better common-subplan elimination, and dynamic predicates with bloom filters — for making those workloads run fast.
To back that up, the team benchmarked Polars SQL on TPC-H and TPC-DS derived data against DuckDB 1.5.6, a DuckDB 2.0 alpha build and DataFusion 54.0.0, on a 16-vCPU, 32GB c7a.4xlarge instance and a 192-vCPU, 384GB c7a.metal machine. Each query ran five times in a hot setting, in a separate process with a 60-second timeout; the best run was taken and engines were compared on both the sum and the geometric mean of query times.
By the team's account, default Polars was fastest on all but one benchmark. The post is candid about a caveat: Polars carries a constant overhead when scaling to 192 threads, which hurts small-data queries, and Polars capped at 32 cores was competitive or winning across the board. The team says the cause has been diagnosed and hopes to ship a fix in the next release. DataFusion timed out on TPC-DS query 72 (and once on query 67) and ran out of memory on TPC-H query 18 on the smaller machine; those queries were excluded for every engine. Because these are vendor-run benchmarks, it matters that the team has shared a public benchmark repository and is encouraging others to replicate the results.
A native Map dtype and stricter semantics
Polars now supports the Arrow MapType directly as a Map dtype — conceptually a Python dictionary mapping keys to values — where previously it was read in as a list of key/value structs. The new type ships with dedicated expressions for dictionary-style access: looking up fixed keys or keys taken from another column, checking whether a key exists, and retrieving lengths, key lists and value lists.
The release also doubles down on strictness. The stated philosophy is that errors should surface up front rather than twenty minutes into a pipeline, and that implicit behavior on data mismatches should be opt-in rather than the default, since mismatches can hide bugs. The team explicitly frames this as valuable for AI-assisted development: agents and humans can call collect_schema() to resolve types and catch schema-level mismatches without materializing any data, shortening feedback loops. Where errors can only be caught against actual data, Polars now defaults to stricter behavior rather than silently producing inconsistent results.
What comes next
According to the post, the roadmap includes better out-of-core support, improved scaling at high CPU counts, a push to make Polars Cloud the fastest distributed engine available, and early work on a GeoPolars project. A migration guide accompanies the release, and problems can be reported on the project's GitHub issue tracker.
Why it matters
Polars has become one of the most widely adopted dataframe libraries in the Python and Rust data ecosystems, and 2.0 changes default behavior in ways that can silently affect existing pipelines — particularly row ordering under the streaming engine, which makes the migration guide essential reading before upgrading. Default out-of-core execution lowers the barrier to working with datasets larger than memory, while first-class SQL puts Polars into more direct competition with DuckDB and DataFusion. The TPC benchmark claims are the team's own, but the published methodology and replication repository invite independent scrutiny. The stricter schema checking, meanwhile, is a notable example of a data library deliberately designing for agentic coding workflows.
- #polars
- #dataframe
- #data-engineering
- #rust
- #sql