mkleen opened a new issue, #25932:
URL: https://github.com/apache/datafusion/issues/25932
### Describe the bug
`SELECT g, sum(DISTINCT x) FROM t GROUP BY g` adds up the unique values of
`x` in each group. DataFusion can answer it in two ways:
1. **The rewrite.** An optimizer rule (`SingleDistinctToGroupBy`) splits the
query in two. First it builds a list of every unique `(group, value)` pair,
then it adds those up. That list can be huge. In the example below it holds 4
million rows for only 2,000 groups.
2. **The direct way.** DataFusion keeps one running total per group and adds
each unique value as it sees it.
Which way you get depends on what else is in the query, not on which way is
cheaper:
- `sum(DISTINCT x)` on its own, or next to a plain `sum`, `min` or `max`,
gets the rewrite.
- Adding `count(*)` switches the rewrite off.
- Next to `avg`, the rewrite never happens.
The rewrite is faster but uses much more memory. On the example below it
needs 2 to 2.3 times as much memory as the direct way, and runs about 1.5 to
1.7 times faster. So adding a `count(*)`, which makes the query do more work,
cuts its memory use in half.
### To Reproduce
Everything below runs in a normal `datafusion-cli`, with no changes, no
special settings and no custom build. The one exception is the last section,
which says so.
Create a test table:
```sql
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);
```
It has 4 million rows, 2,000 groups and about 1 million different values of
`x`. Every row is a different `(g, x)` pair.
Put each query in its own file and run it with a memory limit:
```
datafusion-cli -m <limit> -d 0 --mem-pool-type greedy -f <file>
```
Each file starts with `SET datafusion.execution.target_partitions = 1;`,
then the `CREATE TABLE` above, then one query. Using one partition makes the
numbers the same on any machine. `-d 0` stops DataFusion from writing to disk
when it runs out of memory, so the limit shows how much memory the query really
needs.
#### Memory limit
The smallest memory limit at which each query finishes:
| query | rewritten? | finishes at | fails at |
| --- | --- | --- | --- |
| `SELECT g, sum(DISTINCT x) FROM t GROUP BY g` | yes | 192M | 160M |
| `SELECT g, sum(y), sum(DISTINCT x) FROM t GROUP BY g` | yes | 224M | 192M |
| `SELECT g, count(*), sum(y), sum(DISTINCT x) FROM t GROUP BY g` | no | 96M
| 64M |
The last query does more work than the second one, yet needs less than half
the memory.
`EXPLAIN` shows the difference. With the rewrite there are two aggregate
steps, and the lower one groups by both `g` and `x`:
```
> EXPLAIN FORMAT indent SELECT g, sum(y), sum(DISTINCT x) FROM t GROUP BY g;
logical_plan
Projection: t.g, sum(alias2) AS sum(t.y), sum(alias1) AS sum(DISTINCT t.x)
Aggregate: groupBy=[[t.g]], aggr=[[sum(alias2), sum(alias1)]]
Aggregate: groupBy=[[t.g, t.x AS alias1]], aggr=[[sum(t.y) AS alias2]]
TableScan: t projection=[g, x, y]
physical_plan
ProjectionExec: expr=[g@0 as g, sum(alias2)@1 as sum(t.y), sum(alias1)@2 as
sum(DISTINCT t.x)]
AggregateExec: mode=Single, gby=[g@0 as g], aggr=[sum(alias2), sum(alias1)]
AggregateExec: mode=Single, gby=[g@0 as g, x@1 as alias1],
aggr=[sum(t.y) as alias2]
DataSourceExec: partitions=1, partition_sizes=[489]
```
With `count(*)` added there is only one step, grouped by `g`:
```
> EXPLAIN FORMAT indent SELECT g, count(*), sum(y), sum(DISTINCT x) FROM t
GROUP BY g;
logical_plan
Projection: t.g, count(Int64(1)) AS count(*), sum(t.y), sum(DISTINCT t.x)
Aggregate: groupBy=[[t.g]], aggr=[[count(Int64(1)), sum(t.y), sum(DISTINCT
t.x)]]
TableScan: t projection=[g, x, y]
physical_plan
ProjectionExec: expr=[g@0 as g, count(Int64(1))@1 as count(*), sum(t.y)@2 as
sum(t.y), sum(DISTINCT t.x)@3 as sum(DISTINCT t.x)]
AggregateExec: mode=Single, gby=[g@0 as g], aggr=[count(Int64(1)),
sum(t.y), sum(DISTINCT t.x)]
DataSourceExec: partitions=1, partition_sizes=[489]
```
#### Memory used by the process
The same queries again, with no memory limit, measured with `/usr/bin/time
-l` on macOS (`/usr/bin/time -v` on Linux). Each was run 3 times.
| what runs | most memory used |
| --- | --- |
| create the table, then `SELECT 1` | 302 MiB (every run) |
| create the table, then `SELECT g, sum(DISTINCT x) FROM t GROUP BY g` | 537
to 568 MiB |
| create the table, then `SELECT g, sum(y), sum(DISTINCT x) FROM t GROUP BY
g` | 568 to 599 MiB |
| create the table, then `SELECT g, count(*), sum(y), sum(DISTINCT x) FROM t
GROUP BY g` | 306 MiB (every run) |
The table alone takes 302 MiB. On top of that, the rewritten queries use
about 235 to 266 MiB more, while the `count(*)` query uses only about 4 MiB
more.
The two measurements do not fully agree on the `count(*)` query. It needs a
96M limit to finish, but the process only grows by about 4 MiB. So DataFusion's
own memory tracking counts much more for this query than it really uses. Both
measurements agree that the rewritten queries use far more memory.
#### Time
How long each query took, not counting creating the table. Each was run 5
times with no memory limit.
| query | rewritten? | typical time | range |
| --- | --- | --- | --- |
| `SELECT g, sum(DISTINCT x) FROM t GROUP BY g` | yes | 0.157 s | 0.153 to
0.181 s |
| `SELECT g, sum(y), sum(DISTINCT x) FROM t GROUP BY g` | yes | 0.174 s |
0.167 to 0.175 s |
| `SELECT g, count(*), sum(y), sum(DISTINCT x) FROM t GROUP BY g` | no |
0.310 s | 0.294 to 0.340 s |
So the rewrite trades memory for speed.
#### The same queries with the rewrite turned off
To compare the exact same query with and without the rewrite, I built
`datafusion-cli` from the same commit with the rule changed to do nothing. This
is the only modified build in this report. `EXPLAIN` confirms that none of
these queries are rewritten on it. Everything else is the same as above.
| query | rewrite | finishes at | most memory used | typical time |
| --- | --- | --- | --- | --- |
| `SELECT g, sum(DISTINCT x) FROM t GROUP BY g` | on | 192M | 537 MiB |
0.157 s |
| `SELECT g, sum(DISTINCT x) FROM t GROUP BY g` | off | 96M | 306 MiB |
0.264 s |
| `SELECT g, sum(y), sum(DISTINCT x) FROM t GROUP BY g` | on | 224M | 568
MiB | 0.174 s |
| `SELECT g, sum(y), sum(DISTINCT x) FROM t GROUP BY g` | off | 96M | 305
MiB | 0.267 s |
| `SELECT g, count(*), sum(y), sum(DISTINCT x) FROM t GROUP BY g` | off |
96M | 306 MiB | 0.276 s |
With the rewrite off, all three queries cost the same: a 96M limit, about 4
MiB above the table, and 0.26 to 0.28 seconds. The `count(*)` query gives the
same numbers on both builds, apart from normal timing noise. So `count(*)` is
not doing anything special. It just switches the rewrite off.
For the same query, the rewrite uses 2 to 2.3 times the memory and makes it
1.5 to 1.7 times faster.
### Expected behavior
Whether DataFusion rewrites the query should depend on whether the rewrite
is worth it. It should not depend on whether a `count(*)` happens to be in the
query. Today, adding `count(*)` flips the same work between a fast plan that
uses a lot of memory and a slower plan that uses little.
### Additional context
- #24929 reports the same rewrite using more memory than needed, with the
same table and method.
- #24859 is the change that lets `count(*)` sit next to a `DISTINCT`
aggregate in this rule, and decides when the rewrite is allowed with it.
--
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]