Your Spark cluster is overkill. There, I said it.
While you’re provisioning nodes and debugging executor failures, the Rust-native DataFrame library Polars just shipped version 2.0 with a set of architectural changes that make it the default choice for high-performance data processing, and the benchmarks are frankly uncomfortable for the incumbents.
The release post, authored by Polars creator Ritchie Vink, dropped on October 6th, 2026, and it’s not messing around. Streaming engine as default. Out-of-core spill-to-disk enabled out of the box. SQL as a first-class citizen. A native Map dtype. Stricter type handling.
The headline numbers? Polars beat both DuckDB 1.5.6 and DataFusion 54.0.0 across nearly all TPC-H and TPC-DS derived benchmarks on both a 16-core and 192-core AWS instance. On the 16-core machine, Polars finished TPC-H SF10 in 3.22 seconds versus DuckDB’s 4.44 and DataFusion’s 7.41. At SF100 TPC-DS, Polars clocked 85.83 seconds while DataFusion limped in at 198.48.
Let’s dig into what actually changed and why it matters.

The Streaming Engine Is No Longer Optional
The most consequential change in 2.0 is that calling collect() on a LazyFrame now defaults to the streaming engine. Previously, you had to opt in with .collect(streaming=True). Now it’s the default path, and that’s a major version bump precisely because it changes observable behavior.
Here’s the catch: the streaming engine doesn’t guarantee row order for certain operations, specifically join, group_by, and unpivot. If your downstream code depends on the order rows appear after those operations, you need to explicitly set maintain_order=True.
This is the kind of breaking change that will surprise teams upgrading without reading the migration guide. The trade-off is real: streaming execution dramatically reduces memory footprint and improves performance on most queries. But if you’ve been relying on implicit ordering, your pipeline might produce subtly different results after upgrade.
The community reception on Hacker News has been largely positive, with the 416-point discussion focusing on the benchmark methodology and what it means for the broader data engineering ecosystem. One recurring theme: this is getting harder to dismiss as a niche tool.
For teams still comparing Polars with Pandas for high-performance data processing, the streaming default closes the memory gap that was Pandas’ last advantage.
Spill-to-Disk: Memory Is No Longer the Ceiling
Out-of-core processing is now enabled by default. When memory usage hits roughly 80% of RAM, Polars starts spilling data to disk for supported operations. The default disk budget is 64GB.
This is huge for casual data practitioners running on laptops or shared VMs. Sort operations, window functions, and many expressions can now complete even when the working set exceeds available memory. Previously, you’d hit an OOM error 20 minutes into a pipeline. Now Polars degrades gracefully.
The spilling threshold “may need tuning”, the release notes are refreshingly honest about that. And there are gaps: out-of-core support for joins and group-by operations is on the roadmap but not in this release. If your heavy queries are join-dominated, you’ll still hit memory ceilings.
What’s notable is the resiliance framing. The Polars team isn’t positioning this as “handle petabytes”, they’re positioning it as “don’t crash on your laptop when your data grows.” That’s a pragmatic engineering choice that addresses the actual pain point for most users.
This is part of the rise of DuckDB and Polars in modern data engineering pipelines, where single-node tools are increasingly good enough for workloads that used to require clusters.
The Benchmark Story: Impressive, With Caveats
Here’s where things get spicy. The benchmark methodology matters, and the Polars team was unusually transparent about it:
- Data generated with
tpcgen-cli parquetfrom a specific commit - Queries generated using DuckDB’s
tpch_queries()andtpcds_queries()functions - Each query ran 5 times hot, best time recorded
- File cache cleared between engines, not between queries
- 60-second timeout per query
The results on the c7a.4xlarge (16 vCPUs, 32GB):
| Benchmark | Polars | DuckDB 1.5.6 | DuckDB 2.0 alpha | DataFusion |
|---|---|---|---|---|
| TPC-H SF10 sum | 3.22s | 4.44s | 4.27s | 7.41s |
| TPC-DS SF10 sum | 12.59s | 15.92s | 14.13s | 34.36s |
| TPC-H SF100 sum | 34.63s | 50.17s | 42.12s | 78.04s |
| TPC-DS SF100 sum | 85.83s | 112.70s | 98.59s | 198.48s |
Polars was fastest on every single benchmark on the 16-core machine. On the 192-core c7a.metal, Polars default was fastest on all but one benchmark, but there’s an asterisk.
The 192-core results show a constant overhead when scaling Polars to 192 threads, which hurts small data queries. Polars limited to 32 threads was competitive or winning across all benchmarks, including TPC-DS SF10 where it hit 9.63 seconds versus DuckDB 1.5.6’s 12.25 seconds.
DataFusion timed out on TPC-DS q72, once on q67, and ran out of memory on TPC-H q18 on the smaller machine. Those queries were excluded for all engines, a methodology choice that’s defensible but worth noting.
The Polars team acknowledges the high-thread-count overhead and says they’ve diagnosed the cause, with a fix hoped for in the next release. The benchmark repository is publicly available if you want to replicate the results yourself.
The scaling story across cores is telling. Moving from 16 to 192 vCPUs at SF100 makes Polars 3.8x faster on TPC-H and 2.2x faster on TPC-DS. DuckDB 1.5.6 gets 3.2x and 1.9x. DataFusion only gets 1.7x and 1.0x, literally no speedup on TPC-DS at higher core counts. That’s the parallel execution story: Polars’ Rust-based scheduler with morsel-driven streaming execution scales meaningfully better than its competitors.
What’s Driving the Performance: Architecture Under the Hood
The GitDiagram analysis of the pola-rs/polars repository breaks down the architecture into clean layers:
- User APIs: Python DataFrame API, LazyFrame API, Expressions, Rust API, SQL context
- Query Planning: Subquery rewrites, query optimizer with projection pushdown
- Execution Runtime: Engine selection, execution engines, streaming execution, compute kernels
- Data and I/O: DataFrame core, Arrow arrays, Data I/O, Parquet support
- Interoperability: Python bindings, Arrow FFI, remote engine
The key architectural decisions:
- Arrow as the native memory format, Polars doesn’t convert data structures, it operates directly on Arrow arrays. This eliminates serialization overhead and enables zero-copy sharing across the Rust/Python boundary via Arrow FFI.
- Lazy evaluation by default, The LazyFrame API builds a query plan that gets optimized before execution. Projection pushdown, predicate pushdown, and common-subplan elimination happen before any data is touched.
- Streaming execution with morsels, The streaming engine processes data in small chunks (morsels) that can be distributed across threads without loading everything into memory.
- Join reordering and dynamic predicates/bloom filters, The optimizer improvements in 2.0 make complex joins significantly faster, which is where the TPC-DS gains come from.
- Common-subplan elimination, Repeated computations in a query plan get computed once and reused instead of recomputed per occurrence.
These aren’t incremental tweaks. This is a fundamentally different design philosophy from Pandas, which materializes everything eagerly. And it’s different from Spark, which adds network serialization overhead even for single-node workloads.
For teams contrasting Polars and Spark in the context of lightweight, high-performance ETL, the streaming improvements in 2.0 make the case stronger: why spin up a cluster when a single machine with Polars handles your workload faster?
The New Map dtype: More Than Convenience
Polars now natively supports Arrow’s MapType as a first-class Map dtype. Previously, this was read as List(Struct({"key": ..., "value": ...})), functional but clunky for dictionary-like operations.
The new API is genuinely nicer:
df = pl.DataFrame(
{
"user": ["alice", "bob", "carol"],
"scores": pl.Series(
[{"math": 90, "art": 75}, {"math": 60}, {}],
dtype=pl.Map(pl.String, pl.Int64),
),
"subject": ["art", "art", "math"],
}
)
Key lookups work cleanly:
df.select(
"user",
pl.col("scores").map.get("math").alias("math"),
pl.col("scores").map.get(pl.col("subject")).alias("by_subject"),
pl.col("scores").map.contains_key("art").alias("has_art"),
pl.col("scores").map.len().alias("n"),
pl.col("scores").map.keys().alias("keys"),
pl.col("scores").map.values().alias("values"),
)
This isn’t just a syntax improvement. A dedicated dtype means the query optimizer understands the structure, enabling more aggressive optimizations. Dictionary-like operations that previously required unnesting and re-joining now work directly on the column.
Stricter Polars: Fail Fast, Fail Early
The push toward stricter type handling is the least glamorous but arguably most important change, especially in the age of AI-driven development.
Polars aims to raise errors up front rather than 20 minutes into a pipeline. Implicit behavior on data mismatches is now opt-in, not default. This is a deliberate design choice that prioritizes correctness feedback over convenience.
The collect_schema() method lets you validate a query’s structure without materializing any data. For AI agents iterating on queries, this means schema errors get caught instantly instead of after a full execution cycle. The release notes explicitly frame this as a feature for “faster AI iteration.”
This is the kind of thing that doesn’t show up in benchmarks but matters enormously in practice. Evaluating over-engineered enterprise platforms versus lightweight, efficient tools like Polars often comes down to iteration speed, and stricter validation is a meaningful contributor.
The Verdict: What Should You Do?
Polars 2.0 is a significant step forward, but the upgrade isn’t frictionless.
Do this first:
- Read the migration guide
- Audit your pipelines for row-order dependencies, especially after joins and group-bys
- Test with
maintain_order=Trueif you find ordering issues - Check memory thresholds, the 80% spill trigger may need tuning for your workloads
Watch for:
- Out-of-core joins and group-bys (roadmap, not shipped)
- The 192-thread overhead fix in the next release
- GeoPolars, announced as in progress
The benchmark numbers deserve healthy skepticism until independently replicated. Vendor-run benchmarks have a way of favoring the vendor. But the architecture is sound, the methodology was transparent, and the code is public. That’s more than most projects offer.
The data engineering landscape is shifting. Spark isn’t dead, Pandas isn’t irrelevant, and DuckDB won’t disappear. But the rise of DuckDB and Polars represents a fundamental realignment: for a growing class of workloads, the best tool isn’t a distributed cluster, it’s a well-architected single-node engine that uses your hardware efficiently.
Polars 2.0 makes that case harder to ignore.
Benchmark data and methodology details are derived from the official release post. The benchmarks are derived from TPC-H and TPC-DS and are not comparable to published TPC results.




