mrhhsg commented on code in PR #68651:
URL: https://github.com/apache/doris/pull/68651#discussion_r4144995988


##########
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:
   Valid, fixed in 068c7fd5dcc.
   
   Reproduced on 941f02d6455 with `SELECT kbint, kstr, s, sum(s) OVER 
(PARTITION BY kbint) FROM (SELECT kbint, kstr, sum(kint) AS s FROM t GROUP BY 
kbint, kstr) a` (table bucketed by another column, `agg_phase=0`): the chosen 
plan was `OlapScan -> Exchange(HASH kbint) -> one-phase AggregationNode -> Sort 
-> Analytic`.
   
   **Exemption.** `AggregateUtils.isBucketedHashAggFusible(aggregate, 
childDistribution)` adds the translator's exact-key condition to the shared 
gate. `ChildrenPropertiesRegulator` passes the spec of the distribute it found 
below the aggregate, so the parent-key alternative is no longer exempted from 
the one-phase-with-distribute ban; it is banned again, as it was before 
bucketed aggregation. The output property deriver and the translator use the 
same overload (translator behaviour is unchanged, the key comparison just 
moved). The exact-key check itself has to stay: in the parent-key alternative 
the window consumes the aggregate without an exchange and relies on `HASH(a)`, 
which a fused aggregate does not provide.
   
   **Discount.** I first restricted the discount to the full-key child as well 
(child output properties passed to the cost model) and backed that out, because 
it regressed `nereids_rules_p0/agg_strategy/physical_agg_regulator` 
(`not_skew`): for a dedup aggregate over a CTE consumer the translator fuses 
neither alternative, but the cost model cannot see the child, so only the 
full-key one kept the discount and started to win, turning `Distribute(HASH(d)) 
-> dedup agg -> distinct agg` into `Distribute(HASH(d, a)) -> dedup agg -> 
local agg -> Distribute(HASH(d)) -> global agg`. The two alternatives of one 
aggregate therefore keep the same factor, which preserves their base-branch 
order, and the regulator, which does see the chosen child, rejects the one that 
cannot be fused. After this change a discounted parent-key aggregate over a 
distribute only survives with `agg_phase=1` (no two-phase candidate exists 
there, and it is the plan branch-4.1 already picked) and below a CTE consumer 
(th
 e existing CTE rule, unchanged). A comment in `CostModel` records this.
   
   **Also in this commit.** The aggregate-only gate now checks the phase/mode 
before the environment gate. Besides skipping the alive-backend lookup for 
aggregates that can never be fused, this fixes two 
`ChildOutputPropertyDeriverTest` cases that failed with an NPE on this PR 
(their mocked `ConnectContext` has no `Env`).
   
   **Tests.**
   - FE UT 
`BucketedAggregateTranslatorTest.testParentKeyShuffledAggregateIsNotExemptedAsBucketed`:
 asserts `Analytic <- Exchange <- BucketedAggregation <- OlapScan` and that no 
exchange sits directly on the scan, with the default cost weights, with 
`cbo_net_weight=100` (which favors the parent-key alternative) and with 
`agg_shuffle_use_parent_key=false` as the control. It fails on 941f02d6455 with 
the plan above.
   - Regression `bucketed_hash_agg` Test 11: explain contains `BUCKETED 
AGGREGATE` for a window partitioned by a subset of the GROUP BY keys, and the 
result matches the regular path.
   - Ran locally: 114 FE unit tests in 11 classes (translator / regulator / 
deriver / request deriver / cost / memo), `nereids_rules_p0/agg_strategy` (8 
suites, including `physical_agg_regulator`) and the bucketed merge / sync MV 
suites on a single-BE cluster.
   



-- 
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