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]