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]

Reply via email to