feat: optimize count star materialization - #25849
wudidapaopao wants to merge 5 commits into
Conversation
Codecov Report❌ Patch coverage is Additional details and impacted files@@ Coverage Diff @@
## main #25849 +/- ##
==========================================
- Coverage 82.66% 82.66% -0.01%
==========================================
Files 1147 1147
Lines 446357 446450 +93
Branches 446357 446450 +93
==========================================
+ Hits 368971 369044 +73
- Misses 54997 55000 +3
- Partials 22389 22406 +17 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
| r#"FROM (SELECT sum("total_revenue") AS "alias2", "#, | ||
| r#"date_part('year', "signup_date") AS "group_alias_0", "#, | ||
| r#""customer_id" AS "alias1" "# | ||
| r#"FROM (SELECT date_part('year', "signup_date") AS "group_alias_0", "#, |
There was a problem hiding this comment.
I'm intrigued as to why this has changed, since the ISSUE_23317_QUERY is performing a COUNT(DISTINCT ...).
There was a problem hiding this comment.
SingleDistinctToGroupBy first rewrites COUNT(DISTINCT customer_id) into an inner GROUP BY customer_id and an outer COUNT(alias1). Since customer_id is non-nullable, this PR simplifies the outer COUNT(alias1) to COUNT(), allowing OptimizeProjections to remove alias1 from the inner output.
mhilton
left a comment
There was a problem hiding this comment.
This has buried within it a set of interface enhancements to Accumulator and GroupsAccumulator that add support for functions with no inputs but still have information about the number of rows. I think those changes are important enough to be assessed outside of the context here. Could you make a separate PR for the interface changes so they can be considered aside from the use case please?
| } | ||
| } | ||
|
|
||
| fn retract_batch(&mut self, values: &[ArrayRef]) -> Result<()> { |
There was a problem hiding this comment.
It looks like this will be broken if it was ever used with the zero-column version.
There was a problem hiding this comment.
You're right that retract_batch would fail with nullary input. This PR intentionally does not support nullary aggregate windows: create_window_expr rejects empty arguments, and COUNT(*) used as a window aggregate still uses COUNT(1). So retract_batch cannot receive empty values on the normal path, and the PR does not introduce a regression.
That makes sense. I've split the interface changes into #25886. |
|
@wudidapaopao Thanks for this! I think this is a very useful optimization for some common query shapes. If I understand correctly, there are two distinct optimizations here:
I had Claude Code prototype a quick implementation of caching for const agg args -- based on some quick benchmarks, it matches the performance of this PR for the |
|
@wudidapaopao What do you think of the alternative approach to implementing this optimization that I suggested? |
…aterialization # Conflicts: # datafusion/physical-plan/src/aggregates/aggregate_stream.rs
|
@neilconway Thanks for the suggestion! I tested the aggregate argument caching approach and found that it achieves nearly identical performance. I’ve updated this PR to use the new approach and moved the |
Which issue does this PR close?
Rationale for this change
DataFusion represents
COUNT(*)asCOUNT(1)and expands the scalar1into an array for every input batch, although the accumulator only needs the number of rows.What changes are included in this PR?
DISTINCTCOUNT(*)/COUNT(1)to a nullaryCOUNT()while preserving output names.SingleDistinctToGroupByrewrites after name preservation adds an alias.COUNT()as SQLCOUNT(*).The nullability-based
COUNT(non_nullable_column) -> COUNT(1)rewrite was split into #25990.What is the testing strategy for this PR?
Covered by unit, SQL logic, DataFrame, unparser, and Substrait tests, as well as formatting, Clippy, and extended workspace test suites.
The release-nonlto benchmark used a 5,000,000-row in-memory
Int64table, one thread, one partition, batch size 8192, and ran this query 100 times:Median execution time: 1.263 ms → 0.947 ms (25.05% reduction).
Are there any user-facing changes?
Query results and output schemas are unchanged. Plans may display the internal aggregate as
count(), while SQL unparsing emitsCOUNT(*). New accumulator methods have default implementations, so existing UDAFs remain source-compatible.