apache / apache/gluten

[VL] Optimize the plan tranformation for aggregation with distinct

Open
#10,527 1 comment 0 reactions 0 assignees View on GitHub
enhancement
Dominant language
Scala
Stars
1.6k
Forks
657
Avg merge
2d 21h
Merged PRs (30d)
85

Description

### Description

### Use Case
We have encountered a production use case where the aggregation before `partial(count distinct)` has poor performance because the majority of data has a unique group by keys. But the aggregation of `partial(count distinct)` has significantly smaller cardinality since the group by keys is only a portion.

metrics:

Image

gluten plan:
```
: +- ^(14) FlushableHashAggregateTransformer(keys=[brand_id#605L, gender#606, brand_layered#607, mt_back_city_nation_id#608L, mt_back_province_nation_id#609L, spark_grouping_id#604L], functions=[merge_max(brand_name#504), merge_max(brand_cate2_id#505L), merge_max(brand_cate2#506), merge_max(if (spark_grouping_id#604L IN (11,15)) mt_back_city_nation_name#508 else all), merge_max(if (spark_grouping_id#604L IN (19,23)) mt_back_province_nation_name#510 else all), partial_count(distinct user_id#515)], isStreamingAgg=false, output=[brand_id#605L, gender#606, brand_layered#607, mt_back_city_nation_id#608L, mt_back_province_nation_id#609L, spark_grouping_id#604L, max#737, max#739L, max#741, max#743, max#745, count#748L])
: +- ^(14) HashAggregateTransformer(keys=[brand_id#605L, gender#606, brand_layered#607, mt_back_city_nation_id#608L, mt_back_province_nation_id#609L, spark_grouping_id#604L, user_id#515], functions=[merge_max(brand_name#504), merge_max(brand_cate2_id#505L), merge_max(brand_cate2#506), merge_max(if (spark_grouping_id#604L IN (11,15)) mt_back_city_nation_name#508 else all), merge_max(if (spark_grouping_id#604L IN (19,23)) mt_back_province_nation_name#510 else all)], isStreamingAgg=false, output=[brand_id#605L, gender#606, brand_layered#607, mt_back_city_nation_id#608L, mt_back_province_nation_id#609L, spark_grouping_id#604L, user_id#515, max#737, max#739L, max#741, max#743, max#745])
```

### Proposed Solution
Since Velox "natively" supports distinct aggregation we can add a rule to fuse the two Spark operators to a single aggregation with distinct. There would be performance gain because the new aggregation will have a significantly smaller aggregation buffer. Here is an example:

Velox plan before this optimization:
```
keys=key final_sum(value) final_count(id, distinct=false)
exchange keys=key
keys=key merge_sum(value) partial_count(id, distinct=false)
keys=id,key merge_sum(value)
exchange keys=id,key
keys=id,key partial_sum(value)
```

Velox plan after this optimization:
```
keys=key final_sum(value) final_count(id, distinct=false)
exchange keys=key
keys=key merge_sum(value) partial_count(id, distinct=true)
exchange keys=id,key
keys=id,key partial_sum(value)
```

### Gluten version

None

Contributor guide

Open the contributing guide

Research direction

Start by tracing the two Spark aggregation operators shown in the issue and how Gluten maps them to Velox aggregation support. Compare the current plan with the proposed fused plan, then verify that distinct aggregation preserves the other aggregate results and produces the intended lower-cardinality plan and performance improvement.

Written by the indexing model from the issue text.

Assessment

Tech stack
scala, spark
Domain
performance
Issue type
Feature
Difficulty
5/5
Estimated time
Over a week
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
35/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.