branch-4.1:[enhance](aggregate)Add rewrite rule DecomposeRepeatWithPreAggregation (#59116)#61301
Merged
Merged
Conversation
Contributor
|
Thank you for your contribution to Apache Doris. Please clearly describe your PR:
|
…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
force-pushed
the
pick-branch-4.1
branch
from
March 13, 2026 06:24
7af6f10 to
0c49339
Compare
Contributor
Author
|
run buildall |
Contributor
Author
|
run buildall |
Contributor
FE UT Coverage ReportIncrement line coverage |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
picked from #59116 #60045 #59876 #60091 #60811