+− THE DAILY DIFFdev & AI news
NEEDS REVIEW

Polars 2.0: why joins no longer keep their row order

Polars 2.0 runs collect on a streaming engine by default.

Polars 2.0 runs collect on a streaming engine by default. For join, group_by and unpivot, the release no longer guarantees input row order unless you set maintain_order or sort explicitly. In our run on polars 2.0.0, a 2,000,000-row inner join returned rows out of input order by default and kept the order with maintain_order set. The release also turns on out-of-core spilling and makes the benchmark claims the vendor's own.

Read the written edition (English) ↗

What this video covers

  • Calling collect() on a lazy query now uses the streaming engine by default, which does not guarantee row order for join, group_by and unpivot.
  • The migration guide says the order is no longer the pre-2.0 order and is not guaranteed; it recommends an explicit sort if you rely on it.
  • In our run, a default inner join returned input order False; maintain_order="left_right" on the join returned True. One machine, one run, one join shape.
  • Out-of-core spilling is on by default: it starts at about 80 % of RAM, with a 64 GB default disk budget. Out-of-core join and group_by are on the roadmap.
  • collect_schema() catches a missing column before any data is read, but a cast that fails on one value only fails at collect().
  • The TPC-H and TPC-DS comparison is the vendor's own run on derived data, and its footnote says the results are not comparable to official benchmarks.

Transcript

Why did the default order change?

0:00 You think a join hands your rows back in the order you wrote them. Polars 2.0 stopped promising that, by default, on purpose. In this video: why did the default change? What did it buy? And which of your scripts breaks first? This is The Daily Diff, under the hood. Polars is an open-source DataFrame library: a table engine you call from Python,

0:20 with its core written in Rust. It is free under the MIT license, and you install it with pip install polars. One detail first. The switch that brings the old order back is an argument to the join itself. I'll come back to it at the end. Since 2.0, collect runs the streaming engine by default. Streaming splits the query into chunks that run in parallel, and the chunks finish when they finish.

0:42 For join, group by and unpivot, the release does not promise an order. Who gets hit first? Whoever compares output to a saved file, row by row. That test can pass on your laptop and fail on a different machine. Next, the migration guide.

What does the migration guide promise?

0:57 It says the exact row order shown above is not guaranteed. So if you rely on the order, sort explicitly. Out-of-core is on by default too. It starts spilling to disk at about eighty percent of RAM, with a sixty-four gigabyte budget. Sort, window functions and many expressions can spill now. Joins and group-bys are coming.

Is the benchmark a fair fight?

1:17 Polars says it beats DataFusion and DuckDB on the TPC-H and TPC-DS benchmarks. The numbers come from their own runs, on derived data, and their footnote says they are not comparable to official results. Their own numbers say Polars gets about three point eight times faster from sixteen cores to one hundred ninety-two, on TPC-H. They also say the big machine's extra threads slow down small queries.

What do we see when we run it?

1:41 We ran the same join twice on Polars 2.0.0. The default run shuffles the rows. With the flag, the input order survives. Find the joins that feed a test. Any join without a sort after it is a row order decision you never made. Polars 2.0 is stricter about types.

Does stricter mean earlier?

1:58 The schema check catches a missing column before any data is read. It cannot catch a cast that fails on one row, because that error waits for the data. Now the loop I opened.

How do I get the old order back?

2:08 The switch lives on the join. Set maintain order on the join, or add a sort after it. The guide suggests the sort. The join flag kept the order in our run, with the left side's order kept as left, right. Verdict, under the hood: needs review.

What is the verdict?

2:22 I'd review every join before upgrading. Got a question about this? Put it in the comments. And that's the diff for today. I'm Niko from Axrisi. Merge responsibly.

Sources

  1. Release of Polars 2.0Polars (pola.rs)
  2. Polars 2.0 upgrade guidePolars documentation
  3. Polars homepagePolars (pola.rs)
  4. polars-2.0-benchmark repositoryPolars on GitHub
  5. Release of Polars 2.0 (Hacker News discussion)Hacker News

Related videos