wudidapaopao opened a new pull request, #25849:
URL: https://github.com/apache/datafusion/pull/25849

   ## Which issue does this PR close?
   
   - Closes #25847.
   
   ## Rationale for this change
   
   `COUNT(*)` and `COUNT` over provably non-null arguments count every input 
row, but DataFusion currently evaluates the argument into an array for each 
batch. This adds unnecessary CPU work and can keep otherwise unused columns in 
the scan.
   
   ## What changes are included in this PR?
   
   - Simplify safe, non-`DISTINCT` COUNT arguments to a nullary COUNT while 
preserving output names.
   - Pass row counts to aggregate execution paths that cannot derive them from 
argument arrays.
   - Preserve existing SQL output and optimizer behavior around DISTINCT 
aggregates.
   
   ## What is the testing strategy for this PR?
   
   Added unit coverage for nullary scalar/grouped COUNT, filters, direct state 
conversion, simplification safety, and interaction with 
`SingleDistinctToGroupBy`. Updated SQL logic, DataFrame, object-store, 
unparser, and Substrait expectations.
   
   Validated with:
   
   - `cargo fmt --all -- --check`
   - `cargo clippy --all-targets --all-features -- -D warnings`
   - The extended workspace test suite from the contributor guide
   
   Performance was measured with a release-nonlto microbenchmark over 1,000,000 
in-memory `Int64` rows, one partition, batch size 8192, and:
   
   ```sql
   SELECT COUNT(*) FROM t WHERE v < 900
   ```
   
   The filter returns 90% of the rows. After 50 warmup iterations and 500 
measured iterations per run, across 10 alternating baseline/branch runs:
   
   - baseline (`COUNT(1)`): 224.5 µs median
   - nullary COUNT: 166.2 µs median
   - improvement: 25.97%
   
   Typical 10-million-row ClickBench queries showed no measurable change 
because scan and grouping costs dominate this argument-materialization cost.
   
   ## Are there any user-facing changes?
   
   Query results and output schemas are unchanged. Logical and physical plans 
may display the internal aggregate as `count()`, while SQL unparsing emits 
`COUNT(*)`. The new accumulator methods have default implementations, so 
existing UDAFs remain source-compatible.
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to