Polars 2.0 is a major release that makes SQL a first-class citizen, ships a streaming execution engine and enables out-of-core (spill-to-disk) processing by default, and delivers a suite of core optimizer and engine improvements (join reordering, stronger common-subplan elimination, dynamic predicates/bloom filters). LazyFrame.collect now uses the streaming engine by default for much lower memory use and better performance, with maintain_order=True available when row-order must be preserved. Out-of-core starts spilling at roughly 80% of RAM with a default 64 GB disk budget and currently supports sort, window functions and many expressions; out-of-core joins and group-bys are on the roadmap. Benchmarks on c7a.4xlarge and c7a.metal against DuckDB (1.5.6 and a 2.0 alpha) and DataFusion show Polars fastest on nearly all TPC-H/TPC-DS derived queries, though a scaling overhead at 192 threads was observed and is being addressed.
The release adds a native Map dtype mapping Arrow MapType to dictionary-like operations (key lookups, contains_key, len, keys, values) with dedicated expressions, simplifying and speeding map manipulations. Polars tightens type strictness and fail-fast behavior to surface schema mismatches early - collect_schema() resolves types without materializing data to aid fast iteration and AI-driven workflows. Roadmap items include broader out-of-core support, better multi-core scaling, GeoPolars and faster distributed execution on Polars Cloud; a migration guide and an issue tracker are provided for upgrade feedback.
Summary generated by AI from the linked article. hn.today is not affiliated with Hacker News or Y Combinator.