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]