Polars Main Adds an Out of Core Sort and Tighter Scan Pruning


Polars on main took 94 commits in this window, across 763 files, with 29,263 insertions and 5,229 deletions. Streaming execution and scans are where jobs will feel it: a sort that spills, a scalar window node, and less repeated IO on Parquet, IPC, and Iceberg. A few fixes change results, so a new pin of main needs a retest.

The streaming sort samples the input, builds branchless classification trees, and sorts the batches in parallel. Spill frames hold the state, so the node can sort past RAM. The code is isolated in sort/mod.rs.

Each bucket reserved flush_rows plus the full frame height. It now reserves flush_rows plus the rows that bucket received.

Scalar windows reduce one column to one value per partition and write it back when the morsels replay. At 65536 estimated groups or fewer, each pipeline reduces locally and the partials merge. Above that, hash partitions reduce in blocks of about 4096 rows, and a configured out of core budget caps a block at one eighth of that budget. Morsels can spill. Partition keys stay until the reduction finishes, and the group id of each row stays until that morsel is emitted.

Windows that ignore row order are planned as elementwise work, in the in memory engine as well as in streaming, so a scan can drop order preservation. Those edits sit in lower_ir.rs, lower_expr.rs, and to_graph.rs. Physical nodes now point back at the IR node that query monitoring shows.

Warm Parquet row groups are read by the decode task on Linux, and the prefetch semaphore is skipped. The q14 run cited with the patch cut futex calls from about 630,000 to about 34,000. An environment variable controls the path. Cold row groups still prefetch.

Single file IPC decodes row groups out of order when the query ignores order. A 4.7 GiB S3 object on a 64 vCPU host in eu-west-1 moved from 2.2 s and 3.1 GB peak to 2.1 s and 0.8 GB, because row groups were no longer queued. Remote files follow the same rule. The footer is one suffix range read, so the extra HEAD is gone.

Iceberg manifests go into a process wide LRU (default 64 MiB, and 0 disables it), keyed by path plus a hash of FileIO properties, so different explicit credentials do not share an entry. REST catalogs are excluded, and environment credentials are assumed stable for the process. A laptop select(pl.len()) over 100 S3 manifests went from 649 ms of planning and 867 ms total to 140 ms and 349 ms (test_iceberg.py).

Fallible Hive filters, such as pl.date or a strict cast, stay above the scan and are also tried on the partition values. Success skips files that cannot match. Failure keeps every file, and the original error still surfaces. The try now runs after partition rewriting, in predicate_pushdown/mod.rs, so removing every partition does not force the rewrite to build an empty union. The row filter stays unless no_residual_predicate is set.

scan_lance is unstable. A 59.3 GiB scan of 100 objects went from 19.96 s and 7.35 GB max RSS to 13.44 s and 12.08 GB.

Scalar columns are no longer materialized just to measure them. estimated_size in column/mod.rs takes a boolean. False counts one copy of an unmaterialized scalar. True counts the repeated length, still without allocating. Python DataFrame.estimated_size() passes false, so pl.repeat(1, 2**31) reports 8 bytes. It used to report 8 * 2**31. Spill uses the compact number. Join sampling and sink sizing still ask for the expanded one. Python has no flag for that expanded figure.

Build threads share one runtime bloom, between 64 KiB and 32 MiB, after the buffered row budget. Sum, min, and max store the null mask once per block when 64 or fewer groups are seen and the column has no nulls.

Count guarded sums turn when(x.count() > 0).then(x.sum()).otherwise(null) into a sum that is null on an empty group. Explain prints .sum(null_on_empty=true). SQL SUM needs that null, and the separate count is gone. The reducer is in sum.rs.

Dynamic window k starts at origin + every * k, matching datetime_range. The streaming engine and the in memory engine used to disagree. From a 31 January anchor, the old month starts were Jan 31, Feb 29, Mar 29, Apr 29, May 29. They are now Jan 31, Feb 29, Mar 31, Apr 30, May 31. Empty dynamic group by, including rolling group by with a placement, no longer panics when there are no windows to output.

Rows wider than the first used to lose the extra values.

pl.DataFrame([(1,), (2, 3, 4, 5)], orient="row")
# ShapeError: row at index 1 has length 4 (expected 1)

The old result was (1,) and (2,). Shorter rows already failed, with a message about column height.

list.contains on sliced list chunks read offsets from a rechunked copy and values from the original chunks, so a column needle saw the wrong elements. A case that returned [false, false] now returns [true, true]. is_in on a list column had the same bug. Grouped sort_by sorts per group when reverse or sort moves rows and the group lengths stay the same. is_sorted on List, Array, and Map returns a boolean, using the same row encoding as sort. It used to raise InvalidOperationError.

from_dicts writes straight into column buffers. One million rows of five string columns, on an M4 Max, moved from 432 ms and 740 MB peak to 258 ms and 236 MB. Correlated SQL reuses an outer aggregate the subquery repeats, and a LIKE with only a leading or trailing % lowers to starts_with, ends_with, or contains_literal.

Diff month windows anchored on 31 January, and any budget that reads Python estimated_size() on a repeated column. Rust callers of estimated_size now pass a boolean.

The author plans to revisit the streaming sort in sort/mod.rs. scan_lance was faster on the published scan and used more RAM, and the entry point is still unstable.

Iceberg manifest caching skips REST catalogs, and a rotated environment credential stays invisible until the process exits. Set the budget to 0 when that assumption is wrong. Temporary names gained a crate token (SQL, PLAN, PHYS, OPS, MEM).