github-actions[bot] commented on code in PR #68651:
URL: https://github.com/apache/doris/pull/68651#discussion_r4143580494
##########
fe/fe-core/src/main/java/org/apache/doris/nereids/glue/translator/PhysicalPlanTranslator.java:
##########
@@ -3216,6 +3219,294 @@ private PlanFragment connectJoinNode(HashJoinNode
hashJoinNode, PlanFragment lef
return leftFragment;
}
+ /**
+ * Check whether the one-phase GLOBAL hash aggregate can be fused with its
+ * distribute child into a BucketedAggregationNode. This eliminates
exchange
+ * overhead on single-BE deployments by using in-memory per-bucket merging.
+ */
+ private boolean shouldUseBucketedFusion(PhysicalHashAggregate<? extends
Plan> aggregate,
+ PlanTranslatorContext context) {
+ // Shared eligibility (also used by the regulator, the output property
deriver
+ // and the cost model): session var, single-BE, GROUP BY, spill /
query cache
+ // off, smooth upgrade, no UDAF, one-phase GLOBAL INPUT_TO_RESULT, no
partial
+ // (buffer-producing) function, two-phase capable functions, no pushed
TopN.
+ if (!AggregateUtils.isBucketedHashAggFusible(aggregate)) {
+ return false;
+ }
+ // Child must be PhysicalDistribute with hash distribution matching
group keys
+ Plan child = aggregate.child(0);
+ if (!(child instanceof PhysicalDistribute)) {
+ return false;
+ }
+ // Bucketed fusion bypasses the distribute/exchange and builds
directly on the
+ // child fragment. When the child subtree contains a CTE consumer
(materialized
+ // multicast CTE), the child fragment is the MultiCastPlanFragment; a
parent
+ // distribute would then treat the aggregate output slots as consumer
slots and
+ // fail with "Required producer slot ... doesn't exist". Fall back to
the
+ // regular one-phase path (which keeps the exchange) for such plans.
+ if (containsCTEConsumer(child)) {
+ return false;
+ }
+ // The distribute's child subtree must be a unary pipeline over
exactly one
+ // olap scan. Fusing an aggregate whose input contains a join / set-op
/ CTE
+ // subtree would leave multiple olap scans in a single fragment
(rejected by
+ // UnassignedJobBuilder: "Not supported multiple scan multiple
OlapTable but
+ // not contains colocate join or bucket shuffle join"), and fusing
over a
+ // nested aggregate would break the bucket alignment between stages.
+ if (!isSingleOlapScanPipeline(aggregate.child(0).child(0))) {
+ return false;
+ }
+ // The parent is a fragment-merging node (join / set-op / recursive
union) that
+ // consumes this fragment without an exchange boundary: fusing removes
the
+ // exchange that keeps the scan in its own fragment, so multiple scans
would end
+ // up in the same fragment and the scan-assignment would fail. A
distribute
+ // between the merging node and this aggregate clears the context (see
+ // visitPhysicalDistribute), because its exchange keeps the fused
fragment apart.
+ if (context.isInFragmentMergeChild()) {
+ return false;
+ }
+ DistributionSpec distSpec = ((PhysicalDistribute<?>)
child).getDistributionSpec();
+ if (!(distSpec instanceof DistributionSpecHash)) {
+ return false;
+ }
+ List<ExprId> distKeys = ((DistributionSpecHash)
distSpec).getOrderedShuffledColumns();
+ List<ExprId> groupByKeys = aggregate.getGroupByExpressions().stream()
+ .filter(SlotReference.class::isInstance)
+ .map(SlotReference.class::cast)
+ .map(SlotReference::getExprId)
+ .collect(Collectors.toList());
+ return distKeys.equals(groupByKeys);
Review Comment:
[P2] Cost subset-key aggregates as unfused. A reduced plan is
`Window(PARTITION BY a) -> GlobalAgg(GROUP BY a,b) -> Distribute(HASH(a)) ->
OlapScan`. With the default `agg_shuffle_use_parent_key`, the parent-key
alternative can satisfy the window without an upper exchange and win over the
full-key alternative. `CostModel` already halves this one-phase aggregate's row
cost, but this exact-key check rejects fusion, so translation retains an
exchange on raw rows and the regular aggregate. For large input with few
`(a,b)` groups, the discount can favor this over local preaggregation. Restrict
the discount/one-phase exception to a finally fusible full-key child, or cost
this subset path as an unfused aggregate; add a plan test.
--
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]