Skip to content

branch-4.1:[enhance](aggregate)Add rewrite rule DecomposeRepeatWithPreAggregation (#59116)#61301

Merged
yiguolei merged 7 commits into
apache:branch-4.1from
feiniaofeiafei:pick-branch-4.1
Mar 13, 2026
Merged

branch-4.1:[enhance](aggregate)Add rewrite rule DecomposeRepeatWithPreAggregation (#59116)#61301
yiguolei merged 7 commits into
apache:branch-4.1from
feiniaofeiafei:pick-branch-4.1

Conversation

@feiniaofeiafei

@feiniaofeiafei feiniaofeiafei commented Mar 13, 2026

Copy link
Copy Markdown
Contributor

@hello-stephen

Copy link
Copy Markdown
Contributor

Thank you for your contribution to Apache Doris.
Don't know what should be done next? See How to process your PR.

Please clearly describe your PR:

  1. What problem was fixed (it's best to include specific error reporting information). How it was fixed.
  2. Which behaviors were modified. What was the previous behavior, what is it now, why was it modified, and what possible impacts might there be.
  3. What features were added. Why was this function added?
  4. Which code was refactored and why was this part of the code refactored?
  5. Which functions were optimized and what is the difference before and after the optimization?

feiniaofeiafei and others added 5 commits March 13, 2026 14:13
…on (apache#59116)

This PR add a rewrite rule. This rule will rewrite grouping sets. eg:
     select a, b, c, d, e sum(f) from t group by rollup(a, b, c, d, e);
rewrite to:
   with cte1 as (select a, b, c, d, e, sum(f) x from t group by rollup(a, b, c, d, e))
     select * fom cte1
     union all
     select a, b, c, d, null, sum(x) x from t group by rollup(a, b, c, d)

LogicalAggregate(gby: a,b,c,d,e,grouping_id output:a,b,c,d,e,grouping_id,sum(f))
  +--LogicalRepeat(grouping sets: (a,b,c,d,e),(a,b,c,d),(a,b,c),(a,b),(a),())
->
LogicalCTEAnchor
  +--LogicalCTEProducer(cte)
    +--LogicalAggregate(gby: a,b,c,d,e; aggFunc: sum(f) as x)
  +--LogicalUnionAll
    +--LogicalProject(a,b,c,d, null as e, sum(x))
      +--LogicalAggregate(gby:a,b,c,d,grouping_id; aggFunc: sum(x))
        +--LogicalRepeat(grouping sets: (a,b,c,d),(a,b,c),(a,b),(a),())
          +--LogicalCTEConsumer(aggregateConsumer)
    +--LogicalCTEConsumer(directConsumer)

The scenario where performance optimization is achieved is when
pre-aggregation significantly reduces the amount of data. This rewrite
reduces the number of rows generated by the repeat operator, resulting
in improved performance.
… input expression (apache#60045)

maybe relate PR: apache#21168

BE requires that the repeat node's output slot order should be
inconsistent with its input expressions.
That is output slots = input expressions + GroupingID + other grouping
functions.
But physical translator not ensure this requirement. Then sometimes the
repeat may have bad cast exception.

for sql:
SELECT 100000
FROM db2.table_9_50_undef_partitions2_keys3_properties4_distributed_by5
GROUP BY GROUPING SETS (
        (col_datetime_6__undef_signed, col_varchar_50__undef_signed)
        , ()
        , (col_varchar_50__undef_signed)
        , (col_datetime_6__undef_signed, col_varchar_50__undef_signed)
);

the above sql will have wrong ouput slot order

then BE will have exceptions:
(1105, 'errCode = 2, detailMessage = (172.20.57.146)
[E-7412]assert cast err:[E-7412] Bad cast from type:doris::vectorized::ColumnVector<(doris::PrimitiveType)26> to doris::vectorized::ColumnStr<unsigned int>
0#doris::Exception::Exception(int, std::basic_string_view<char, std::char_traits<char> > const&, bool) at /home/zcp/repo_center/doris_master/doris/be/src/common/exception.cpp:0\n\t1#
...
… shuffle key choosing (apache#59876)

Problem Summary:
add var agg_shuffle_use_parent_key, default value is true, when set
false, can make agg shuffle by all gby key instead of parent shuffle
key.
…ation in DecomposeRepeatWithPreAggregation (apache#60091)

Related PR: apache#59116 apache#60045

**Problem Summary:**
There are 2 problems:

1. When `DecomposeRepeatWithPreAggregation` optimization removes the
maximum grouping set, grouping functions (e.g., `grouping()`,
`grouping_id()`) that reference expressions only existing in the removed
grouping set would fail because of index lookup failure: When computing
grouping function values, `GroupingSetShapes.indexOf()` would return -1
for expressions not in `flattenGroupingSetExpression`, causing incorrect
behavior.

Example scenario:
```sql
SELECT a, b, c, grouping(c) 
FROM t1 
GROUP BY GROUPING SETS ((a, b, c), (a, b), (a), ());
```
After optimization removes the maximum grouping set `(a, b, c)`, the
expression `c` no longer exists in the remaining grouping sets. When
computing `grouping(c)`, the index lookup would fail.

2. The `repeatSlotIdList` calculation was incorrect because it relied on
the assumption that the order of `physicalRepeat.getOutputExpressions()`
matches the order of `outputTuple`. However, this assumption can be
broken during rule rewriting (e.g., in
`DecomposeRepeatWithPreAggregation.constructRepeat()` at line 502-503,
where `replacedRepeatOutputs` is constructed by `new
ArrayList<>(child.getOutput())` and then appends grouping functions,
potentially changing the order). This mismatch causes incorrect slot ID
mapping and wrong query results.

Example scenario:
When `DecomposeRepeatWithPreAggregation` rewrites a repeat plan, the
output expressions order may differ from the output tuple order, leading
to incorrect slot ID mapping in `repeatSlotIdList`.


**Solution:**

1. For problem 1: Modified `Repeat.toShapes()` method to ensure all
expressions referenced by grouping functions are included in
`flattenGroupingSetExpression`, even if they are not in any grouping
set. For expressions that only exist in the removed maximum grouping
set, they are correctly marked as `shouldBeErasedToNull = true` in all
remaining grouping sets, which is the correct SQL semantics.

Key changes:
- `Repeat.toShapes()`: Enhanced to collect all expressions referenced by
grouping functions using `ExpressionUtils.collectToList()` and merge
them into `flattenGroupingSetExpression`
- For expressions not in any grouping set, they are correctly marked as
`shouldBeErasedToNull = true` in all `GroupingSetShape` instances
- This ensures `GroupingSetShapes.indexOf()` always returns a valid
index for grouping function arguments, preventing index lookup failures

2. For problem 2: Modified `Repeat.computeRepeatSlotIdList()` method to
accept an additional `outputSlots` parameter. Instead of relying on the
order of `outputExpressions` matching `outputTuple`, the method now uses
`outputSlots` to correctly map grouping set expressions to their indices
in the output tuple. This ensures correct slot ID calculation even when
rule rewriting changes the order of output expressions.

Key changes:
- `computeRepeatSlotIdList(List<Integer> slotIdList, List<Slot>
outputSlots)`: Added `outputSlots` parameter
- `getGroupingSetsIndexesInOutput(List<Slot> outputSlots)`: Modified to
use `outputSlots` instead of `getOutputExpressions()`
- `indexesOfOutput(List<Slot> outputSlots)`: Modified to build index map
from `outputSlots` instead of `getOutputExpressions()`
- `PhysicalPlanTranslator.visitPhysicalRepeat()`: Updated to pass
`outputSlots` to `computeRepeatSlotIdList()`

3. Another task for this pull request is to set the shuffle key for the
rewritten bottom aggregation. Choosing a suitable shuffle key can
improve the aggregation rate of the top aggregation's local aggregation.
…after construct cteProducer (apache#60811)

Related PR: apache#59116 

Reproduce:
```sql
    select a,b,c,c1 from (
    select a,b,c,d,sum(d) c1 from t1 group by grouping sets((a,b,c),(a,b,c,d),(a),(a,b,c,c))
    ) t group by rollup(a,b,c,c1);
```
Error:
```shell
java.lang.IllegalArgumentException: Stats for CTE: CTEId#0 not found
```

Problem Summary:

Root cause:
This SQL triggers DecomposeRepeatWithPreAggregation twice: first for the
inner LogicalAggregate, then for the outer LogicalAggregate.
When rewriting the inner aggregate, a CTE structure is created
(CTEAnchor, CTEProducer, CTEConsumer).
When rewriting the outer aggregate, choosePreAggShuffleKeyPartitionExprs
calls StatsDerive to derive statistics for the current subtree.
At that point, the child of the outer LogicalAggregate is only the
consumer subtree; the producer subtree is not part of the plan being
visited.
StatsDerive.visitLogicalCTEConsumer looks up producer statistics by
CTEId, but the producer was never derived in this context, so the lookup
fails with Stats for CTE: CTEId#0 not found.
How to fix:
In DecomposeRepeatWithPreAggregation.constructProducer(), derive
statistics for the CTE producer when it is created. This ensures
producer stats are stored in the context before the consumer is later
visited. When StatsDerive visits LogicalCTEConsumer during the outer
rewrite, it can find the producer stats correctly.

Another fix: This pr increases the ndv requirement for the key when
selecting a shuffle key. Only when the ndv > instance number * 512 will
it be selected as the shuffle key. This is to avoid scenarios where low
ndv leads to data skew.
@feiniaofeiafei

Copy link
Copy Markdown
Contributor Author

run buildall

@feiniaofeiafei

Copy link
Copy Markdown
Contributor Author

run buildall

@feiniaofeiafei feiniaofeiafei changed the title pick decomposeRepeat related pr branch-4.1: pick decomposeRepeat related pr Mar 13, 2026
@feiniaofeiafei feiniaofeiafei changed the title branch-4.1: pick decomposeRepeat related pr branch-4.1:[enhance](aggregate)Add rewrite rule DecomposeRepeatWithPreAggregation (#59116) Mar 13, 2026
@hello-stephen

Copy link
Copy Markdown
Contributor

FE UT Coverage Report

Increment line coverage 91.93% (376/409) 🎉
Increment coverage report
Complete coverage report

@yiguolei
yiguolei merged commit f3ac2d6 into apache:branch-4.1 Mar 13, 2026
12 of 14 checks passed
yiguolei pushed a commit that referenced this pull request Mar 16, 2026
…eAggregation (#59116) (#61301)

picked from #59116 #60045 #59876 #60091 #60811

---------

Co-authored-by: yujun <yujun@selectdb.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants