Repository navigation
[EPIC] Improved Externalized / Spilling / Large than Memory Hash Aggregation #13123
Description
Activity
@2010YOUY01 says in #13090 (comment)
Really nice paper, we can implement the same benchmark and compare in the future 😄 They implemented a unified buffer pool for both table data cache and operator (like aggregation) intermediate results, to easily support spilling in various operators. I think they didn't mention any optimization specific to the spilling part of aggregation, and just use simple LRU policy in the buffer pool. Maybe there are some spilling and merging specific optimizations we can explore (all of memory-limited aggregate/SortMergeJoin/Sort can benefit from)
DF doesn't have a buffer pool in the traditional sense, and the way arrow-rs allocates memory directly from the system allocator makes it quite hard to implement. However, I think the fact that we have arrow-rs and the arrow IPC offers lots of opportunity.
Also, are you interested in improving DataFusion's external aggregation capabilities? I think it is a non trivial gap at the moment and would be great to improve (and I would be interested in helping do so).
if you are, I can start organizing the work into some tickets to see if we can get some others to check it out tooYes, I'm start to look at related components now. Perhaps we can start with making memory-limited SQL queries more stable (e.g. more tests, make sure TPCH-SF1000 is able to run on laptop correctly), and later optimize.
I think starting with stability and then optimizing is a great idea 💯
Note that one challenge of TPCH specifically is that it contains many joins and is largely focused on that, so in order to run TPCH-SF1000 we would also need to implement spilling joins
Another potential option would be to work on running clickbench with a very small memory (100MB)?
Or maybe we could figure out another large dataset 🤔
Note that one challenge of TPCH specifically is that it contains many joins and is largely focused on that, so in order to run TPCH-SF1000 we would also need to implement spilling joins
Maybe @comphead 's work to get SMJ working in #13111 will help this (e.g. we could always use SMJ for the large TPCH queries 🤔 )
Another potential option would be to work on running clickbench with a very small memory (100MB)?
This is a good idea, we should get clickbench work under memory constraints before TPCH
Reacted by Andrew Lamb- changed the title
[-][EPIC] Improved Externalized / Spilling / Out of core Hash Aggregation[/-][+][EPIC] Improved Externalized / Spilling / Large than Memory Hash Aggregation[/+]on Jan 10, 2025 Here is a PR to optimize the spill format:
- 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 As a followup to #19287 I'm thinking of working on spilling support for
GroupOrderingvalues other thanNone. Does anyone know if there are specific pitfalls to look out for when implementing that?Reacted by Andrew LambAs a followup to #19287 I'm thinking of working on spilling support for
GroupOrderingvalues other thanNone. Does anyone know if there are specific pitfalls to look out for when implementing that?Not that I know of
Not that I know of
I ended up implementing this in the linked PR. The only gotcha I came across is that in case of partial or full ordering, the output ordering is reported as being that ordering. The implication for spilling is that the spill files need to be sorted and merge accordingly.
Reacted by KAZUYUKI TANIMURA- added a commit that references this issue
on Dec 20, 2025
This is a collection of items to improve external (spilling) aggregation
Background
in the Solid State Age (DuckDB external aggregation paper))
DataFusion has supported memory limited / spilling hash aggregation since @kazuyukitanimura added it last year in #7400.
We can likely improve this feature and @2010YOUY01 is considering working on it
Tasks the solution you'd like