github-actions[bot] commented on code in PR #68651:
URL: https://github.com/apache/doris/pull/68651#discussion_r4141692879


##########
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:
   [P1] Preserve a single multicast projection for projected CTE consumers. 
With `agg_phase=1`, a twice-referenced materialized CTE can give this branch 
`Agg -> Distribute(HASH(g)) -> Project(k+1 AS g, a, b) -> CTEConsumer`; 
`SUM(a+b)` and `MAX(a+b)` make it insert another project above the existing 
one. `MergeProjectPostProcessor` has already run. Both projects then visit the 
same `MultiCastPlanFragment`: the inner one fills its 
`DataStreamSink.projections`, and the outer one throws `generate invalid plan` 
in `visitPhysicalProject`. Merge the CSE expressions into the existing project, 
or otherwise ensure the CTE sink receives one projection; add a planning case 
for this shape.



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