Repository navigation
[EPIC] Improved performance in H2O.ai benchmarks #13548
Description
Activity
For posterity, here is a link to the discord chat: https://discord.com/channels/885562378132000778/1309883046886903870/1309887744595595324
Would like to note that the DataFusion performance really starts to lag when the dataset size grows.
Take a look at this query:
select id2, id4, power(corr(v1, v2), 2) as r2 from x group by id2, id4.When the dataset is 10 million rows, then Polars takes 3 seconds and DataFusion takes 3.6 seconds, so pretty similar.
When the dataset is 100 million rows, then Polars takes 126 seconds and DataFusion takes 2,100 seconds.
Reacted by Raz LuvatonReacted by Yongting YouWhen the dataset is 100 million rows, then Polars takes 126 seconds and DataFusion takes 2,100 seconds.
What version are you working with?
@Rachelint has some ideas of how to improve this:
Hm this seems something quadratic in nature?
When the dataset is 100 million rows, then Polars takes 126 seconds and DataFusion takes 2,100 seconds.
What version are you working with?
@Rachelint has some ideas of how to improve this:
Does it fully explain the dramatic difference? @MrPowers how do you generate the 10M vs 100M rows?
I would also expect this to help (but it was merged and depends on when it is merged)
- Skipping partial aggregation when it is not helping for high cardinality aggregates #11627
I think the right thing to do is to get the query / dataset and profile it
- Skipping partial aggregation when it is not helping for high cardinality aggregates #11627
@Dandandan - thanks to the great work by @SemyonSinchenko, it's easy to generate these datasets with falsa.
Here's the command to generate the 10 million row dataset:
falsa groupby --path-prefix=~/data --size SMALL --data-format PARQUET. Just useMEDIUMto generate the 100 million row dataset.@Dandandan - thanks to the great work by @SemyonSinchenko, it's easy to generate these datasets with falsa.
Here's the command to generate the 10 million row dataset:
falsa groupby --path-prefix=~/data --size SMALL --data-format PARQUET. Just useMEDIUMto generate the 100 million row dataset.Thanks, I will profile and see what happen about the so long time cost in datafusion.
When the dataset is 100 million rows, then Polars takes 126 seconds and DataFusion takes 2,100 seconds.
What version are you working with?
@Rachelint has some ideas of how to improve this:
* [Sketch for aggregation intermediate results blocked management #11943](https://github.com/apache/datafusion/pull/11943) * [Manage group values and states by blocks in aggregation #11931](https://github.com/apache/datafusion/issues/11931)🤔 I guess it may be caused by the similar reason of what we encountered during benchmarking in #11827
When the dataset is 100 million rows, then Polars takes 126 seconds and DataFusion takes 2,100 seconds.
What version are you working with?
@Rachelint has some ideas of how to improve this:* [Sketch for aggregation intermediate results blocked management #11943](https://github.com/apache/datafusion/pull/11943) * [Manage group values and states by blocks in aggregation #11931](https://github.com/apache/datafusion/issues/11931)🤔 I guess it may be caused by the similar reason of what we encountered during benchmarking in #11827
Specifically that
powerandcorrneed to supportconvert_to_state?When the dataset is 100 million rows, then Polars takes 126 seconds and DataFusion takes 2,100 seconds.
What version are you working with?
@Rachelint has some ideas of how to improve this:* [Sketch for aggregation intermediate results blocked management #11943](https://github.com/apache/datafusion/pull/11943) * [Manage group values and states by blocks in aggregation #11931](https://github.com/apache/datafusion/issues/11931)🤔 I guess it may be caused by the similar reason of what we encountered during benchmarking in #11827
Specifically that
powerandcorrneed to supportconvert_to_state?I am not sure, but I think it maybe really related to
GroupAccumulatorAdapteras #11827?
I am running and profiling it to find the answer.Reacted by Andrew LambI rerun the H2O Q9 with
GroupsAccumulatorforcorr(). See #13581
h2o dataset is inparquetformatResult ---- main, h2o_10m: 0.8s main, h2o_100m: 12s pr, h2o_10m: 0.2s pr, h2o_100m: 4sI didn't reproduce the drastic slowdown in
mainbranch🤔When the dataset is 10 million rows, then Polars takes 3 seconds and DataFusion takes 3.6 seconds, so pretty similar.
When the dataset is 100 million rows, then Polars takes 126 seconds and DataFusion takes 2,100 seconds.
Reacted by Daniël Heres, Andrew Lamb and kamilleThe h2o benchmarks are run on a Intel(R) Xeon(R) Platinum 8375C CPU @ 2.90GHz machine with 128 cores and 250 GB of RAM.
DataFusion groupby queries perform well on the 100 million row dataset (~5GB of data in a CSV file):
Some don't run with the 1 billion row dataset (~50GB of data in an uncompressed CSV file):
I am using a M3 Macbook with 16 GB of RAM. How much RAM does your machine have? Perhaps DataFusion only struggles with query 9 when the machine doesn't have lots of extra RAM.
Reacted by Tai Le ManhGroupsAccumulatorfor median/corr reduces memory usage (should be by quite a bit).Looking at the benchmark results, I think query 8 is worth analyzing / optimizing as well:
#13548I am using a M3 Macbook with 16 GB of RAM. How much RAM does your machine have? Perhaps DataFusion only struggles with query 9 when the machine doesn't have lots of extra RAM.
This explains 👍🏼 I ran the benchmark on a macbook with 48G of ram.
It is likely Q9 requires > 16G RAM, and OS memory swapping caused the performance regression.We should also take a look at how much memory does DataFusion consume for those queries, comparing to other systems. Thanks for the report.
Reacted by Daniël Heres and kamilleI am using a M3 Macbook with 16 GB of RAM. How much RAM does your machine have? Perhaps DataFusion only struggles with query 9 when the machine doesn't have lots of extra RAM.
This explains 👍🏼 I ran the benchmark on a macbook with 48G of ram. It is likely Q9 requires > 16G RAM, and OS memory swapping caused the performance regression.
We should also take a look at how much memory does DataFusion consume for those queries, comparing to other systems. Thanks for the report.
Yes, I run it today, and my machine has only 16GB memory too... and I found the query very very slow due to swapping, too...
I think making DataFusion work better in lower memory situations would certainly be nice
Reacted by Tai Le Manh and kamille- changed the title
[-][EPIC] Improved aggregate function performance[/-][+][EPIC] Improved aggregate function performance (faster H20.ai benchmarks)[/+]on Dec 13, 2024 Update here is that @2010YOUY01 has made
corrmuch better:I also dug up and connected the task to add the H20.ai queries to bench.sh. Check out
- changed the title
[-][EPIC] Improved aggregate function performance (faster H20.ai benchmarks)[/-][+][EPIC] Improved performance (driven by H20.ai benchmarks)[/+]on Dec 16, 2024 - changed the title
[-][EPIC] Improved performance (driven by H20.ai benchmarks)[/-][+][EPIC] Improved performance in H20.ai benchmarks[/+]on Dec 16, 2024 Very nice! The current slowest contains a window function:
SELECT id6, largest2_v3 FROM (SELECT id6, v3 AS largest2_v3, ROW_NUMBER() OVER (PARTITION BY id6 ORDER BY v3 DESC) AS order_v3 FROM x WHERE v3 IS NOT NULL) sub_query WHERE order_v3 <= 2;
- changed the title
[-][EPIC] Improved performance in H20.ai benchmarks[/-][+][EPIC] Improved performance in H2O.ai benchmarks[/+]on May 16, 2025 @MrPowers Could you use
.collect(engine="streaming")for Polars?- addedEPICA larger project, actively underway, with sub tasksA larger project, actively underway, with sub tasksPROPOSAL EPICA proposal being discussed that is not yet fully underwayA proposal being discussed that is not yet fully underwayand removedEPICA larger project, actively underway, with sub tasksA larger project, actively underway, with sub tasks
on Aug 12, 2025

Is your feature request related to a problem or challenge?
The basic aggregate functions like
COUNTandSUMin DataFusion are very fast (see Apache DataFusion is now the fastest single node engine for querying Apache Parquet files)However, many of the other aggregate functions are not particularly fast, and this shows up specifically on some of the H20 benchmarks
We saw this in the results in the 2024 DataFusion SIGMOD paper

(BTW we have made median faster)
@MrPowers has also observed similar results on discord (link):
See his version of the benchmarks here
https://github.com/MrPowers/mrpowers-benchmarks
Testing
dfbench#7209Functions
corrfunction #13549medianfunction #13550Other improvements
Describe the solution you'd like
DataFusion has two APIs ways to implement Aggregate functions like
SUMandCOUNTAccumulator(api docs)GroupsAccumulator(api docs)The basic aggregates are implemented using
GroupsAccumulatorand are part of DataFusions performanceThis ticket tracks the effort to improve the performance of these for these "more advanced" aggregate functions, likely by implementing
GroupsAccumulatorDescribe alternatives you've considered
For each function listed above, ideally we would:
GroupsAccumulatorfor the relevant aggregate function in a second PR (along with tests for correctness). We would use the benchmark to verify the performanceHere is a pretty good example of how @eejbyfeldt did this for
STDDEV:Additional context
No response