Optimize `count distinct` to 3 stages of distributed agg
- Dominant language
- Go
- Stars
- 40.5k
- Forks
- 6.2k
- PR merge metrics
- PR metrics pending
Description
## Enhancement
Steps to reproduce:
```
create table log_link_visit_action(idvisit varchar(100), idvisitor varchar(100), time_dom_completion int, time_on_load int);
create table log_visit(idvisit varchar(100), visit_last_action_time date, idsite int);
alter table log_link_visit_action set tiflash replica 1;
alter table log_visit set tiflash replica 1;insert into log_link_visit_action values('1', '2', 1,1);
insert into log_link_visit_action values('1', '3', 1,1);
insert into log_link_visit_action values('1', '4', 1,1);
insert into log_visit values('1', now(), 1);
EXPLAIN
SELECT
COUNT(DISTINCT log_link_visit_action.idvisit) AS `2`,
COUNT(DISTINCT log_link_visit_action.idvisitor) AS `1`,
SUM(log_link_visit_action.time_dom_completion + log_link_visit_action.time_on_load
) AS page_load_total,
COUNT(*) AS `3`
FROM log_visit AS log_visit
INNER JOIN log_link_visit_action AS log_link_visit_action ON log_link_visit_action.idvisit = log_visit.idvisit
WHERE log_visit.idsite = 1
AND log_visit.visit_last_action_time >= '2022-02-01 00:00:00'
AND log_visit.visit_last_action_time <= '2023-04-28 23:59:59';
+--------------------------------------------------+----------+-------------------+-----------------------------+-----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
| id | estRows | task | access object | operator info |
+--------------------------------------------------+----------+-------------------+-----------------------------+-----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
| TableReader_43 | 1.00 | root | | data:ExchangeSender_42 |
| └─ExchangeSender_42 | 1.00 | batchCop[tiflash] | | ExchangeType: PassThrough |
| └─Projection_38 | 1.00 | batchCop[tiflash] | | Column#10, Column#11, Column#12, Column#13 |
| └─HashAgg_39 | 1.00 | batchCop[tiflash] | | funcs:count(distinct test.log_link_visit_action.idvisit)->Column#10, funcs:count(distinct test.log_link_visit_action.idvisitor)->Column#11, funcs:sum(Column#14)->Column#12, funcs:sum(Column#15)->Column#13 |
| └─ExchangeReceiver_41 | 1.00 | batchCop[tiflash] | | |
| └─ExchangeSender_40 | 1.00 | batchCop[tiflash] | | ExchangeType: PassThrough |
| └─HashAgg_12 | 1.00 | batchCop[tiflash] | | group by:Column#17, Column#18, funcs:sum(Column#16)->Column#14, funcs:count(1)->Column#15 |
| └─Projection_44 | 0.31 | batchCop[tiflash] | | cast(plus(test.log_link_visit_action.time_dom_completion, test.log_link_visit_action.time_on_load), decimal(41,0) BINARY)->Column#16, test.log_link_visit_action.idvisit, test.log_link_visit_action.idvisitor |
| └─HashJoin_37 | 0.31 | batchCop[tiflash] | | inner join, equal:[eq(test.log_visit.idvisit, test.log_link_visit_action.idvisit)] |
| ├─ExchangeReceiver_20(Build) | 0.25 | batchCop[tiflash] | | |
| │ └─ExchangeSender_19 | 0.25 | batchCop[tiflash] | | ExchangeType: Broadcast |
| │ └─Selection_18 | 0.25 | batchCop[tiflash] | | eq(test.log_visit.idsite, 1), ge(test.log_visit.visit_last_action_time, 2022-02-01 00:00:00.000000), le(test.log_visit.visit_last_action_time, 2023-04-28 23:59:59.000000), not(isnull(test.log_visit.idvisit)) |
| │ └─TableFullScan_17 | 10000.00 | batchCop[tiflash] | table:log_visit | keep order:false, stats:pseudo |
| └─Selection_22(Probe) | 9990.00 | batchCop[tiflash] | | not(isnull(test.log_link_visit_action.idvisit)) |
| └─TableFullScan_21 | 10000.00 | batchCop[tiflash] | table:log_link_visit_action | keep order:false, stats:pseudo |
+--------------------------------------------------+----------+-------------------+-----------------------------+-----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
```
By now, optimizer will rewrite `select count(distinct a) from t` => `select count(distinct a) from (select a from t group by a)`. If there are a big number of distinct groups, the final stage about `count distinct` will still occupy a lot of memory at a single TiFlash node.
If using `group by` to implement the `count distinct`, like
```
EXPLAIN
SELECT
-- Replacement subqueries
(SELECT COUNT(*) FROM (
(SELECT COUNT(*) FROM log_link_visit_action va
INNER JOIN log_visit v ON va.idvisit = v.idvisit
WHERE v.idsite = 1
AND v.visit_last_action_time >= '2022-02-01 00:00:00'
AND v.visit_last_action_time <= '2023-04-28 23:59:59'
GROUP BY va.idvisit) AS s
)
) AS `2`,
(SELECT COUNT(*) FROM (
(SELECT COUNT(*) FROM log_link_visit_action va
INNER JOIN log_visit v ON va.idvisit = v.idvisit
WHERE v.idsite = 1
AND v.visit_last_action_time >= '2022-02-01 00:00:00'
AND v.visit_last_action_time <= '2023-04-28 23:59:59'
GROUP BY va.idvisitor) AS s
)
) AS `1`,
-- Replacement subqueries
SUM(log_link_visit_action.time_dom_completion + log_link_visit_action.time_on_load
) AS page_load_total,
COUNT(*) AS `3`
FROM log_visit AS log_visit
INNER JOIN log_link_visit_action AS log_link_visit_action ON log_link_visit_action.idvisit = log_visit.idvisit
WHERE log_visit.idsite = 1
AND log_visit.visit_last_action_time >= '2022-02-01 00:00:00'
AND log_visit.visit_last_action_time <= '2023-04-28 23:59:59';
+---------------------------------------------+---------+--------------+-----------------------------+-----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
| id | estRows | task | access object | operator info |
+---------------------------------------------+---------+--------------+-----------------------------+-----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
| Projection_220 | 1.00 | root | | 1->Column#69, 3->Column#89, Column#48, Column#49 |
| └─StreamAgg_224 | 1.00 | root | | funcs:sum(Column#94)->Column#48, funcs:count(1)->Column#49 |
| └─Projection_252 | 0.00 | root | | cast(plus(test.log_link_visit_action.time_dom_completion, test.log_link_visit_action.time_on_load), decimal(41,0) BINARY)->Column#94 |
| └─TableReader_235 | 0.00 | root | | data:ExchangeSender_234 |
| └─ExchangeSender_234 | 0.00 | cop[tiflash] | | ExchangeType: PassThrough |
| └─HashJoin_225 | 0.00 | cop[tiflash] | | inner join, equal:[eq(test.log_visit.idvisit, test.log_link_visit_action.idvisit)] |
| ├─ExchangeReceiver_231(Build) | 0.00 | cop[tiflash] | | |
| │ └─ExchangeSender_230 | 0.00 | cop[tiflash] | | ExchangeType: Broadcast |
| │ └─Selection_229 | 0.00 | cop[tiflash] | | eq(test.log_visit.idsite, 1), ge(test.log_visit.visit_last_action_time, 2022-02-01 00:00:00.000000), le(test.log_visit.visit_last_action_time, 2023-04-28 23:59:59.000000), not(isnull(test.log_visit.idvisit)) |
| │ └─TableFullScan_228 | 1.00 | cop[tiflash] | table:log_visit | keep order:false, stats:pseudo |
| └─Selection_233(Probe) | 3.00 | cop[tiflash] | | not(isnull(test.log_link_visit_action.idvisit)) |
| └─TableFullScan_232 | 3.00 | cop[tiflash] | table:log_link_visit_action | keep order:false, stats:pseudo |
+---------------------------------------------+---------+--------------+-----------------------------+-----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
```
There will be 2 extra uncorrelated subqueries, which means there are at least 3 times of table-scans.
We hope to optimize `count distinct` by 3 stage distracted agg and require one time of table-scans.
1. shuffle `group by`
2. shuffle and compute `count distinct`
3. collect and compute result.
Contributor guide
Assessment
This issue has not been assessed yet.