Skip to content
chornous.dev

Type two or more letters. Esc closes.

Index of sheets
Theme

Sheet 03 · Writing

All notes

Polars 2.0: lazy queries run on the streaming engine, and row order is no longer a given

The Polars team tagged Python Polars 2.0.0 on October 6, 2026. LazyFrame.collect() now defaults to the streaming engine with out-of-core spilling at 80% of RAM, read_csv dispatches to scan_csv, and mixed signed and UInt64 math returns Int128.

Pl. 62 · pipeline drawing generated from the slug “polars-2-0-streaming-default”

The Polars team tagged Python Polars 2.0.0 on October 6, 2026, and the announcement reached the Hacker News front page the next day with 423 points. The release notes list six breaking changes. The upgrade guide in the Polars repository lists several dozen, and the first one changes how every lazy query you own executes.

Streaming is the default

LazyFrame.collect() and collect_async() with engine="auto" now run on the streaming engine. Before 2.0, auto meant the in-memory engine. Eager DataFrame operations still run in memory, and sink_* and explain() behave as before.

Two defaults came with it. Polars now spills to disk once a query passes 80% of available RAM, and it caps the spill at a 64 GB disk budget. A new POLARS_OOMKILL_THRESHOLD_MB setting sits beside them. The release also adds a streaming, out-of-core sort, so sort no longer forces a large query back into memory.

The cost shows up in row order. The streaming engine doesn't guarantee order for group_by, joins or unpivot. If a test compares a collected frame to a fixture row by row, expect it to fail at random.

import polars as pl

lf = pl.scan_parquet("events/*.parquet")

# 2.0: streaming engine, output order not guaranteed
out = lf.group_by("user_id").agg(pl.len()).collect()

# Keep order where your code depends on it
out = lf.group_by("user_id", maintain_order=True).agg(pl.len()).collect()

# Or opt out of streaming for one query
out = lf.group_by("user_id").agg(pl.len()).collect(engine="in-memory")

You can opt out for a whole process with pl.Config.set_engine_affinity("in-memory") or the POLARS_ENGINE_AFFINITY=in-memory environment variable. Use that as a bridge while you fix tests, and plan to remove it.

Changes that alter results without an error

Polars 2.0 changes that return different data
Operation2.0 behavior
signed int combined with UInt64Int128, where 1.x returned Float64
explode() on an empty listzero rows by default
Parquet ENUM columnread as String
Arrow, Parquet and Iceberg map columnsnew Map dtype
SQL numeric literalsexact Decimal, with truncating % and DIV
to_struct() on null outer valueskeeps the outer null

The explode() change hits joins downstream. If you explode a list column and expect one row per parent, parents with empty lists now drop out of the frame.

The SQL changes move result types. Window functions over grouped rows now run on the aggregated rows, QUALIFY runs before the projection, and Polars resolves SQL queries at collect() time instead of up front.

Calls that raise now

Polars 2.0 turns many silent coercions into errors. The upgrade guide lists these among others:

  • casts from integer to Enum or Categorical, and from string to temporal types
  • &, | and ^ between booleans and integers
  • comparisons between naive and time-zone-aware datetimes
  • shift(None) and passing a Python list to search_sorted()
  • pl.concat(how="horizontal") on frames of different heights

For the last one, how="horizontal_extend" restores the old padding. is_in() coercion is strict now, and a Decimal needle against float data fails.

I/O moved onto the lazy path

pl.read_csv dispatches to scan_csv and pl.read_ipc dispatches to scan_ipc. The IPC reader dropped memory_map and rechunk. CSV schemas match columns by name instead of position, and headerless files name columns from column_0. Polars no longer rewinds file-like objects before it reads them, so it starts from wherever your code left the cursor.

The guide also removes LazyFrame.profile() and the DataFrame Interchange Protocol, and it marks read_avro() and write_avro() unstable. cut() and qcut() are deprecated in favor of bin_intervals(), bin_quantiles() and bin_ranks().

New features

The release adds struct.eval, scan_external_reader, an unstable scan_lance with filter pushdown, and GROUPING SETS, ROLLUP and CUBE in SQL. The performance list runs long and focuses on the streaming engine: faster null tracking in streaming group-by aggregations, prefetched hash table lookups, and cached Iceberg manifests across scans. The release notes publish no benchmark figures, so measure your own workloads before you claim a speedup.

This week

  1. Pin polars<2 in production until your tests pass on 2.0.
  2. On a branch, install 2.0 and run your suite with no config changes. Order-dependent failures point at queries that need maintain_order=True or an explicit sort.
  3. Search for explode(, read_csv(, concat( with how="horizontal", and UInt64 columns.
  4. Diff a sample of outputs between 1.x and 2.0 for queries that feed reports, since dtype changes like Int128 don't raise.

If a pipeline runs on a machine with little disk, set the out-of-core disk budget before you deploy. The 64 GB default assumes room you may not have.

Volodymyr Chornous