mkleen opened a new issue, #25936:
URL: https://github.com/apache/datafusion/issues/25936

   ### Describe the bug
   
   `SELECT g, count(DISTINCT x) FROM t GROUP BY g` can run in two ways:
   
   1. **Rewrite.** The `SingleDistinctToGroupBy` rule first builds a list of 
every different `(g, x)` pair, then counts per group.
   2. **Count directly.** One pass over the rows. For numbers,there is a fast 
way to do this.
   
   For numbers, counting directly is about 20% faster and needs the same 
memory. But currently DataFusion always uses the rewrite.
   
   ### To Reproduce
   
   Stock `datafusion-cli`:
   
   ```sql
   SET datafusion.execution.target_partitions = 1;
   
   CREATE TABLE t AS
   SELECT v % 2000 AS g, (v * 48271) % 999983 AS x, v % 1000 AS y
   FROM (SELECT unnest(generate_series(0, 3999999)) AS v);
   
   SELECT g, count(DISTINCT x) FROM t GROUP BY g;
   ```
   
   ```
   datafusion-cli -m <limit> -d 0 --mem-pool-type greedy -f <file>
   ```
   
   The table has 4 million rows, 2,000 groups, and 4 million different `(g, x)` 
pairs (same as #24929).
   
   | rewrite | memory limit needed | time |
   | --- | --- | --- |
   | yes | 192M | 0.152 s |
   | no | 192M | 0.124 s |
   
   "yes" is stock `datafusion-cli`. "no" is a test build with the rule switched 
off.
   
   `EXPLAIN` with the rule on. The inner aggregate builds one row per `(g, x)` 
pair (4 million), and the outer one counts them per `g`:
   
   ```
   > EXPLAIN FORMAT indent SELECT g, count(DISTINCT x) FROM t GROUP BY g;
   logical_plan
   Projection: t.g, count(alias1) AS count(DISTINCT t.x)
     Aggregate: groupBy=[[t.g]], aggr=[[count(alias1)]]
       Aggregate: groupBy=[[t.g, t.x AS alias1]], aggr=[[]]
         TableScan: t projection=[g, x]
   physical_plan
   ProjectionExec: expr=[g@0 as g, count(alias1)@1 as count(DISTINCT t.x)]
     AggregateExec: mode=Single, gby=[g@0 as g], aggr=[count(alias1)]
       AggregateExec: mode=Single, gby=[g@0 as g, x@1 as alias1], aggr=[]
         DataSourceExec: partitions=1, partition_sizes=[489]
   ```
   
   With the rule off there is one aggregate, which counts directly:
   
   ```
   > EXPLAIN FORMAT indent SELECT g, count(DISTINCT x) FROM t GROUP BY g;
   logical_plan
   Aggregate: groupBy=[[t.g]], aggr=[[count(DISTINCT t.x)]]
     TableScan: t projection=[g, x]
   physical_plan
   AggregateExec: mode=Single, gby=[g@0 as g], aggr=[count(DISTINCT t.x)]
     DataSourceExec: partitions=1, partition_sizes=[489]
   ```
   
   For text (`CAST(... AS VARCHAR)`), the rewrite is more than twice as fast 
and needs less memory (224M vs 384M), so the rewrite is the right choice there.
   
   ### Expected behavior
   
   For number columns, `count(DISTINCT x)` should count directly instead of 
using the rewrite. 
   
   ### Additional context
   
   - #24929


-- 
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