TL;DR
Get the latest gadgets delivered free — and shop member deals
- Fast, free delivery on millions of items
- Access to Prime Big Deal Days deals on October 6–7
- Prime Video, Amazon Music and more included
Polars 2.0 makes its streaming engine the default for LazyFrame.collect() and enables initial out-of-core spilling, which can let supported queries use disk when memory fills. The project also reports strong results in its own TPC-H and TPC-DS tests, but those benchmarks have specific test conditions and the release changes row-order guarantees for some operations.
Polars has released version 2.0, making its streaming engine the default when users collect a LazyFrame and enabling an initial form of out-of-core processing that can spill supported operations to disk. The update also elevates SQL support and introduces a Map data type, while changing the default row-order behavior of some operations—an important compatibility consideration for existing workloads.
Under the new default, calling collect() on a LazyFrame uses the streaming engine. Polars says this can bring memory and performance improvements for many queries. However, the engine does not preserve observable row order by default for certain operations, including joins, group-bys and unpivots. Users who need order preserved can set maintain_order=True for relevant operations.
Version 2.0 also enables spill-to-disk support by default. Polars says spilling starts at about 80% of system RAM, a threshold the project says may need tuning, and the default disk budget is 64 GB. Current support covers operations such as sorts, window functions and many expressions; joins and group-bys are not yet supported for out-of-core execution, though the project says they are on its roadmap.
The release adds a native Map dtype for Arrow MapType data, representing key-value mappings in a form Polars can work with directly. The project also cites optimizer and engine changes, including join reordering, improved common-subplan elimination, dynamic predicates and bloom filters. It describes SQL as a first-class interface in this release, broadening the workloads it intends the engine to serve.
Memory Use and Row-Order Changes
The default shift could affect both how much memory a query needs and whether its output appears in the same order as before. Streaming and disk spilling may allow some larger or memory-intensive jobs to complete without requiring all intermediate data to remain in RAM. That is useful for analysts and developers working with datasets that approach available memory, although the benefit depends on whether the specific operations in a query are supported.
The row-order change is a practical migration issue: users whose downstream steps depend on a particular order may need to request it explicitly. Polars presents the release as a step toward faster and more resilient workloads, but its own announcement does not quantify performance improvements across all users or datasets. Its benchmark results are evidence for the tested setup, not a guarantee for every workload.
As an affiliate, we earn on qualifying purchases.
SQL Benchmarks and Test Conditions
To support its performance claims, Polars tested its SQL engine on data derived from TPC-H and TPC-DS, comparing it with DuckDB 1.5.6, a DuckDB 2.0 alpha build and DataFusion 54.0.0. The tests ran on two AWS machine configurations: a c7a.4xlarge with 16 vCPUs and 32 GB of memory, and a c7a.metal with 192 vCPUs and 384 GB.
According to the release post, each query ran five times in a hot setting, with a separate process per query and a 60-second timeout; the best of five runs was used. Polars says it was fastest on all but one of the reported benchmark comparisons. DataFusion timed out on TPC-DS query 72, once on query 67, and ran out of memory on TPC-H query 18 on the smaller machine; those queries were excluded from results for all engines. Polars also reports overhead when scaling to 192 threads, which hurt small-data queries, and says it hopes to address that in a later release. The project published a repository to support replication.
“Calling collect on a LazyFrame will now default to the streaming engine.”
— Polars, in its version 2.0 release post
As an affiliate, we earn on qualifying purchases.
Limits Still Affect Disk Spilling
Out-of-core coverage is incomplete: the release post says joins and group-bys are not yet supported for spilling, and it does not give a date for adding them. The stated spill threshold of about 80% of RAM may also need tuning. It is not clear from the source how performance or disk use will vary across different hardware, query plans and datasets.
The benchmark findings come from tests designed and reported by Polars. Although the project provides a repository for replication, results may differ under other settings; the release itself notes scaling overhead on the 192-vCPU machine. The source also does not provide a detailed migration guide or quantify how many existing applications rely on the prior row-order behavior.
As an affiliate, we earn on qualifying purchases.
Roadmap for Broader Query Support
Polars says it plans to add out-of-core support for joins and group-bys, which would extend disk-backed execution to more query patterns. It also says it has diagnosed the overhead affecting some workloads at 192 threads and hopes to fix it in the next release, without specifying a release date in the supplied post.
For now, users evaluating the upgrade can review the affected operations, check whether downstream logic requires stable row order, and test representative queries against their own data. Polars encourages readers to reproduce its published SQL benchmarks; actual results and migration needs remain workload-dependent.
As an affiliate, we earn on qualifying purchases.
Key Questions
What is the main change in Polars 2.0?
LazyFrame.collect() now defaults to the streaming engine, and initial spill-to-disk support is enabled by default for supported operations.
Does Polars 2.0 always preserve row order?
No. Polars says the streaming engine does not guarantee observable row order by default for some operations, including joins, group-bys and unpivots. Users can opt in with maintain_order=True where needed.
Which operations can spill to disk?
The release post lists sorts, window functions and many expressions as supported. Joins and group-bys are not yet included, though Polars says it plans to add them.
Are Polars’ benchmark results independently verified?
The supplied results were run and reported by Polars. The project has shared a benchmark repository for replication, but the source does not describe the results as an independent evaluation.
Source: hn
Fall Picks
fall essentials
As an affiliate, we earn on qualifying purchases.
