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


##########
fe/fe-core/src/main/java/org/apache/doris/nereids/processor/post/ProjectAggregateExpressionsForCse.java:
##########
@@ -151,9 +170,62 @@ public Plan 
visitPhysicalHashAggregate(PhysicalHashAggregate<? extends Plan> agg
             PhysicalProperties projectPhysicalProperties = 
ChildOutputPropertyDeriver.computeProjectOutputProperties(
                     project.getProjects(), ((PhysicalPlan) 
project.child()).getPhysicalProperties());
             project = 
project.withPhysicalPropertiesAndStats(projectPhysicalProperties, 
project.getStats());
-            aggregate = (PhysicalHashAggregate<? extends Plan>) aggregate
-                    .withAggOutput(aggOutputReplaced)
+            return (Plan) aggregate.withAggOutput(aggOutputReplaced)
                     .withChildren(project);
+        } else if (aggregate.child() instanceof PhysicalDistribute) {
+            // One-phase (INPUT_TO_RESULT) aggregate over a distribute
+            // (aggregate -> distribute -> scan): insert the CSE project 
between
+            // the distribute and its child, instead of between the aggregate 
and
+            // the distribute. This keeps the aggregate's child as a distribute
+            // (so bucketed fusion and the property machinery still see the 
same
+            // shape), and the project lands inside the scan
+            // fragment, so the common aggregate argument is computed once per 
row
+            // before the exchange. After bucketed fusion bypasses the 
distribute,
+            // the executed plan is BucketedAgg(sum(x), max(x)) -> Project(a+b 
AS x)
+            // -> scan.
+            //
+            // Only the one-phase shape reaches here with complex aggregate
+            // arguments: two-phase GLOBAL aggregates (BUFFER_TO_RESULT) 
reference
+            // the local phase's intermediate slots, so no CSE candidate is
+            // extracted for them anyway. Guard explicitly anyway to keep the
+            // intent clear and to stay safe if a future aggregate function
+            // surfaces a non-slot argument on the GLOBAL phase.
+            if (!(aggregate instanceof PhysicalHashAggregate)) {
+                return aggregate;
+            }
+            PhysicalHashAggregate<? extends Plan> hashAggregate =
+                    (PhysicalHashAggregate<? extends Plan>) aggregate;
+            if (hashAggregate.getAggPhase() != AggPhase.GLOBAL
+                    || hashAggregate.getAggMode() != AggMode.INPUT_TO_RESULT) {
+                return aggregate;
+            }
+            PhysicalDistribute<?> distribute = (PhysicalDistribute<?>) 
aggregate.child();
+            List<NamedExpression> projections = new ArrayList<>();
+            projections.addAll(inputSlots);
+            projections.addAll(cseCandidates.values());
+            List<Slot> projectOutput = new ImmutableList.Builder<Slot>()
+                    .addAll(inputSlots)
+                    .addAll(slotMap.values())
+                    .build();
+            LogicalProperties projectLogicalProperties = new LogicalProperties(
+                    () -> projectOutput,
+                    () -> DataTrait.EMPTY_TRAIT
+            );
+            AbstractPhysicalPlan distributeChild = ((AbstractPhysicalPlan) 
distribute.child());
+            PhysicalProperties projectPhysicalProperties = 
ChildOutputPropertyDeriver.computeProjectOutputProperties(
+                    projections, distributeChild.getPhysicalProperties());
+            PhysicalProject<? extends Plan> project = new 
PhysicalProject<>(projections, Optional.empty(),

Review Comment:
   Fixed in fee2f89e12c.
   
   Reproduced on a local cluster with the shape you describe (`agg_phase=1`, a 
materialized CTE referenced twice, `SELECT g, SUM(a+b), MAX(a+b) FROM (SELECT 
k+1 g, a, b FROM c) x GROUP BY g UNION ALL ...`): the plan was `Agg -> 
Distribute -> Project(a, b, g, a+b) -> Project(k+1 AS g, a, b) -> CTEConsumer` 
and the translation failed with `generate invalid plan`.
   
   `ProjectAggregateExpressionsForCse` now merges the CSE project into the 
project that is already below the distribute instead of stacking a second one. 
It uses the regular project-merge rules (`Project.canMergeChildProjections` / 
`mergeProjections`, the same ones `MergeProjectPostProcessor` applies before 
this processor runs), so CSE expressions that reference an alias of the 
existing project are rewritten to its input and the output expr ids stay the 
same. When the two projects cannot be merged (e.g. a volatile alias would be 
evaluated twice), the aggregate is left as it is rather than stacking a 
project. The CTE sink therefore always receives one projection; as a side 
effect the non-CTE shape no longer gets an extra `VSELECT` above the scan (the 
scan projection computes `a+b`), which also holds for the fused bucketed plan.
   
   Tests:
   - FE UT `ProjectAggregateExpressionsForCseTest` (new): merge into the scan 
project, CSE referencing an alias of the existing project, the projected 
materialized CTE consumer case (one multicast fragment, two consumer sinks, 
each with a single projection that computes both `k+N` and `a+b`), and the 
unmergeable volatile case. All four fail on the previous head; the CTE one 
fails with `generate invalid plan`.
   - Regression `cse_agg_distribute`: new projected materialized CTE case 
(explain + `order_qt`, results equal to `enable_aggregate_cse=false`); the join 
case now asserts `notContains("VSELECT")` and that `a+b` is evaluated only in 
the two scan projections.
   
   This push also addresses the residual of 
https://github.com/apache/doris/pull/68651#discussion_r4140769120 raised in the 
review summary: the pin is no longer checked on an earlier resolution of the 
replica. `OlapScanNode.addScanRangeLocations` now compares the pinned backend 
with the backend id resolved inside the location loop, i.e. the very value 
written into the `TScanRangeLocation` (`getBackendIdWithClusterId` for a 
`CloudReplica`), and skips the replica otherwise; a tablet left without a 
location on the pinned backend fails the query with a retry hint. 
`BucketedAggregateMultiBackendTest.testPinHoldsForTheBackendResolvedWhenTheLocationIsBuilt`
 simulates a replica whose two resolutions differ (the second one landing on 
the recovered backend): on the previous head a location was published on the 
recovered backend, now the query fails with the retry hint.
   



##########
fe/fe-core/src/main/java/org/apache/doris/nereids/glue/translator/PhysicalPlanTranslator.java:
##########
@@ -3216,6 +3200,262 @@ 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) 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. Only fuse when the
+        // parent chain keeps an exchange boundary (e.g. a top-level 
aggregate).
+        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);
+    }
+
+    /** Returns true if the plan subtree contains a physical CTE consumer. */
+    private boolean containsCTEConsumer(Plan plan) {
+        if (plan instanceof PhysicalCTEConsumer) {
+            return true;
+        }
+        for (Plan child : plan.children()) {
+            if (containsCTEConsumer(child)) {
+                return true;
+            }
+        }
+        return false;
+    }
+
+    /**
+     * Returns true if the plan subtree is a unary pipeline over exactly one 
olap
+     * scan, i.e. it translates into a single-scan fragment that bucketed 
fusion
+     * can safely build upon. Subtrees containing fragment-merging or
+     * distribution-changing nodes (join / set-op / CTE / nested aggregate /
+     * storage-layer aggregate) are rejected.
+     */
+    private boolean isSingleOlapScanPipeline(Plan plan) {
+        if (plan instanceof PhysicalOlapScan) {
+            return true;
+        }
+        if (plan instanceof PhysicalHashJoin
+                || plan instanceof PhysicalNestedLoopJoin
+                || plan instanceof PhysicalSetOperation
+                || plan instanceof PhysicalCTEConsumer
+                || plan instanceof PhysicalCTEAnchor
+                || plan instanceof PhysicalHashAggregate
+                || plan instanceof PhysicalStorageLayerAggregate) {
+            return false;
+        }
+        if (plan.children().size() == 1) {
+            return isSingleOlapScanPipeline(plan.child(0));
+        }
+        return false;
+    }
+
+    /**
+     * Fuse a one-phase GLOBAL hash aggregate and its PhysicalDistribute child
+     * into a BucketedAggregationNode, skipping the exchange node entirely.
+     * Visits the distribute's child directly to keep everything in one 
fragment.
+     */
+    private PlanFragment visitBucketedFusion(
+            PhysicalHashAggregate<? extends Plan> aggregate,
+            PlanTranslatorContext context) {
+        // Visit the distribute's direct child, bypassing the distribute 
entirely.
+        // This avoids creating an ExchangeNode that bucketed agg does not 
need.
+        Plan distributeChild = aggregate.child(0).child(0);
+        PlanFragment inputPlanFragment = distributeChild.accept(this, context);

Review Comment:
   Follow-up for the residual raised in the latest review summary (cloud 
multi-replica: the pinned-backend filter resolved a `CloudReplica` once and the 
scan-range construction resolved it again): fixed in fee2f89e12c.
   
   The separate pre-filter (`keepReplicasOnPinnedBackend`, which used 
`getBackendIdWithoutException()`) is removed. The pin is now checked inside the 
location loop of `OlapScanNode.addScanRangeLocations`, against the backend id 
that was just resolved for the location 
(`getBackendIdWithClusterId(cachedClusterId)` for a `CloudReplica`) and that is 
written into the `TScanRangeLocation` right after. There is a single resolution 
per replica, so the published location cannot differ from the one the pin was 
checked on, whatever `CloudReplica.getBackendIdImpl` returns on another call. A 
replica that resolves to another alive backend is skipped; if no location on 
the pinned backend remains, the query fails with `... is served by backend X 
instead of backend Y, the only alive backend when the query was planned; ... 
please retry` instead of splitting groups across backends.
   
   Test: 
`BucketedAggregateMultiBackendTest.testPinHoldsForTheBackendResolvedWhenTheLocationIsBuilt`
 (2 BEs, `replication_num=2`) makes the resolution used for the location return 
the recovered backend while every other resolution still returns the replica's 
original backend. On the previous head the query planned fine and published a 
location on the recovered backend (the test fails with "nothing was thrown"); 
now it fails with the retry hint. The existing 
`testFusedScanIsPinnedToTheBackendSeenBySingleBackendGate` still passes (one 
location per tablet, on the pinned backend). Not verified on a real cloud 
deployment.
   



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