← Back to context

Comment by niltecedu

8 hours ago

A bit surprised about the datafusion results from the post, I have tried it time and time again, but datafusion has always been the leading/trading blowers with polars for our workfloads with duckdb being vastly slower.

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.