Skip to content

Allow the HLLAccumulator used by approx_distinct() to consume state produced by HllGroupsAccumulator - #25659

Open
masonh22 wants to merge 5 commits into
apache:mainfrom
coralogix:hll-grouped-merge
Open

masonh22 wants to merge 5 commits into
apache:mainfrom
coralogix:hll-grouped-merge

Conversation

@masonh22

@masonh22 masonh22 commented Sep 23, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

  • Closes #.

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

@codecov-commenter

codecov-commenter commented Sep 23, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 68.88889% with 14 lines in your changes missing coverage. Please review.
✅ Project coverage is 82.72%. Comparing base (7570366) to head (be2076e).
⚠️ Report is 190 commits behind head on main.

Files with missing lines Patch % Lines
...afusion/functions-aggregate/src/approx_distinct.rs 68.88% 11 Missing and 3 partials ⚠️
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.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

@2010YOUY01

Copy link
Copy Markdown
Contributor

we will occasionally feed the state from an AggregateUDFImpl's groups accumulator into it's regular, non-grouped accumulator.

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.

@masonh22

Copy link
Copy Markdown
Contributor Author

we will occasionally feed the state from an AggregateUDFImpl's groups accumulator into it's regular, non-grouped accumulator.

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?

@avantgardnerio

avantgardnerio commented Sep 29, 2026 •

Copy link
Copy Markdown
Contributor

this isn't a supported pattern in DataFusion

I'd push back on this a bit. This file treats the two accumulators' state as interchangeable:

  1. GroupHll::merge_serialized is documented to accept state produced by the per-group Accumulator
  2. The dense sketch is documented as wire-compatible with the per-group Accumulator
  3. Both accumulators share a single state field, declared as hll_registers. Since perf: improve approx_distinct performance 100x when there are fewer distinct values with many groups #22768 (first released in 55.0.0), that column sometimes holds a list of hashes instead of registers, and the upgrade guide doesn't mention it.

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 merge_serialized implementations should be de-duplicated. With that change, I don't think this adds any maintenance burden.

@masonh22

Copy link
Copy Markdown
Contributor Author

@avantgardnerio I factored out the deserialization into a SerializedHll in abf16b7.

@avantgardnerio avantgardnerio left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Given the refactor, I think this is reasonable.

@2010YOUY01

Copy link
Copy Markdown
Contributor

The issue is that the requirement for GroupsAccumulator to have the same physical representation as Accumulator is currently neither documented nor tested, the only tests should be the UTs added in this PR. For simpler aggregate functions like avg(), I believe they already happen to use the same representation.

I understand that this is doable, and the PR looks quite clean now. What I’m still unclear about is why Accumulator <---> GroupsAccumulator compatibility is necessary in the first place. For example, would it be possible for your app to always use GroupsAccumulator?

If there are use cases that require this stronger guarantee, I’d suggest:

  • Explain from the end usage side why is it necessary
  • documenting it explicitly as part of the Accumulator / GroupsAccumulator contract;
  • adding a test harness to make the invariant easy to verify for aggregate implementations.

@masonh22

Copy link
Copy Markdown
Contributor Author

The issue is that the requirement for GroupsAccumulator to have the same physical representation as Accumulator is currently neither documented nor tested, the only tests should be the UTs added in this PR. For simpler aggregate functions like avg(), I believe they already happen to use the same representation.

I understand that this is doable, and the PR looks quite clean now. What I’m still unclear about is why Accumulator <---> GroupsAccumulator compatibility is necessary in the first place. For example, would it be possible for your app to always use GroupsAccumulator?

If there are use cases that require this stronger guarantee, I’d suggest:

  • Explain from the end usage side why is it necessary
  • documenting it explicitly as part of the Accumulator / GroupsAccumulator contract;
  • adding a test harness to make the invariant easy to verify for aggregate implementations.

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.

@Dandandan

Copy link
Copy Markdown
Contributor

I propose to merge this in the next 24 hours if there is no objections.
It seems a non-controversial improvement and makes this aggregation similar to other aggregation functions.

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 2010YOUY01 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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();

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This seems not to be the expected usage of accumulator 🤔

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I removed these tests to avoid enforcing this undocumented usage

@thinkharderdev

Copy link
Copy Markdown
Contributor
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.

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.

If this is the main use case, could Q1 also use GroupsAccumulator?

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.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

functions Changes to functions implementation

Projects

None yet

Development

Successfully merging this pull request may close these issues.

6 participants