The DataFrame library written in Rust launches version 2.0 with the streaming engine as the default mode and out-of-core enabled out of the box, reducing RAM consumption in analytical pipelines.
The organisation pola-rs, based in the Netherlands, has released Polars 2.0, the new major version of its DataFrame library written in Rust and backed by Apache Arrow. The release, which according to ecosistemastartup.com comes with benchmarks against DuckDB and DataFusion, introduces two changes that directly affect data teams: the streaming engine becomes the default mode and spill to disk (out-of-core) is enabled out of the box.
Until version 1.x, the in-memory engine was the default, and streaming required manual activation via pl.Config.set_engine_affinity(engine="streaming"). From now on, calling .collect() on a LazyFrame automatically uses streaming. In practice, most queries consume less RAM and execute faster without modifying the code. The major version jump reflects a change in contract: the streaming engine no longer guarantees the order of rows in operations like join, group_by, or unpivot. If the logic depends on observable order, maintain_order=True must be explicitly activated.
The out-of-core begins to spill data to disk when RAM consumption exceeds approximately 80% of the system's capacity, with a default disk budget of 64 GB. Operations that already support spilling include sort, window functions, and most expressions; joins and group-bys are on the immediate roadmap. This makes Polars an alternative to Spark for datasets that fit on a single large machine, according to the statement.
Version 2.0 also consolidates SQL as the official interface. In the TPC-H and TPC-DS benchmarks published by the team, Polars leads in most tests against DuckDB 1.5.6, DuckDB 2.0 alpha, and Apache Arrow DataFusion 54.0.0. On AWS instances c7a.metal (192 vCPUs, 384 GB RAM), Polars scales from 16 to 192 vCPUs, being 3.8 times faster in TPC-H and 2.2 times faster in TPC-DS (measured by total time). In the same comparison, DuckDB 1.5.6 achieves 3.2x and 1.9x; DuckDB 2.0 alpha, 2.2x and 1.5x; DataFusion, 1.7x and 1.0x.
DuckDB 1.5.6 still has an advantage in small datasets when Polars runs with its 192 threads. The team identified a parallelization overhead and recommends limiting Polars to 32 cores in those cases until it is resolved in the next release. In a previous benchmark of Polars itself (May 2025, PDS-H on c7a.24xlarge, SF-100), the streaming engine of Polars 1.30.0 completed the suite in 23.94 seconds compared to 19.65 s of DuckDB 1.3.0 and 152.27 s of the in-memory engine of Polars. Dask finished in 548.52 s and PySpark 4.0.0 in 312.43 s. The repository with the methodology is public at github.com/pola-rs/polars-2.0-benchmark.
Among the minor new features, Polars 2.0 incorporates the new dtype Map, which supports Arrow MapType natively instead of representing it as List(Struct), and adds greater strictness in error detection.

