← Back to context

Comment by f311a

8 hours ago

There are some benchmarks for the previous version https://benchmark.clickhouse.com/#system=+ti%20rud|Dusa,s|PD...

Doing the benchmarks for 2.0 on the large AWS metal machines at small data sizes (SF=10) really opened my eyes that we have some low-hanging fruit in Polars when it comes to optimizing our constant overhead for smaller queries.

For example our join currently does a full partition into T partitions, for each of the T threads. Overall we create T^2 partitions, which on a 192-core machine is non-trivial. Great if you have a ton of data to feed that with, but if you 'only' have a few dozen million rows it becomes rather small. This is the primary reason we saw in the benchmarks that Polars pinned to 32 threads beats 192 thread Polars at SF=10.

I'll be working on improving that soon. I expect that to have a big impact on SF=10, and a decent impact on ClickBench, which sits between SF=10 and SF=100 in terms of rows.