Data & documents

Polars 2.0 brings out-of-core processing and faster SQL queries

Polars 2.0 enables spill-to-disk processing by default and leads DuckDB in SQL benchmarks, offering stricter type safety for AI workflows.

A metallic data cube splitting into blocks with glowing circuit lines.
Illustration generated for this article

The Polars team has released version 2.0 of their data processing library, marking a significant shift in how the engine handles memory and SQL workloads. Published on October 6, 2026, this major update introduces out-of-core processing as a default feature and positions Polars as a top performer in standard SQL benchmarks against competitors like DuckDB and DataFusion.

What happened

This release focuses on resilience and performance rather than just adding new features. The most impactful change is that calling collect on a LazyFrame now defaults to the streaming engine. This shift allows Polars to handle datasets larger than available RAM by spilling temporary data to disk. The system begins this spill-to-disk process when memory usage hits approximately 80% of the available RAM, with a default disk budget of 64GB. Currently, operations like sorting, window functions, and many expressions support this out-of-core behavior, with joins and group-by operations planned for future updates.

Because the streaming engine does not guarantee row order for operations such as join, group_by, and unpivot, this change required a major version bump. Users who depend on specific row ordering must now explicitly set maintain_order=True. This default behavior aims to make Polars more robust for casual data practitioners working with high-memory workloads, preventing crashes that previously occurred when datasets exceeded physical memory limits.

In addition to memory management, Polars 2.0 treats SQL as a first-class citizen. The library has dramatically increased its SQL coverage and improved its optimizer with better join reordering, common-subplan elimination, and dynamic predicates. These enhancements allow Polars to execute complex SQL queries more efficiently, bridging the gap between programmatic data manipulation and traditional database interactions.

How it works

The performance gains in SQL execution come from deep optimizations in the query engine. Polars now leverages bloom filters and dynamic predicates to reduce the amount of data processed during query execution. By eliminating common subplans and reordering joins effectively, the engine minimizes redundant computations. These technical improvements enable Polars to compete directly with established analytical databases.

To validate these claims, the team ran benchmarks using TPC-H and TPC-DS derived data on AWS c7a instances. They compared Polars against DuckDB 1.5.6, DuckDB 2.0 alpha, and DataFusion 54.0.0. The tests involved running each query five times in a hot setting, clearing the file cache between engines, and taking the best runtime. The results showed that Polars was the fastest engine in almost all benchmarks on both 16-vCPU and 192-vCPU machines, although it exhibited some overhead on small data queries when scaled to 192 threads.

Key details

  • Out-of-core processing is enabled by default, spilling to disk at ~80% RAM usage with a 64GB default disk budget.
  • The streaming engine is now the default for collect, which may change row order unless maintain_order=True is set.
  • Polars led DuckDB and DataFusion in TPC-H and TPC-DS1 benchmarks on both c7a.4xlarge and c7a.metal instances.
  • A new Map dtype supports Arrow MapType directly, allowing dictionary-like key lookups and iterations.
  • Stricter type checking and collect_schema() enable faster feedback loops for AI agents and developers.
  • DataFusion timed out or ran out of memory on several queries where Polars and DuckDB completed successfully.

Why it matters

For software engineers and data scientists, the default out-of-core support means greater reliability when processing large datasets. Previously, exceeding memory limits would crash the process, requiring manual chunking or external tools. Now, Polars can gracefully handle larger-than-memory workloads by using disk space, making it easier to build resilient data pipelines without extensive infrastructure tuning. This is particularly valuable for teams that do not have dedicated data engineering resources to manage complex distributed systems.

The stricter type system and schema validation also address a growing need in AI-driven development. As more developers use AI agents to generate code, early error detection becomes critical. By failing fast on schema mismatches through collect_schema(), Polars helps agents and humans iterate quicker. This reduces the time spent debugging silent failures or incorrect data types deep in a pipeline, leading to more maintainable and correct data processing code.

What you can do

  • Upgrade to Polars 2.0 and review the migration guide provided by the team to handle breaking changes.
  • Test your existing pipelines for row-order dependencies and add maintain_order=True where necessary.
  • Experiment with the new Map dtype for handling nested key-value data structures more efficiently.
  • Run SQL queries directly in Polars to leverage the improved optimizer and benchmark performance against your current stack.
  • Use collect_schema() in your development workflow to catch type errors before executing heavy data operations.
  • Monitor memory usage and adjust the spill-to-disk threshold if your workload requires different resource allocation.

Tools from the Bytechap store

$89

DocBento

Self-hosted document management that reads every scan and answers with page citations.

Live demo

Keep reading

All stories