Repository navigation
Allow the HLLAccumulator used by approx_distinct() to consume state produced by HllGroupsAccumulator - #25659
Allow the HLLAccumulator used by approx_distinct() to consume state produced by HllGroupsAccumulator#25659masonh22 wants to merge 5 commits into
HLLAccumulator used by approx_distinct() to consume state produced by HllGroupsAccumulator#25659Conversation
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## main #25659 +/- ##
==========================================
+ Coverage 82.45% 82.72% +0.26%
==========================================
Files 1140 1147 +7
Lines 436589 448186 +11597
Branches 436589 448186 +11597
==========================================
+ Hits 359996 370750 +10754
- Misses 54838 54906 +68
- Partials 21755 22530 +775 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
Thank you for this PR, this looks like an interesting use case. However, this isn't a supported pattern in DataFusion, and I'd suggest not introducing special cases for it, since they add complexity. WDYT about doing it downstream? It seems quite doable, since it only requires forking one aggregate function. |
Fair point, but I feel like this will inevitably lead to us forking/wrapping every aggregate function we use, which would not be ideal. I did file #25660 to discuss making this a supported pattern. I implemented a fuzz test I can share that is able to enforce this requirement. As an alternative to this, would you be open to a PR that applies the same optimization used by the groups accumulator to the regular accumulators? That would solve the problem I'm running into and potentially benefit the regular accumulators. What do you think? |
I'd push back on this a bit. This file treats the two accumulators' state as interchangeable:
So I read this PR as restoring the other half of a symmetry the file already documents, rather than adding a special case. I agree the two |
|
@avantgardnerio I factored out the deserialization into a |
avantgardnerio
left a comment
There was a problem hiding this comment.
Given the refactor, I think this is reasonable.
|
The issue is that the requirement for I understand that this is doable, and the PR looks quite clean now. What I’m still unclear about is why If there are use cases that require this stronger guarantee, I’d suggest:
|
I agree with you on this, and going forward this is something I want to do. I filed #25660 to do what you describe, and if you have feedback I would appreciate if you would share it there. I also have a prototype test harness for testing the invariant in #25710. The reason I made this separate PR is that I believed the fix here was small and wouldn't be too controversial. I thought that what I proposed in #25660 would be met with much more skepticism. |
|
I propose to merge this in the next 24 hours if there is no objections. I agree with @2010YOUY01 we should somehow document this and some further testing / integration to avoid breakage in the future. I am also wondering if we can somehow integrate this feature or document have it as example. |
2010YOUY01
left a comment
There was a problem hiding this comment.
I think the structural changes in this PR are a nice cleanup on their own, but the UTs also encode an additional contract that does not currently exist.
I'm leaning toward removing those UTs before merging, then documenting it in the public APIs first before enforcing it.
Regarding the proposed contract, I still don't fully understand the motivation. I read the issue and the discussion in this PR, but the main justification I found is roughly: "our use case requires Accumulator <-> GroupsAccumulator compatibility, so we want to enforce this as a new invariant."
Every new invariant adds some maintenance overhead, even if it is small. Ideally, that extra engineering effort should buy us clear functionality or performance. So I think it is important to first explain the end use case that requires this contract and agree that DataFusion should support it. Otherwise, anyone can propagate a downstream hack to upstream and turn it into a permanent invariant.
Below is just my guess at the actual use case, based on the UT added here:
// Quoted UT in this PR
let mut grouped = HllGroupsAccumulator::new();
grouped
.update_batch(std::slice::from_ref(&values), &group_indices, None, 3)
.unwrap();
let state = grouped.state(EmitTo::All).unwrap();
let mut direct = HLLAccumulator::new();
direct.update_batch(std::slice::from_ref(&values)).unwrap();
let mut merged = HLLAccumulator::new();
merged.merge_batch(&state).unwrap();
assert_eq!(
distinct_count(&mut merged),
distinct_count(&mut direct),
);Mapped to an end use case:
Suppose Q1 and Q2 have already been computed, and we cached their partial aggregate states:
Q1: select avg(v) from t1
Q2: select k, avg(v) from t2 group by k
Then we want to compute Q3 directly from Q1 and Q2's partial states, without recomputing the original inputs:
Q3: select avg(v) from (t1 union t2)
If this is the main use case, could Q1 also use GroupsAccumulator? It seems doable with some refactor, and we can avoid a new contract to all aggregate function implementation.
| direct.update_batch(std::slice::from_ref(&values)).unwrap(); | ||
|
|
||
| let mut merged = HLLAccumulator::new(); | ||
| merged.merge_batch(&state).unwrap(); |
There was a problem hiding this comment.
This seems not to be the expected usage of accumulator 🤔
There was a problem hiding this comment.
I removed these tests to avoid enforcing this undocumented usage
Not exactly this, but something conceptually similar (inspired by https://github.com/datafusion-contrib/datafusion-query-cache although ultimately quite different in implementation). We can cache and reuse partial aggregations to avoid costly scans. This requires adding a "cache key" to partial aggregates group by expressions. In the case of approx distinct with no group by (in the user query) we then end up with partial aggregation grouped by only the cache key and final aggregation with no group by expression which blows up.
This was our first thought as well but using the groups accumulator for approx distinct for a single group is a ~30% performance regression so its a non-starter. |
Which issue does this PR close?
Rationale for this change
At Coralogix, we have an optimization to share partial aggregations between queries where possible. Because of this, we will occasionally feed the state from an
AggregateUDFImpl's groups accumulator into it's regular, non-grouped accumulator.#22768 added an optimized groups accumulator for
approx_distinct()that uses a different state format for groups with few distinct values. Trying to feed this state into the non-grouped accumulator gives an internal error:"Impossibly got invalid binary array from states".What changes are included in this PR?
This change updates the non-grouped accumulators for
approx_distinct()so they can correctly parse the state from the grouped accumulator.What is the testing strategy for this PR?
The existing tests pass. I confirmed through testing no longer included in this PR that the groups and non-groups accumulators are compatible, but these were removed to avoid enforcing undocumented requirements.
Are there any user-facing changes?
No