Polars 2.0 makes streaming the default and adds spill-to-disk queries
The Amsterdam-based Polars says initial disk spilling covers sorting, window functions and many expressions; joins and group-bys remain in-memory.
By RuntimeWire Staff · Published
Primary source: Polars
Why it matters
Polars 2.0 sets streaming as the default, a change meant to make memory-efficient execution the ordinary path, and adds limited spill-to-disk support as Vink works to extend one DataFrame API toward larger production queries. Developers also need to handle row order explicitly, while joins and group-bys still cannot spill to disk.
Ritchie Vink, who started Polars as an open-source project in 2020, has made the library's streaming engine the default in Polars 2.0, released on October 6th. The change is meant to make memory-efficient execution the ordinary path for users, while initial spill-to-disk support lets some queries continue when data no longer fits in RAM.
Vink's release post describes a project that began as a way to learn about Rust and query engines. He later wrote that frustration with pandas' performance and design helped turn the side project into a faster DataFrame library. In August, he pitched the coming default change in practical terms: "your code will be faster without doing anything" in a post on X.
That promise comes with a behavior change developers need to account for. Calling collect on a LazyFrame now uses the streaming engine by default, and operations such as joins, group-bys and unpivoting no longer guarantee row order unless users request it with maintain_order=True. Polars had outlined the change and migration implications in its September release-candidate announcement; the final release makes it the default.
From a faster DataFrame to a query engine
Vink announced Polars as a company with co-founder Chiel Peters in 2023 in a post. The first commit to the open-source project came in June 2020. Vink said it started as a pet project to learn about query engines, Apache Arrow and Rust. Peters, Polars' other co-founder, previously served as CTO of Amsterdam data consultancy Xomnia.
Polars 2.0 adds first-class SQL support alongside changes to the optimizer, including join reordering, common-subplan elimination and dynamic predicates and bloom filters. It also introduces an Arrow-compatible Map data type and tighter rules around data types and explicitness. These changes are intended to make the same Rust-based engine more useful to SQL users and catch mistakes earlier in development.
The release also enables out-of-core processing for supported operations by default. Polars says it begins spilling to disk at about 80% of RAM use and sets a default disk budget of 64 GB. The initial support covers sorting, window functions and many expressions. Joins and group-bys remain outside that support for now, so the release does not yet make every large query independent of available memory.
A benchmark claim with a company-run setup
Polars says its SQL engine was fastest in all but one of its TPC-H and TPC-DS benchmark comparisons against DataFusion, DuckDB 1.5.6 and a DuckDB 2.0 alpha. Its tests used two AWS machine types and five hot runs per query, with the best run reported. DataFusion timed out or ran out of memory on several queries; Polars excluded those queries from the results for every engine.
Those are Polars-run comparisons on data derived from TPC benchmarks, not official TPC-certified results. The company has published a benchmark repository so others can inspect and reproduce the setup. The release also notes a scaling weakness: Polars' default configuration was not fastest in every comparison on the 192-vCPU machine, where limiting it to 32 threads was competitive or faster. Polars says it has diagnosed the overhead and hopes to address it in a future release.
Polars has also been building beyond the local library. In 2025, it raised an €18 million Series A led by Accel, with Bain Capital Ventures participating again. Vink said the funds would support a distributed engine and Polars Cloud, bringing the same API to managed and customer-run infrastructure. Polars launched its distributed engine on Kubernetes in June 2026.
With streaming as the default, queries can make better use of memory, and supported workloads can spill to disk. Out-of-core joins and group-bys are still planned work, and the benchmark's large-machine results show that more cores do not automatically translate into faster queries.
RuntimeWire covered the release candidate in September, when Vink's team set out the streaming default and the row-order change. The final release now puts those trade-offs into users' day-to-day code.