💥 Breaking changes
- Run SQL window functions over grouped rows on the aggregated rows, and evaluate QUALIFY before the projection (#29746)
- Type exact SQL numeric literals as Decimal and truncate SQL
%/DIV (#29536) - Allow deterministic expression plugins to opt into CSE/CSPE (#29428)
- Read Parquet ENUM type as pl.String (#29331)
- More map operations (#29296)
- Deprecate cut/qcut (#29329)
🚀 Performance improvements
- Fix phases ending too soon due to join sampling (#29752)
- Group record batch fetches of remote IPC scans (#29750)
- Mark group-by node as memory intensive pipeline blocker (#29749)
- Don't inline slow OOC memory drift path (#29745)
- Only prefetch on unspill when there is memory headroom (#29701)
- Pushdown filters to scan_lance (#29538)
- Do not materialize equal
ScalarColumns inColumn::append/extend(#28989) - Evaluate scalar windows in the streaming window node (#29698)
- Don't materialize scalar columns in
DataFrame::estimated_size(#29711) - Defer cached pread from prefetch to decode for Parquet (#29702)
- Decode IPC scans out of order when order is not observed (#29675)
- IPC metadata suffix fetch (#29692)
- Speed up null tracking in streaming group-by sum, min and max (#29691)
- Cache Iceberg manifest files across scans (#29623)
- Improve pushdown support for fallible Hive predicates (#29594)
- Streamline
BytecodeParsermethod rewriting (#29632) - Stream
from_dictsrecords straight into column buffers (#29656) - Read correlated SQL aggregates from the outer join when the subquery repeats it (#29685)
- Emit files unordered for remote IPC files when order not observed (#29678)
- Don't keep input order for order-insensitive windows in all engines (#29680)
- Fix parallelization in top-k reducer (#29684)
- Share one bloom filter between the build threads of a runtime filter (#29677)
- Add streaming (out-of-core) sort (#29657)
- Insert keys that start a run of equal keys into the hot group-by table right away (#29673)
- Prefetch key-row hash table lookups in blocks (#29663)
- Lower count-guarded sums to a sum that is null on empty input (#29668)
- Use single PUT for small cloud uploads (#29603)
- Don't cache plain scans in the streaming engine (#29654)
- Change single-key group-by hash table layout and add prefetching (#29661)
- Use the integer ranges of sampled parquet row groups in filter estimates (#29655)
- Raise HTTP rate-limiter read init default (#29584)
- Change join hash table layout and add prefetching (#29653)
- Add a fast-path for
as_listwith only a single input (#29592) - Use one runtime filter builder per thread for a finished sample (#29646)
- Compute decimal addition, subtraction and multiplication in i64 when values fit (#29638)
- Spread decimal sums and means over lanes and check the sum's 38 digits once (#29637)
- Compute decimal addition, subtraction and multiplication without per-row validity (#29634)
- Skip or reuse UDF disassembly in
BytecodeParser(#29625) - Rescale a single decimal operand once in addition and subtraction (#29631)
- Dynamic hot table size in streaming GroupBy (#29629)
- Lower gathering from literal as elementwise (#29605)
- Execute fixed-size rolling windows on streaming engine (#29602)
- Use the finite rule of equi joins to determine build side of semi/anti (#29617)
- Drop the byte limit on runtime filter build sides (#29618)
- Add eviction-free path to GroupedReduction, with lanes (#29613)
- Improve window functions (#29591)
- Push down
is_nan/is_not_nanfilters to Iceberg (#29579) - Compute shared group-by agg input subexpressions once (#29574)
- Emit files unordered for remote Parquet files when order not observed (#29581)
- Filter rows by the runtime join key range (#29576)
- Update few-group sums, means and counts in lanes (#29573)
- Loop over 4k arrays in group-by update (#29572)
- Decode parquet scans out of order when order is not observed (#29546)
- Don't always row-encode multi-column group-by (#29570)
- Push down
str.starts_withfilters to Iceberg (#29545) - Trace join runtime filters through unions (#29565)
- Fold SQL temporal literals into plain values (#29562)
- Don't rehash stored key on every TotalIndexMap probe (#29561)
- Fuse drops into all filters on streaming engine (#29550)
- Join probe on Arrays directly (#29558)
- Better build side selection in equi-join sampler (#29524)
- Large optimisation for
BytecodeParserrewrite/dispatch (#29548) - Parallelize expressions in in-memory map fallback (#29535)
- Push down
!=filters to Iceberg (#29455) - Improve predicate stats (#29495)
- Scale bloom filters from keys seen (#29477)
- Optimize simple like queries (#29459)
- Insert head slice below select(len() <cmp> n) to enable slice pushdown (#29209)
- Parallelize row encoding in the row_encode expression (#29434)
- Degrade outer join optimization (#29432)
- Optionally reorder semi join (#29443)
- Borrow the cached regex instead of cloning it per page (#29422)
- Lower SQL LIKE '%x%' to a literal substring match (#29431)
- Optimize
is_inexpression (#29425) - Pushdown bloom filters in probe side (#29423)
- Narrow the decimal rescale to 64 bits when the value fits (#29396)
- Use a thread-local copy of a pushed-down regex predicate (#29411)
- Sample anti-join build side (#29416)
- Improve performance importing from arrow (#29412)
- Sample table size to determine build side in semi joins (#29375)
- Optimize row-encoding (#29405)
- Order pushed parquet predicate columns by measured selectivity (#29397)
- Coerce float literals to decimal instead of casting the column (#29395)
- Fix OOM on TPCH SQL and fix fuzzing errors (#29389)
- Restrict correlated SQL aggregates to requested keys (#29383)
- Prepare Parquet scans for row-group splitting in Polars Cloud (#29295)
- Skip row groups by join runtime ranges without a statistics frame (#29370)
- Reduce copy in
scan_lines(#29310) - Increase HTTP read rate-limit default (#29363)
- Evaluate a pushed parquet predicate sequentially per conjuct (#29352)
- Disable system certificates for
CloudScheme::Httpsources (#29284) - Reduce rechunk in
sort_in_place(#29343) - Dynamic predicates for hash joins (#29312)
- Use HTTP suffix range for Parquet size and footer (#29308)
- Derive predicates from join conditions (#29304)
- Lower uncorrelated subqueries to semi joins and push semi/anti joins below inner joins (#29289)
- Make leaf name iterator unique (#29291)
- Don't clone the full frame per arm in when/then/otherwise (#29258)
- Push inner joins before outer joins and rewrite left-join-is-null to anti join (#29277)
- Use stats to decide cross join buffering side (#29270)
- Improve cache-removal and join-order cost estimates (#29263)
- Improve CSPE cost evaluation (#29250)
- Fuse group-by pre-select into node after partition (#29251)
- Inline hot small functions (#29244)
- Fix plan-time regressions in projection pushdown for wide frames (#28724)
- Rechunk before selecting in group-by pre-select (#29219)
- Reuse Iceberg data file sizes (#29063)
✨ Enhancements
- Python Polars 2.0 (#29723)
- Enable OOC by default with 80% of available RAM as threshold (#29741)
- Support window frames in SQL FIRST_VALUE and LAST_VALUE, and add NTH_VALUE and a LAG/LEAD default (#29740)
- Set default OOC disk budget to 64 GB (#29734)
- Add POLARS_OOMKILL_THRESHOLD_MB (#29735)
- Improve estimated DataFrame memory usage (#29712)
- Extend pipe_with_dtype for multiple expressions (#29676)
- Attribute physical nodes to IR nodes (#29522)
- Make approx_quantile sketch states mergeable across processes (#29660)
- Add unstable scan_lance (#29413)
- New Expression and function: pipe_with_dtype (#29547)
- Resolve
is_inand Map lookup coercion from dtypes alone (#29487) - Cast the needle of
is_inand Map lookups exactly or not at all (#29486) - Support
Mapin JSON and NDJSON reading and writing (#29478) - Add
scan_external_reader(#29232) - Add
struct.eval(#29454) - Improve
BytecodeParserUDF variable resolution (#29553) - Type exact SQL numeric literals as Decimal and truncate SQL
%/DIV (#29536) - Make Iceberg sink commits idempotent across re-executions (#29512)
- Allow deterministic expression plugins to opt into CSE/CSPE (#29428)
- Improve cloud object_store IO error types and messages (#29471)
- Add erf(c) function (#29501)
- Update
BytecodeParserfor Python 3.15 (#29491) - Expose the registered source as
scan_fn.io_source(#28897) - Support selectors in
joinkeys (#29233) - Honor Iceberg sort orders in native sinks (#29318)
- Add
APPROX_QUANTILEto the SQL frontend (#29288) - Add support for
approx_quantilein the streaming engine (#29237) - Read Parquet ENUM type as pl.String (#29331)
- Expose more
scan_iceberg/delta-related attributes in the visitor for cudf_polars (#29297) - More map operations (#29296)
- Binning functions (#28888)
- Support
collectandcollect_batchesusingRemoteEngine(#28914) - Support GROUP BY GROUPING SETS, ROLLUP, CUBE and GROUPING() (#29278)
- Fix tpch SQL issues (#29269)
- Fuse filters in (inner) join operation (#29218)
- Add in-memory support for approximate quantile (#29206)
- Disable casts from
StringtoTime(#29215)
🐞 Bug fixes
- Find SQL aggregates by their function name, and raise on nested aggregate calls (#29751)
- Run SQL window functions over grouped rows on the aggregated rows, and evaluate QUALIFY before the projection (#29746)
- Match categories added after
is_inprepares a string haystack (#29744) - Give tied rows the same running aggregate in SQL windows, and support more window frames (#29736)
- Fix scan_lance CI error (#29737)
- Rank NULLs and ties in SQL window functions, and add PERCENT_RANK, CUME_DIST and NTILE (#29729)
- Read from database without requiring SQLAlchemy
asyncioextra (#29713) - Don't slice the input of a select that has a window expression (#29725)
- Fix named SQL windows and raise for window shapes that gave wrong results (#29720)
- Preserve projections across nested caches (#29710)
- Allow the same equi-join key pair twice in a join condition (#29715)
- Defer fallible Hive pruning until after partition rewriting (#29700)
- Make dynamic window bounds data-independent (#29706)
- Panic on empty windows in dynamic group-by (#29704)
- Only reserve builder room for the rows a sort bucket receives (#29696)
- Fix handling of pruned input nodes during graph traversal (#29532)
- Minor fixes for
DataFrame"orient" behaviour (#29679) - Fix
list.containson multi-chunk sliced lists (#29635) - Make
is_sortedorder nested types assortdoes (#29683) - Fix compile errors after attributing physical nodes to IR nodes (#29687)
- Raise on ragged rows instead of silently dropping values (#29658)
- Exclude IPC dictionary batches from record batch stats row count (#29672)
- Handle reordered groups in sort_by (#29639)
- Don't recompile regex (#29597)
- Fix null and mixed-dtype edge cases in
is_in(#29593) - Errors in str.to_lowercase and str.to_titlecase (#29607)
- Make str.zfill consistent with pad for unicode (#29604)
- Empty group-by argmin/max on scalar (#29608)
- Support nested by in grouped
min_by/max_byin streaming engine (#29628) - Guard
is_nan/is_not_nanpushdown against Decimal columns in Iceberg (#29610) - Fix read_excel from_arrow error (#29492)
- Improve default S3 endpoint resolution (#29470)
- Don't let a SQL derived table alias replace a registered table (#29577)
- Write a single Avro header when a DataFrame has multiple chunks (#29569)
- Fix bug in
group_bypredicate pushdown (#29499) - Preserve literal state in list/arr eval under group_by (#29557)
- Don't evict when updating an existing
LRUCachekey (#29551) - Use SQL decimal result scales for
*and/(#29520) - Compute mixed-scale Decimal operations without a lossy common cast (#29519)
- Keep Decimal results within their declared precision and round Decimal to float casts correctly (#29539)
- Suggest
str.containsfor string containment inmap_elementsUDFs (#25472) - Keep a landed Iceberg snapshot's manifests when a commit response is lost (#29511)
- Allow deterministic expression plugins to opt into CSE/CSPE (#29428)
- Keep requested length in streaming negative slice (#29399)
- Match engine output dtype in planner for
Decimaland integer arithmetic (#29515) - Push down
!=filters to Iceberg (#29455) - Reject non-numeric input and preserve
Decimalscale inentropy(#29228) - Honor env var in Parquet/IPC pipeline budget (#29518)
- Stop casting data to a Decimal needle in
is_in, allowArray(Null)casts (#29485) - Clamp inflight budget from HTTP rate-limiter (#29497)
- Don't propagate RUSTFLAGS into dependencies in coverage (#29510)
- Read legacy Parquet LIST structures according to the spec (#29469)
- Address Python 3.15 compatibility issue with
_is_empty_method(#29489) - Insert head slice below select(len() <cmp> n) to enable slice pushdown (#29209)
- Fix loading negative decimals from iceberg statistics (#29452)
- Support
NaNcomparison pushdown to Iceberg (#29453) - Stop parse_version from splicing a pre-release suffix into the number (#29401)
- Produce correct parquet statistics for enum columns (#29364)
- Don't restart predicate pushdown if the predicate was not pushable (#29418)
- Add regression test for Decimal/integer arithmetic schema (#29120)
- Dynamic boundary schema mismatch (#29448)
- Avoid panic on Enum and Categorical literals in Parquet predicates (#29421)
- Fix SQL count star on CSV performance regression (#29427)
- Respect include_bounaries on in-mem empty dynamic (#29440)
- Fix over ordering for multiple cols (#29417)
- Enable dtype-i128 alongside dtype-decimal in polars-stream (#29410)
- Prevent bias in
'req_double'at the median by actually merging the sketches (#29326) - Fix OOM on TPCH SQL and fix fuzzing errors (#29389)
- Size row-index table statistics from the statistics frame (#29381)
- Preserve categories object in pyo3-polars (#29385)
- Strip the leading slash from Windows Delta table roots (#29386)
- Isolate expanded Python dataset scans (#29378)
- Make Iceberg bucket sort keys serializable (#29376)
- Apply "schema_overrides" in
read_databasefor Arrow-based drivers (#29273) - Remove duplicated word in rate-limit comment (#29372)
- Check whether Datetime is monotonically increasing dynamically (#29293)
- Keep input order of unmatched build rows in ordered streaming equi join (#29371)
- Unaliased constants in SQL SELECT with GROUP BY (#29367)
- Consistent
DateandDecimalmeans between the streaming and in-memory engines (#29359) - Respect string statistics for enum columns during parquet scanning (#29366)
- Support selectors properly in DataFrame n_unique (#29360)
- Decimal Parquet statistics for decimal/i128/f16 (#29350)
- Early check for converting integer map keys (#29348)
- Lowering for input-independent filter (#29340)
- Correlation of constant column returning non-NaN for larger inputs (#29319)
- Incorrect height in multi-input GroupBy (#29332)
- Fix panic in scan_iceberg for snapshot_id before a schema change (#28895)
- Resolve arithmetic
Structsupertypes per-field (#29261) - Fix comparison expression method comment (#29257)
- Don't panic on an empty or null quantile expression input (#29240)
- Detect list-valued quantile literals in approx_quantile
auto(#29236) - Various map issues (#29147)
- Rolling quantile should respect window trimming (#29200)
- Ensure ignored columns are excluded from dtype
Wildcardselector (#29220)
📖 Documentation
- Add engines and other missing items to API reference (#29493)
- Update migration guide with rc2 changes (#29415)
- Clarify that ambiguous parameter refers to DST transitions (#28873)
- Document gzip/zstd compression for read_csv and scan_csv (#29387)
- Add LazyFrameResolver to reference guide (#29379)
- Add user-guide for new enable_monitoring feature (#29358)
- Improve join_where engine tag (#29234)
🛠️ Other improvements
- Hold the category mapping in the categorical
is_inlookup (#29747) - Derive serde for moment states (#29742)
- Give each crate its own temporary column name counter (#29689)
- Apply
A-io-partitioninglabel for Hive-related issues and PRs (#29659) - Use explicit match arms in
LogSeriesoperations & improve error messages (#29575) - Improve reliability of a slightly flaky test (#29587)
- Split
has_joins_or_unionsintohas_joinsandhas_unions(#29338) - Remove scheduled cache cleaning (#29517)
- Add window placement to group by rolling (#29529)
- Add window placement to group by dynamic (#29527)
- Fix flaky test (#29530)
- Improve cloud object_store IO error types and messages (#29471)
- Fix iceberg-related test failures (#29533)
- Add IR structs for temporal group-by options (#29526)
- Workflow tweaks (#29502)
- Share the temporal group-by index space (#29498)
- Compute the first dynamic window start in one place (#29484)
- Fix anti/semi join test row with invalid result order expectation (#29467)
- Additional negative Decimal coverage for Iceberg (#29468)
- Add regression test for Decimal/integer arithmetic schema (#29120)
- Tidy the staged parquet predicate (#29408)
- Add dymanicpredicates to preferred build sides (#29335)
- Deprecate cut/qcut (#29329)
- Mark test as slow (#29336)
- Bump maturin (#29282)
- Bump object_store crate to 0.14.2 (#29317)
- Update rustls dependency to version
0.23.45(#29303) - Bump build deps used in ARM64 Windows release pipeline (#29280)
- Fix the stalling
test_fused_many_morsels_and_skewtest (#29290) - Use
SpillFrames inDataFrameSearchBuffers(#29207) - More obvious
MapChunkedstorage handling (#29248) - Centralize hoisted aggregate bookkeeping and clarify grouping predicates (#29287)
- Update analytics endpoint config (#29222)
- Add a blanket
lf.collect().schema == lf.collect_schema()check (#29224) - Add expand_paths parameter to toggle path expansion (#29210)
Thank you to all our contributors for making this release possible!
@ATL2001, @Aidavdw, @AlessandroKuz, @JakubValtar, @Kevin-Patyk, @MarcoGorelli, @Punisheroot, @Rodrigo-Palma, @SatvikMishra08, @TNieuwdorp, @aarushkandukoori, @abokhalill, @alexander-beedie, @atharva7905k, @ayushh0110, @borchero, @c-peters, @cBournhonesque, @carnarez, @dancsi, @dominikandreasseitz, @dsprenkels, @fsimkovic, @gautamvarmadatla, @hadrian-reppas, @jonasdedden, @kafka1991, @kdn36, @krithikashreeL, @lun3x, @madsbk, @matthewbayer, @mikhail5555, @mroeschke, @nameexhaustion, @orlp, @r-brink, @ritchie46, @shaneraphel, @vgvr0 and @wtn