note

The cluster tax

Polars and DuckDB spent this year moving the ceiling on what one machine can do. The Polars 1.41 release has my favourite optimisation detail of the year in it.

· polars / duckdb / spark

The thing nobody warns you about with Spark is that the compute is the cheap part.

You wanted to process some Parquet. What you got was a scheduler, a shuffle service, executors that die for reasons you’ll be reading about in a stack trace at half past eleven, serialisation costs you can’t see, a memory model with three different kinds of memory, and a web UI you will come to know intimately. The job runs. It’s just that “the job runs” now involves a distributed system you’re on call for.

The threshold where you have to take that deal moved a long way this year.

My favourite detail of the year

Polars 1.41 replaced their Parquet metadata parsing with a hand-written Thrift decoder.

The bottleneck wasn’t reading the data. It was reading the description of the data — the footer, the schema, the column statistics you read to work out which parts of the file you can skip.

Anyone who’s watched a Spark job burn minutes listing files and cracking open Parquet footers before doing a byte of real work knows that shape. It’s the “why has nothing happened yet” phase. Polars looked at where the time actually went, found a generated Thrift parser, and wrote one by hand.

Things pandas never did for you

pandas has no query optimiser. It never did. Every operation happens right now, in the order you wrote it, whether or not that order made any sense. Filter after the join instead of before it and congratulations, you materialised the whole join. The optimisation pass was you, squinting at your own code, moving lines around.

Polars 1.42 taught its optimiser to throw out contradictory filter predicates — ask for rows where x < 10 and x > 20 and it reads nothing at all. 1.41 added nested common subplan elimination, so a subquery you referenced four times gets computed once. 1.43 sped up joins on hive-partitioned data.

Together that’s what makes the lazy API worth the weirdness. The first time the optimiser deletes half your pipeline because you never used those columns, it stops feeling like ceremony.

DuckDB deleted a stall

DuckDB 1.5.0 made checkpointing non-blocking and picked up about 17% on TPC-H throughput.

Checkpointing used to be a thing that stopped — a periodic window where the database does its own bookkeeping instead of answering you. Now reads and writes carry on through it.

Half of what gets marketed as performance work is really this: someone found a stop-the-world pause and got rid of it. Those pauses are what make a system feel twitchy under load, and twitchy systems are how you end up buying a cluster because nobody could explain the p99.

The obvious objection

Polars is building the distributed thing too. Polars Cloud 0.9.0 landed in July with distributed expressions, an Iceberg sink, Kubernetes deployment. Credit to them, they published a straight single-node versus distributed comparison in mid-July.

My complaint is with the reflex. Somebody says “this won’t fit on one box” early, usually before anyone has measured, and from that moment it’s your problem forever. A cluster sized against a 2019 assumption doesn’t tell you it’s now processing data that would fit in RAM on a laptop. It just keeps billing you and paging you.

Go and find out what your actual single-machine ceiling is on this year’s engines. If it turns out you’re under it, that’s a whole lot of spark complexity you get to not have.

← all notes