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


##########
fe/fe-core/src/main/java/org/apache/doris/nereids/util/AggregateUtils.java:
##########
@@ -161,4 +166,83 @@ public static Set<NamedExpression> 
getDistinctNamedExpr(LogicalAggregate<? exten
                 .map(NamedExpression.class::cast)
                 .collect(ImmutableSet.toImmutableSet());
     }
+
+    /**
+     * Check if order keys are identical to group-by keys (1-1 mapping, same 
order).
+     * Shared utility used by both PushTopnToAgg and SplitAggWithoutDistinct.
+     */
+    public static boolean isOrderKeysMatchGroupKeys(List<OrderKey> orderKeys,
+            List<Expression> groupByKeys) {
+        if (orderKeys.size() != groupByKeys.size()) {
+            return false;
+        }
+        for (int i = 0; i < groupByKeys.size(); i++) {
+            if (!groupByKeys.get(i).equals(orderKeys.get(i).getExpr())) {
+                return false;
+            }
+        }
+        return true;
+    }
+
+    /**
+     * Check the basic environmental conditions for bucketed hash aggregation.
+     * This is the shared eligibility gate used by ChildrenPropertiesRegulator
+     * (to allow the one-phase-GLOBAL+distribute pattern), CostModel (for cost
+     * discount), and PhysicalPlanTranslator (for fusion into 
BucketedAggregationNode).
+     *
+     * @return true if the session variable is enabled, there is exactly one 
alive BE,
+     *         spill is disabled, no smooth upgrade is in progress, the 
aggregate has
+     *         GROUP BY keys and contains no user-defined aggregate function.
+     */
+    public static boolean isBucketedHashAggEnabled(Aggregate<? extends Plan> 
aggregate) {
+        ConnectContext ctx = ConnectContext.get();
+        if (ctx == null) {
+            return false;
+        }
+        if (!ctx.getSessionVariable().enableBucketedHashAgg) {
+            return false;
+        }
+        // Must have GROUP BY keys (without-key aggregation not supported)
+        if (aggregate.getGroupByExpressions().isEmpty()) {
+            return false;
+        }
+        // Bucketed agg has no spill support. Keep the regular (spillable) 
aggregation
+        // when spill is enabled, otherwise a high-cardinality GROUP BY could 
hit the
+        // memory limit instead of spilling.
+        if (ctx.getSessionVariable().enableSpill || 
ctx.getSessionVariable().enableForceSpill) {
+            return false;
+        }
+        // be_number_for_test can only disable bucketed agg (to test the 
multi-BE plan),
+        // never bypass the single-BE gate below.
+        int beNumberForTest = ctx.getSessionVariable().getBeNumberForTest();
+        if (beNumberForTest > 0 && beNumberForTest != 1) {
+            return false;
+        }
+        // Correctness gate: single-BE only (cross-BE in-memory merge is 
impossible).
+        // Scan ranges always go to the real alive backends, so count them 
directly
+        // (getBackendsNumber() would return be_number_for_test).
+        // Note: do not clamp to 1 — with zero backends bucketed agg must not 
be enabled.
+        SystemInfoService clusterInfo = ctx.getEnv().getClusterInfo();
+        List<Long> aliveBackendIds = 
clusterInfo.getAllBackendByCurrentCluster(true);
+        if (aliveBackendIds.size() != 1) {
+            return false;
+        }
+        // Smooth upgrade safety net: old BE processes do not recognize
+        // BUCKETED_AGGREGATION_NODE plan node type
+        for (Long beId : aliveBackendIds) {

Review Comment:
   Not changed, with the reasoning below.
   
   - The supported rolling-upgrade order for Doris is BE first, then FE, so an 
upgraded FE is not expected to plan for a BE older than itself. New plan node 
types rely on this order, and the planner has no per-backend capability check 
for them.
   - There is no backend capability or exec-version report to check today: 
`be_exec_version` is an FE config that is sent down to the BE in the query 
options, and the heartbeat only carries a free-form version string. A real 
capability gate needs a new heartbeat field reported by the BE, which is a 
protocol change that should be designed on master first rather than introduced 
in a backport.
   - This gate is byte-identical to master. The `isSmoothUpgradeSrc` check 
covers the one case where old and new BE processes are expected to coexist 
behind a new FE.
   - The path can be switched off per session or globally with 
`enable_bucketed_hash_agg=false` if a deployment has to run FE ahead of BE.



##########
be/src/exec/operator/bucketed_aggregation_sink_operator.h:
##########
@@ -0,0 +1,147 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+#pragma once
+
+#include <stdint.h>
+
+#include "exec/operator/operator.h"
+#include "runtime/exec_env.h"
+#include "runtime/runtime_profile.h"
+
+namespace doris {
+
+class BucketedAggSinkOperatorX;
+
+/// Sink-side local state for bucketed hash aggregation.
+/// Each pipeline instance builds 256 per-bucket hash tables (two-level hash 
table).
+/// No locking: each instance writes to per_instance_data[_instance_idx].
+class BucketedAggSinkLocalState : public 
PipelineXSinkLocalState<BucketedAggSharedState> {
+public:
+    ENABLE_FACTORY_CREATOR(BucketedAggSinkLocalState);
+    using Base = PipelineXSinkLocalState<BucketedAggSharedState>;
+    BucketedAggSinkLocalState(DataSinkOperatorXBase* parent, RuntimeState* 
state);
+    ~BucketedAggSinkLocalState() override = default;
+
+    Status init(RuntimeState* state, LocalSinkStateInfo& info) override;
+    Status open(RuntimeState* state) override;

Review Comment:
   Fixed in 5eb6a1de61c.
   
   - `BucketedAggSinkLocalState::is_blockable()` now returns 
`BucketedAggSinkOperatorX::has_blockable_aggregate()`, which checks 
`AggFnEvaluator::is_blockable()` (function and argument expressions). It reads 
the operator's own evaluators instead of the local/shared clones, because those 
are only created in `open()`, after the first submit.
   - The source merges and finalizes the same functions, so 
`BucketedAggSourceOperatorX::is_blockable(state)` now also follows the paired 
sink operator (wired in `pipeline_fragment_context.cpp`).
   - Added BE UT 
`AggOperatorBlockableTest.bucketed_agg_follows_blockable_aggregate` covering 
sink and source, before the local state is opened.



##########
fe/fe-core/src/main/java/org/apache/doris/nereids/glue/translator/PhysicalPlanTranslator.java:
##########
@@ -3216,6 +3202,351 @@ 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: session var, single-BE, GROUP BY, smooth 
upgrade, no UDAF
+        if (!AggregateUtils.isBucketedHashAggEnabled(aggregate)) {
+            return false;
+        }
+        // Must be one-phase: GLOBAL + INPUT_TO_RESULT
+        if (aggregate.getAggPhase() != AggPhase.GLOBAL
+                || aggregate.getAggMode() != AggMode.INPUT_TO_RESULT) {
+            return false;
+        }
+        // BucketedAggregationNode always finalizes into the output tuple slot
+        // types (need_finalize=true, isPartial=false), so fusing an aggregate
+        // whose functions produce buffers (partial) would fail the BE
+        // result-type check: the slot type of a buffer-producing function is
+        // Varchar while the function's final return type (e.g. DOUBLE for
+        // stddev) is what insert_result_into writes. The one-phase GLOBAL
+        // dedup aggregate of a 3-phase DISTINCT plan has exactly this shape —
+        // the node itself is INPUT_TO_RESULT but its non-distinct functions
+        // run in INPUT_TO_BUFFER mode — and must stay on the regular
+        // AggregationNode path, which serializes when isPartial.
+        if (containsPartialAggFunction(aggregate)) {
+            return false;
+        }
+        // Exclude one-phase-only aggregates (e.g. GROUP_CONCAT with ORDER BY).
+        // BucketedAggregationNode has no sort-info field, so fusing would drop
+        // the aggregate ORDER BY contract. Only aggregates supporting 
two-phase
+        // execution can be safely fused.
+        if (!supportsTwoPhaseAgg(aggregate)) {
+            return false;
+        }
+        // BucketedAggregationNode does not support sortByGroupKey 
(PushTopnToAgg
+        // optimization). Regular AggregationNode fills sort info; fusing 
would drop it.
+        if (aggregate.getTopnPushInfo() != null) {
+            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;
+    }
+
+    /**
+     * Check whether all aggregate functions in this physical hash aggregate
+     * support two-phase execution. One-phase-only aggregates (e.g. 
GROUP_CONCAT
+     * with ORDER BY) cannot be bucketed because BucketedAggregationNode does 
not
+     * carry sort-info metadata (aggSortInfos); fusing them would drop the
+     * aggregate ORDER BY contract and produce unordered results.
+     */
+    private boolean supportsTwoPhaseAgg(PhysicalHashAggregate<? extends Plan> 
aggregate) {
+        for (NamedExpression o : aggregate.getOutputExpressions()) {
+            AtomicBoolean foundOnePhaseOnly = new AtomicBoolean(false);
+            o.foreach(c -> {
+                if (c instanceof OrderExpression) {
+                    // Any aggregate function with an internal ORDER BY
+                    // (e.g. GROUP_CONCAT(... ORDER BY ...)) needs sort-info
+                    // metadata, which BucketedAggregationNode does not carry.
+                    foundOnePhaseOnly.set(true);
+                    return true;
+                }
+                if (c instanceof AggregateExpression) {
+                    AggregateFunction func = ((AggregateExpression) 
c).getFunction();
+                    if (!func.supportAggregatePhase(AggregatePhase.TWO)) {
+                        foundOnePhaseOnly.set(true);
+                        return true;
+                    }
+                }
+                return false;
+            });
+            if (foundOnePhaseOnly.get()) {
+                return false;
+            }
+        }
+        return true;
+    }
+
+    /**
+     * Check whether the aggregate's output contains any buffer-producing
+     * (partial) aggregate function, i.e. an AggregateExpression in a mode with
+     * productAggregateBuffer=true. BucketedAggregationNode cannot carry such
+     * functions: it always finalizes into the output tuple slot types, while a
+     * buffer-producing function's slot type is the serialized Varchar type
+     * (AggregateExpression.getDataType) — writing the final result (e.g. 
DOUBLE
+     * for stddev) into that String column fails the BE result-type check.
+     */
+    private boolean containsPartialAggFunction(PhysicalHashAggregate<? extends 
Plan> aggregate) {
+        for (NamedExpression o : aggregate.getOutputExpressions()) {
+            AtomicBoolean foundPartial = new AtomicBoolean(false);
+            o.foreach(c -> {
+                if (c instanceof AggregateExpression) {
+                    if (((AggregateExpression) 
c).getAggregateParam().aggMode.productAggregateBuffer) {
+                        foundPartial.set(true);
+                    }
+                    return true;
+                }
+                return false;
+            });
+            if (foundPartial.get()) {
+                return true;
+            }
+        }
+        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);
+
+        List<Expression> groupByExpressions = 
aggregate.getGroupByExpressions();
+        List<NamedExpression> outputExpressions = 
aggregate.getOutputExpressions();
+
+        // 1. generate slot reference for each group expression
+        List<SlotReference> groupSlots = 
collectGroupBySlots(groupByExpressions, outputExpressions);
+        ArrayList<Expr> execGroupingExpressions = 
translateGroupByExprs(groupByExpressions, context);
+
+        // 2. collect agg expressions and generate agg function to slot 
reference map.
+        //    Reuse the shared helper from visitPhysicalHashAggregate; the 
bucketed
+        //    path passes null for hasPartialOut (never partial, always 
needsFinalize).
+        Pair<List<Slot>, ArrayList<FunctionCallExpr>> aggResult =
+                collectAggFunctions(outputExpressions, null, context);
+        List<Slot> aggFunctionOutput = aggResult.first;
+        ArrayList<FunctionCallExpr> execAggregateFunctions = aggResult.second;
+
+        // 3. generate output tuple
+        Pair<TupleDescriptor, List<Integer>> tupleAndIds =
+                buildAggOutputTuple(groupSlots, aggFunctionOutput, context);
+        TupleDescriptor outputTupleDesc = tupleAndIds.first;
+        List<Integer> aggFunOutputIds = tupleAndIds.second;
+
+        // Bucketed agg uses AggPhase.FIRST (update semantics): raw input -> 
final result.
+        // Not partial — always needsFinalize.
+        AggregateInfo aggInfo = AggregateInfo.create(execGroupingExpressions, 
execAggregateFunctions,
+                aggFunOutputIds, false /* isPartial */, outputTupleDesc,
+                AggregateInfo.AggPhase.FIRST);
+
+        BucketedAggregationNode bucketedAggNode = new BucketedAggregationNode(
+                context.nextPlanNodeId(), inputPlanFragment.getPlanRoot(), 
aggInfo, true);

Review Comment:
   Fixed in 5eb6a1de61c.
   
   `AggregateUtils.isBucketedHashAggEnabled` now returns false when 
`enable_query_cache` is on, so the regular plan with the cacheable LOCAL 
`AggregationNode` is kept. The gate is shared by the regulator, the cost model 
and the translator, so they stay consistent.
   
   Tests: FE UT 
`BucketedAggregateTranslatorTest.testQueryCacheKeepsRegularAggregation` (no 
`BucketedAggregationNode`, and an `AggregationNode` that is a query cache 
candidate), plus an explain check in the `bucketed_hash_agg` regression suite.



##########
regression-test/suites/nereids_rules_p0/agg_strategy/cse_agg_distribute.groovy:
##########
@@ -0,0 +1,69 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+suite("cse_agg_distribute") {
+    sql "SET enable_nereids_planner=true"
+    sql "SET enable_fallback_to_original_planner=false"
+    sql "SET runtime_filter_mode=OFF"
+
+    sql "DROP TABLE IF EXISTS cse_agg_distribute_tbl"
+    sql """
+        CREATE TABLE cse_agg_distribute_tbl (
+            id int,
+            grp varchar(20),
+            a int,
+            b int
+        ) DUPLICATE KEY(id)
+        DISTRIBUTED BY HASH(id) BUCKETS 3
+        PROPERTIES('replication_num' = '1')
+    """
+    sql """ INSERT INTO cse_agg_distribute_tbl VALUES
+        (1, 'g1', 1, 2),
+        (2, 'g2', 3, 4),
+        (3, 'g1', 5, 6),
+        (4, 'g2', 7, 8),
+        (5, 'g1', 9, 10)
+    """
+
+    // SUM(a+b) and MAX(a+b) share the same argument, so the aggregate-argument
+    // CSE must extract "a+b" into a project node and make both functions
+    // reference the extracted slot, instead of re-evaluating a+b per function.
+    String query = "SELECT grp, SUM(a+b), MAX(a+b) FROM cse_agg_distribute_tbl 
GROUP BY grp"
+
+    // ---------------------------------------------------------------------
+    // one-phase aggregate over a distribute (the aggregate is a join child,
+    // so the distribute is required by the join): the CSE project must be
+    // inserted below the distribute, keeping the distribution-key slots
+    // intact. Both aggregates must reference the extracted slot (4
+    // occurrences: SUM/MAX of each side).
+    // ---------------------------------------------------------------------
+    sql "set agg_phase=1"
+    sql "set enable_bucketed_hash_agg=false"
+    String joinQuery = """
+        SELECT t1.grp, t1.s, t1.m, t2.s2, t2.m2 FROM
+         (SELECT grp, SUM(a+b) s, MAX(a+b) m FROM cse_agg_distribute_tbl GROUP 
BY grp) t1
+         JOIN (SELECT grp, SUM(a+b) s2, MAX(a+b) m2 FROM 
cse_agg_distribute_tbl GROUP BY grp) t2
+         ON t1.grp = t2.grp
+    """
+    explain {
+        sql("${joinQuery}")
+        contains("VEXCHANGE")
+        contains("VSELECT")
+        multiContains("cast(a as BIGINT) + cast(b as BIGINT))[#", 4)

Review Comment:
   Not changed: the statement is well formed.
   
   `"cast(a as BIGINT) + cast(b as BIGINT))[#"` is a single string literal; the 
`)` and `[#` are inside the quotes. It matches the explain text `(cast(a as 
BIGINT) + cast(b as BIGINT))[#N]`, and the call is `multiContains(String, int)` 
with `4` as the second argument. The suggested pattern without `)[#` would also 
match other occurrences of the expression and change the expected count.
   
   The suite parses and runs: 
`nereids_rules_p0/agg_strategy/cse_agg_distribute` passed locally on this head 
(`Test 2 suites, failed 0 suites` together with `bucketed_hash_agg`). The line 
is identical to master.



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