TL;DR
Get privacy and security gear delivered free with Prime
- Fast, free delivery on millions of items
- Prime Video, Amazon Music and more included
- Member-only deals all year
Polars says version 2.0 makes its streaming engine the default for LazyFrame collection and enables initial spill-to-disk support. The project also reports strong results in its TPC-H and TPC-DS SQL benchmarks, while noting limitations in supported out-of-core operations and scaling to many CPU threads.
Polars has shipped version 2.0, making its streaming engine the default for collecting lazy queries and enabling initial spill-to-disk support. The release also expands SQL support and adds a Map data type, changes that the project says are intended to improve memory use and broaden the workloads its query engine can handle.
In Polars 2.0, calling collect on a LazyFrame uses the streaming engine by default. The project says this can improve memory use and performance on many queries. However, streaming does not preserve observable row order by default for some operations, including joins, group-bys and unpivots. Users who need that behavior can set maintain_order=True.
The release enables out-of-core execution, which spills supported work to disk when memory use reaches about 80% of RAM; the default disk budget is 64 GB. The initial support covers operations such as sorting, window functions and many expressions. Polars says joins and group-bys are not yet covered, but are planned for later work.
Polars also reports that its SQL coverage has grown and describes SQL as a first-class interface in version 2.0. The project highlights optimizer and engine changes including join reordering, improved common-subplan elimination, dynamic predicates and bloom filters. The release adds a Map dtype for Arrow MapType data, represented previously in Polars as a list of key-value structs.
More Queries Can Run Within Memory Limits
Making streaming the default changes how existing lazy queries execute, so users may see different performance and memory behavior after upgrading. The row-order caveat also means some workflows must explicitly request order preservation and check that their results meet expectations.
Spilling supported operations to disk could help users complete work that would otherwise exceed available memory. The current limits matter: joins and group-bys are not yet included, and Polars sets a default disk budget of 64 GB. Teams handling large datasets will need to consider both the supported operations and local storage capacity.
The SQL changes may make Polars relevant to workloads currently run with other analytical engines. But the performance evidence comes from benchmarks conducted and reported by Polars, under specified hardware and query conditions; it is not a guarantee of results on other datasets or systems.
high performance data analysis laptop
As an affiliate, we earn on qualifying purchases.
As an affiliate, we earn on qualifying purchases.
How Polars Tested Its SQL Engine
To compare SQL performance, Polars ran queries derived from TPC-H and TPC-DS against Polars, DuckDB 1.5.6, DuckDB 2.0 alpha and DataFusion 54.0.0. Tests used an AWS c7a.4xlarge machine with 16 virtual CPUs and 32 GB of RAM, and a c7a.metal machine with 192 virtual CPUs and 384 GB of RAM. Each query ran five times in a hot setting, and the best run was used in comparisons based on total and geometric-mean query times.
Polars says its default configuration was fastest in all but one benchmark in those tests. It also reports that 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 says its own performance has a constant overhead when scaled to 192 threads, hurting small queries, and that it hopes to address this in a future release. The project published a repository so others can replicate the benchmark.
Polars characterizes 2.0 as a release focused on enabling SQL and improving core performance, even though it did not set out to make the version a major feature release. The major version change also accommodates the altered row-order behavior of the default streaming engine.
“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 of the First Spill Support
The release post does not specify a calendar date for the launch. It also does not say when out-of-core support for joins and group-bys will arrive. The stated threshold of about 80% of RAM may need tuning, according to Polars, and the post does not establish how the default behaves across different machines or workloads.
The reported benchmark results have defined test conditions, including hot runs, best-of-five timing and excluded queries. The post encourages replication, but results for other data, hardware, settings or query mixes remain unknown. Polars also says its 192-thread overhead may affect smaller queries and hopes to address it in a later release.
out-of-core data processing software
As an affiliate, we earn on qualifying purchases.
As an affiliate, we earn on qualifying purchases.
Joins and Group-Bys Remain on the Roadmap
Polars says it plans to extend out-of-core execution to joins and group-bys. It also says it has diagnosed the overhead seen when scaling to 192 threads and hopes to fix it in the next release, though it gives no date or detailed schedule. Users can consult the benchmark repository to reproduce the published SQL comparisons and evaluate the new defaults against their own workloads.
As an affiliate, we earn on qualifying purchases.
Key Questions
What changed in Polars 2.0?
LazyFrame collection now defaults to the streaming engine. Version 2.0 also enables initial spill-to-disk support, expands SQL features and adds a Map dtype.
Does the streaming engine preserve row order?
Not by default for some operations, including joins, group-bys and unpivots. Polars says users can request order preservation with maintain_order=True.
Which operations can spill to disk?
The release post lists sorting, window functions and many expressions among the supported operations. Out-of-core joins and group-bys are planned, but are not included in the initial support described.
Did Polars prove it is faster than DuckDB and DataFusion?
Polars reports that its default configuration was fastest in all but one of its TPC-H and TPC-DS benchmarks. Those results reflect the project’s stated test setup and do not establish performance for every workload or system.
Source: hn
Halloween Picks
halloween
As an affiliate, we earn on qualifying purchases.
