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


##########
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:
   [P2] Mark the bucketed sink blockable for blocking aggregate functions. 
`BucketedAggSinkLocalState` inherits `is_blockable() == false`, while the 
regular aggregate sink checks its evaluator. An eligible grouped `AI_AGG` query 
can issue a synchronous HTTP request from `add()` after its context fills, yet 
`PipelineTask::is_blockable()` will route this task to the CPU scheduler. 
Implement the evaluator check here (and account for source-side 
merge/finalization) so the request cannot occupy a CPU worker for its timeout.



##########
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:
   [P1] Gate the new plan node on the backend's actual capability. This check 
only excludes `isSmoothUpgradeSrc`, which is set for cloud smooth upgrade and 
remains false for a shared-nothing BE. After an FE-first upgrade with the sole 
BE still on the old version, the default-enabled bucketed path sends 
`BUCKETED_AGGREGATION_NODE` to that BE; its dispatcher has no case for the new 
enum, so eligible grouped queries fail. Check a BE capability or minimum 
version before enabling this path.



##########
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:
   [P2] Preserve query-cache admission for cache-enabled grouped queries. The 
regular two-phase plan has a cacheable LOCAL `AggregationNode` over the scan, 
but this fused path installs a `BucketedAggregationNode` without that cache 
point. `QueryCacheNormalizer` recognizes only `AggregationNode`, and the BE 
creates cache operators only for that node type, so with 
`enable_query_cache=true` an eligible grouped OLAP query loses cache hits. Skip 
fusion when query cache is enabled until both FE normalization and BE cache 
operators support this node.



##########
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:
   [P2] Fix the malformed `multiContains` assertion. The call closes after the 
first argument and leaves `[#", 4)` with mismatched delimiters, so Groovy 
cannot parse this new suite and none of its CSE checks run. Use 
`multiContains("cast(a as BIGINT) + cast(b as BIGINT)", 4)`.



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